欢迎光临
我们一直在努力

Flink基础之Source API详解及代码实战演练:从入门到自定义

摘要

讲透 Flink Source API 的完整体系与实战用法:环境内置方法、KafkaSource / FileSource 等生产级 Connector、新 Source API 的 SourceReader 与 SplitEnumerator 架构、自定义 Source 的两种写法,附 4 个可直接运行的实战案例与 6 个真实踩坑点。

关键词

Flink、Source API、KafkaSource、FileSource、SourceReader、SplitEnumerator、SourceFunction、fromSource、数据分片、并行度


前面十一篇把 Flink 的架构、资源、部署都讲了一遍,但一直没认真聊过数据是从哪进来的。env.addSource(…) 可能是很多人的第一行 Flink 代码,但很少有人追问:Source 到底有几种写法、新老 API 差在哪、并行度为什么上不去。这篇把 Source 一次讲透——从一行代码的 fromElements 到需要实现三个接口的自定义 Source,配 4 个能直接跑起来的案例。

先说结论:官方 Connector 能覆盖的场景,永远不要自己写 Source。自定义 Source 的学习成本高、坑多,只该作为兜底手段。


一、Source API 分类全景

Flink 的 Source 有三类入口,按「省事程度」排序:

在这里插入图片描述

  • 环境方法:env.fromElements / fromCollection / fromSequence / socketTextStream / readTextFile,一行代码出流,适合开发调试。
  • Connector Source:KafkaSource、FileSource、PulsarSource 等官方实现,生产主力。
  • 自定义 Source:SourceFunction(旧)或新 API 的 Source + SourceReader + SplitEnumerator(新),兜底手段。
  • 选型路径:先找官方 Connector → 没有就找环境方法 → 都不行才自定义。

    还有一个背景必须了解:Flink 1.12 起把 Source 统一到新 Source API(env.fromSource(source, watermarkStrategy, name)),旧的 env.addSource(…) 和 SourceFunction 被标记废弃。老代码还能跑,新项目一律用 fromSource。


    二、实战一:环境内置方法,一行代码出流

    最快的上手方式,适合本地调试、单元测试、跑通链路:

    import org.apache.flink.streaming.api.datastream.DataStream;
    import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment;

    public class EnvSourceDemo {
    public static void main(String[] args) throws Exception {
    StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();

    // 从固定元素创建流(内存集合,数据量小)
    DataStream<String> elements = env.fromElements("a", "b", "c");

    // 从 Java 集合创建流
    DataStream<Integer> collection = env.fromCollection(
    java.util.Arrays.asList(1, 2, 3, 4));

    // 从数字序列创建流(常用于压力测试/模拟数据)
    DataStream<Long> sequence = env.fromSequence(1, 1000);

    // 从 Socket 读取(nc -lk 9999 起一个端口服务)
    DataStream<String> socket = env.socketTextStream("localhost", 9999);

    // 各自打印验证
    elements.print("elements");
    collection.print("collection");
    sequence.print("sequence");
    socket.print("socket");

    env.execute("env-source-demo");
    }
    }

    注意 socket 是无界流——作业会一直挂着等数据,本地测试完记得手动取消(或加 -d 后 flink cancel)。


    三、实战二:KafkaSource,生产最常用

    生产环境 80% 的 Source 是 Kafka。Flink 1.14+ 用新 API 的 KafkaSource 替代了老旧的 FlinkKafkaConsumer:

    import org.apache.flink.api.common.eventtime.WatermarkStrategy;
    import org.apache.flink.api.common.serialization.SimpleStringSchema;
    import org.apache.flink.connector.kafka.source.KafkaSource;
    import org.apache.flink.connector.kafka.source.enumerator.initializer.OffsetsInitializer;
    import org.apache.flink.streaming.api.datastream.DataStream;
    import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment;

    public class KafkaSourceDemo {
    public static void main(String[] args) throws Exception {
    StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();

    // ① 构建 KafkaSource:分区自动发现、offset 策略、反序列化
    KafkaSource<String> source = KafkaSource.<String>builder()
    .setBootstrapServers("localhost:9092")
    .setTopics("orders")
    .setGroupId("flink-order-group")
    // 首次从最早开始读(历史回填);committedOffsets 则从上次提交续读
    .setStartingOffsets(OffsetsInitializer.earliest())
    // 消息 value 反序列化为 String
    .setValueOnlyDeserializer(new SimpleStringSchema())
    .build();

    // ② 统一入口 fromSource:watermark 策略在这里传入
    DataStream<String> stream = env.fromSource(
    source,
    WatermarkStrategy.noWatermarks(), // 无事件时间时用 noWatermarks
    "Kafka Source");

    // 每行 "orderId,amount,category"
    stream.map(line -> {
    String[] p = line.split(",");
    return "order=" + p[0] + " amount=" + p[1];
    })
    .print();

    env.execute("kafka-source-demo");
    }
    }

    两个关键点:

  • fromSource 是唯一入口:KafkaSource 对象本身不直接产生流,必须通过 env.fromSource(source, watermarkStrategy, name) 接入。Watermark 策略(事件时间相关)也在这里统一传入。
  • offset 策略要按场景选:历史回填用 earliest,断点续跑用 committedOffsets,只关心新数据用 latest。首次上线配错 latest 会静默丢掉历史数据。

  • 四、实战三:FileSource,读文件/目录

    FileSource 是新 API 的文件 Source,支持读目录、通配符、分块并行:

    import org.apache.flink.api.common.eventtime.WatermarkStrategy;
    import org.apache.flink.api.common.serialization.SimpleStringSchema;
    import org.apache.flink.connector.file.src.FileSource;
    import org.apache.flink.core.fs.Path;
    import org.apache.flink.streaming.api.datastream.DataStream;
    import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment;

    public class FileSourceDemo {
    public static void main(String[] args) throws Exception {
    StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();

    // 按行读取 hdfs:///data/logs 目录下所有文件(可加通配符 *.log)
    FileSource<String> fileSource = FileSource.forRecordStreamFormat(
    new SimpleStringSchema(),
    new Path("hdfs:///data/logs"))
    .build();

    DataStream<String> lines = env.fromSource(
    fileSource, WatermarkStrategy.noWatermarks(), "File Source");

    lines.print();

    env.execute("file-source-demo");
    }
    }

    FileSource 的两个实用特性:

    • 有界/无界两种模式:默认有界(读完即止);配合 ContinuousFileMonitoringFunction 或文件监听可做无界(新文件到达即读)。
    • 分块并行:大文件自动按块切分(Split),并行度可以超过文件数。

    五、新 Source API 架构:SourceReader 与 SplitEnumerator

    理解了 Connector 用法,再看它底层的设计——新 Source API 由三个核心角色组成:

    在这里插入图片描述

    • Split(分片):可独立读取的数据单元。Kafka 一个分区 = 一个 Split,文件一个块 = 一个 Split。
    • SplitEnumerator(运行在 JobManager):发现所有 Split、分发给各 SourceReader、处理失败重分配、并行度调整时重新均衡。
    • SourceReader(运行在 TaskManager):每个并行实例一个,从分配的 Split 读数据、产出记录和水印、上报 Split 消费进度。

    这套设计解决了旧 SourceFunction 的三个硬伤:

  • 并行度可动态调整:Split 可以重新分配,旧 API 的并行实例绑定固定数据;
  • 精确一次有保证:每个 Split 的读取进度随 Checkpoint 持久化,故障后精确恢复;
  • 水印与空闲检测内置:没有数据的分区不会卡住整体水印推进。

  • 六、实战四:自定义 Source 的两种写法

    官方没有现成 Connector 时(比如对接内部消息系统),才有必要自定义。两种写法:

    6.1 旧 API:SourceFunction(简单,但已废弃)

    import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment;
    import org.apache.flink.streaming.api.functions.source.SourceFunction;

    public class CustomSourceDemo {
    public static void main(String[] args) throws Exception {
    StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();

    // 旧 API:实现 SourceFunction,每秒产出一条模拟数据
    env.addSource(new SourceFunction<String>() {
    private volatile boolean running = true;

    @Override
    public void run(SourceContext<String> ctx) throws Exception {
    int i = 0;
    while (running) {
    ctx.collect("event-" + i++);
    Thread.sleep(1000);
    }
    }

    @Override
    public void cancel() {
    running = false; // 取消时置位,让 run 循环退出
    }
    }).print();

    env.execute("custom-source-demo");
    }
    }

    必须知道的坑:SourceFunction 没有自动的 Checkpoint 支持——ctx.collect 的数据进度不会随 Checkpoint 保存,故障重启会从 run() 开头重新执行,造成数据重复或丢失。只在模拟数据、无状态演示场景用它,生产慎用。

    6.2 新 API:三件套(复杂,但语义完整)

    新 API 需要实现 Source、SourceReader、SplitEnumerator 三个接口,代码量百行起步,日常开发几乎不会手写——这也是「先找官方 Connector」这条原则的底气:官方已经帮你把最难的框架部分做好了。


    七、Source 并行度与数据分片

    Source 的并行度有一个被广泛误解的约束:Source 实际并行度 = min(配置并行度, Split 数)。

    在这里插入图片描述

    • Kafka:并行度 ≤ 分区数。设并行度 8、topic 只有 4 个分区 → 4 个实例空转。想提吞吐先扩分区,改并行度没用——这是最常见的无效调优。
    • 文件:并行度 ≤ 文件数/分块数。一个文件一个 Split,文件少并行度上不去;大文件自动分块后并行度可提升。
    • 集合/socket:fromElements 并行度恒为 1,socket 恒为 1,多并行无意义。

    记住:Split 数才是 Source 吞吐的天花板。


    八、六个真实踩坑

  • 还在用 FlinkKafkaConsumer + addSource。旧 API 在新版本会打废弃警告,且不享受新 Source API 的 Split 管理、动态并行度能力。迁到 KafkaSource + fromSource。
  • socket/无界源测试完不取消,作业一直挂着。本地调试用 socket 记得 flink cancel 或直接 Ctrl+C。
  • fromElements 塞大数据集。它是内存集合直接转流,数据量大时 JVM 直接 OOM。大数据量用 fromSequence 模拟或直接上 Kafka/文件。
  • 自定义 SourceFunction 当成生产 Source 用。没有 Checkpoint 进度管理,故障恢复会重复/丢失数据。生产级自定义必须走新 API 三件套。
  • fromSource 忘传 Watermark 策略。事件时间作业里,watermark 策略必须在 fromSource 传入;用 noWatermarks 时窗口永远不触发(事件时间语义下)。
  • FileSource 读「单个文件」失败。FileSource 面向目录/通配符读取;只读一个已知文件用 env.readTextFile(旧)或明确指定文件路径。另外要确认路径是 HDFS 还是本地——本地路径在集群模式读不到。

  • Source 是 Flink 作业的「水龙头」,拧开方式决定了水的品质:环境方法适合调试,Connector 是生产主力,自定义是最后手段。把「新 API 的 Split 机制」「并行度受 Split 数约束」「SourceFunction 无 Checkpoint」这三件事记牢,数据入口这一环就不会再出低级问题。
    ce 面向目录/通配符读取;只读一个已知文件用 env.readTextFile(旧)或明确指定文件路径。另外要确认路径是 HDFS 还是本地——本地路径在集群模式读不到。


    赞(0)
    未经允许不得转载:171主机测评 » Flink基础之Source API详解及代码实战演练:从入门到自定义
    分享到: 更多 (0)

    评论 抢沙发

    • 昵称 (必填)
    • 邮箱 (必填)
    • 网址