摘要
讲透 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 有三类入口,按「省事程度」排序:

选型路径:先找官方 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");
}
}
两个关键点:
四、实战三: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 的三个硬伤:
六、实战四:自定义 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 吞吐的天花板。
八、六个真实踩坑
Source 是 Flink 作业的「水龙头」,拧开方式决定了水的品质:环境方法适合调试,Connector 是生产主力,自定义是最后手段。把「新 API 的 Split 机制」「并行度受 Split 数约束」「SourceFunction 无 Checkpoint」这三件事记牢,数据入口这一环就不会再出低级问题。
ce 面向目录/通配符读取;只读一个已知文件用 env.readTextFile(旧)或明确指定文件路径。另外要确认路径是 HDFS 还是本地——本地路径在集群模式读不到。




