Apache Flink DataStream 连接器全景指南:预定义源汇、官方连接器与接入方式详解
【免费下载链接】flink项目地址: https://gitcode.com/gh_mirrors/fli/flink
导读
本文以 Apache Flink DataStream API 的连接器体系为核心,系统梳理 Flink 提供数据接入/输出的三种主要途径:始终可用的预定义数据源与数据汇、随 Flink 项目源码发布的官方连接器(Kafka、Kinesis、FileSystem、JDBC、DataGen、Hybrid Source 等),以及 Apache Bahir 社区维护的扩展连接器。读完本文,你将掌握如何为自己的 DataStream 作业选择正确的连接器、如何引入对应依赖、如何在无外部系统的情况下用 DataGen 快速构造测试数据、如何用 Hybrid Source 平滑完成"先读历史数据、再接实时流"的典型切换,以及不同连接器的端到端容错保证差异。
一、连接器的三种来源
Flink DataStream 作业的数据输入与输出,按来源可以分为三类:
- 预定义源与汇(Predefined Sources and Sinks):内置于 Flink 运行时,始终可用,无需额外依赖;
- Flink 项目官方连接器(Flink Project Connectors):由 Apache Flink 项目维护,随源码发布,用于对接各类第三方系统;
- Apache Bahir 连接器:由 Apache Bahir 社区维护发布,覆盖 ActiveMQ、Flume、Redis 等更多外部系统。
对应地,DataStream 编程指南中的数据源与数据汇章节指出,DataStream初始数据由各类 source 创建,结果通过 sink 写回外部系统,例如文件或标准输出;程序可以在本地 JVM 中执行,也可以提交到集群。理解这三类接入方式,是设计任何 DataStream 应用的第一步。
二、预定义数据源(Predefined Data Sources)
Flink 内建的基础数据源"开箱即用",覆盖了文件、目录、Socket、集合与迭代器四类场景。它们全部定义在StreamExecutionEnvironment上,对应实现位于 StreamExecutionEnvironment.java(其中readTextFile/readFile系列在 L1635-L1925,fromCollection/fromElements/fromSequence系列在 L1371-L1620,socketTextStream在 L1927-L2010,addSource在 L2145 之后)。
2.1 文件类数据源
| 方法 | 说明 |
|---|---|
readTextFile(path) | 按行读取符合TextInputFormat规范的文本文件,逐行返回String |
readFile(fileInputFormat, path) | 按指定的文件输入格式读取(一次性)文件 |
readFile(fileInputFormat, path, watchType, interval, pathFilter, typeInfo) | 前两个方法内部调用的完整版本,支持持续监控或一次性处理 |
第三种readFile是核心实现:根据watchType决定处理模式——FileProcessingMode.PROCESS_CONTINUOUSLY会每隔interval毫秒周期性扫描目录中的新数据;FileProcessingMode.PROCESS_ONCE则只处理路径下当前已有的数据后退出。pathFilter可进一步排除不需要处理的文件。
底层实现机制:Flink 将文件读取拆分为"目录监控"与"数据读取"两个子任务。监控由一个**非并行(并行度=1)**的任务完成,负责周期(或一次性)扫描目录、发现待处理文件、将文件划分为 splits 并分发给下游读取任务;读取由多个与作业并行度相等的任务并行执行,每个 split 只会被一个 reader 读取,而一个 reader 可以依次读取多个 split。
两个必须注意的语义陷阱:
PROCESS_CONTINUOUSLY模式下,文件一旦被修改,其内容会被整体重新处理——即使只是在文件末尾追加数据,也会导致全部内容被重复处理,从而破坏 exactly-once 语义;PROCESS_ONCE模式下,source 扫描完路径即退出,不会等待 readers 读完文件内容(readers 会继续读到全部数据)。此后不再产生新的 checkpoint,节点故障后作业只能从最近一次 checkpoint 恢复,恢复速度可能变慢。
2.2 Socket 数据源
socketTextStream(hostname, port)从 Socket 读取数据,支持自定义元素分隔符(delimiter参数),是本地联调时最常用的实时输入方式之一。
2.3 集合与迭代器数据源
| 方法 | 说明 |
|---|---|
fromCollection(Collection) | 从java.util.Collection创建流,集合内元素必须类型一致 |
fromCollection(Iterator, Class) | 从迭代器创建流,Class指明元素类型 |
fromElements(T...) | 从给定对象序列创建流,对象类型必须一致 |
fromParallelCollection(SplittableIterator, Class) | 从可拆分迭代器并行创建流 |
fromSequence(from, to) | 并行生成区间内的数字序列 |
集合类数据源特别适合测试:先在本地用fromElements/fromCollection验证逻辑,再无缝替换为读取外部系统的真实连接器。注意,集合数据源要求元素类型及迭代器实现Serializable,且不支持并行执行(并行度固定为 1)。
2.4 自定义源:addSource / fromSource
通过addSource(new SomeSourceFunction<>(...))可以挂载任意自定义 source function。实现自定义源时,非并行源实现SourceFunction,并行源实现ParallelSourceFunction或继承RichParallelSourceFunction。新一代连接器(如 Kafka、FileSystem)则通过env.fromSource(source, watermarkStrategy, "sourceName")接入——这也是 DataGen、FileSystem、Hybrid Source 等文档中推荐的标准用法。
三、预定义数据汇(Predefined Data Sinks)
预定义 sink 支持写入文件、标准输出/标准错误、Socket,全部封装为DataStream上的操作:
| 方法 / 输出格式 | 说明 |
|---|---|
writeAsText()/TextOutputFormat | 逐行将元素以toString()结果写出 |
writeAsCsv(...)/CsvOutputFormat | 将 Tuple 写为逗号分隔文件,行列分隔符可配置 |
print()/printToErr() | 将toString()值打印到标准输出/标准错误;可指定前缀区分多个 print;并行度大于 1 时输出会带任务标识 |
writeUsingOutputFormat()/FileOutputFormat | 自定义文件输出,支持自定义对象到字节的转换 |
writeToSocket | 按SerializationSchema将元素写入 Socket |
addSink | 调用自定义 sink function,Flink 自带的连接器(如 Kafka)即以 sink function 形式实现 |
重要提醒:write*()系列方法主要面向调试,它们不参与 Flink 的 checkpoint 机制,通常只有 at-least-once 语义——数据是否及时冲刷到目标系统取决于OutputFormat的实现,故障场景下部分记录可能丢失。若需要可靠的、端到端 exactly-once 的文件写入,请使用FileSink(详见 FileSystem 连接器文档);通过.addSink(...)实现的自定义 sink 若正确接入 checkpoint,同样可以获得 exactly-once 语义。
四、Flink 项目官方连接器一览
连接器为对接各类第三方系统提供了现成代码。作为 Apache Flink 项目的一部分,当前官方支持以下系统(标注其是作为 source 还是 sink 使用):
| 连接器 | 类型 | 对接系统 |
|---|---|---|
| Apache Kafka | source / sink | 分布式消息队列 |
| Apache Cassandra | source / sink | 分布式 NoSQL 数据库 |
| Amazon DynamoDB | sink | AWS 托管 NoSQL 数据库 |
| Amazon Kinesis Data Streams | source / sink | AWS 实时数据流服务 |
| Amazon Kinesis Data Firehose | sink | AWS 流数据投递服务 |
| DataGen | source | 内置测试数据生成器 |
| Elasticsearch | sink | 分布式搜索与分析引擎 |
| Opensearch | sink | 开源搜索与分析引擎 |
| FileSystem | source / sink | 本地/分布式文件系统 |
| RabbitMQ | source / sink | AMQP 消息队列 |
| Google PubSub | source / sink | Google 云消息服务 |
| Hybrid Source | source | 多源顺序切换组合源 |
| Apache Pulsar | source | 云原生消息流平台 |
| JDBC | sink | 各类关系型数据库 |
| MongoDB | source / sink | 文档型 NoSQL 数据库 |
其中 DataGen、FileSystem、Hybrid Source 的详细文档位于本仓库 connectors/datastream 目录下(datagen.md、filesystem.md、hybridsource.md),其余连接器的接入说明以对应版本的官方发行文档为准。
4.1 使用官方连接器的注意事项
- 通常需要额外的第三方组件:Kafka 连接器需要可访问的 Kafka 集群,JDBC 连接器需要可访问的数据库服务,文件/消息队列类连接器同样如此;
- 不在二进制发行版中:尽管这些流式连接器属于 Flink 项目、包含在源码发行版中,但它们不包含在官方预编译的二进制发行版里,需要用户按各连接器子章节的说明自行添加对应依赖。
五、重点连接器实战
5.1 DataGen:无外部系统时的数据生成利器
DataGen 连接器文档 介绍了一个内置、无需额外依赖的 Source 实现DataGeneratorSource,非常适合在本地开发或演示时模拟输入数据,避免依赖 Kafka 等外部系统。其实现位于 DataGeneratorSource.java,配合 GeneratorFunction.java 使用。
工作原理:DataGeneratorSource并行产生 N 条数据,它会将序号序列切分为与源子任务数量相等的并行子序列,并把类型为Long的"索引"值交给用户提供的GeneratorFunction,由它将子序列映射为任意类型的生成事件。例如下面这段代码生成["Number: 0", "Number: 1", ..., "Number: 999"]:
GeneratorFunction<Long, String> generatorFunction = index -> "Number: " + index; long numberOfRecords = 1000; DataGeneratorSource<String> source = new DataGeneratorSource<>(generatorFunction, numberOfRecords, Types.STRING); DataStreamSource<String> stream = env.fromSource(source, WatermarkStrategy.noWatermarks(), "Generator Source");输出顺序与并行度相关:每个子序列内部按序产出;如果并行度限制为 1,则整体呈现从"Number: 0"到"Number: 999"的严格有序输出。
限速(Rate Limiting):DataGeneratorSource内置限速能力。下面的代码让所有源子任务合计每秒不超过 100 条:
GeneratorFunction<Long, Long> generatorFunction = index -> index; double recordsPerSecond = 100; DataGeneratorSource<String> source = new DataGeneratorSource<>( generatorFunction, Long.MAX_VALUE, RateLimiterStrategy.perSecond(recordsPerSecond), Types.STRING);RateLimiterStrategy还提供按 checkpoint 限制产出记录数等其他策略(定义于org.apache.flink.api.connector.source.util.ratelimit.RateLimiterStrategy)。
有界性:该 source 本质上始终有界,但把记录数设为Long.MAX_VALUE后,从实际效果看就变成了"永不结束"的无界源;对有限序列,文档建议在BATCH执行模式下运行作业。
确定性要求:若GeneratorFunction对相同输入Long始终输出相同结果(即输出相对于输入确定),则该源可用于构建 at-least-once 与端到端 exactly-once 语义的作业;同时也可以基于生成事件与自定义WatermarkStrategy在源端直接产出确定性的 watermark。
5.2 Hybrid Source:异源顺序切换的单一输入流
Hybrid Source 文档 介绍了一种包含多个具体 source 的组合源,解决"从异构来源顺序读取、汇成单一输入流"的问题,其核心实现位于 HybridSource.java。
典型场景是启动引导(bootstrap):先读取 S3 上数天的有界历史数据,再无缝切换到 Kafka 的最新无界实时流。HybridSource会在有界文件输入结束后自动从FileSource切换到KafkaSource,整个过程不中断应用。在HybridSource出现之前,用户必须在拓扑中自行创建多个 source 并手工实现切换机制,既增加运维复杂度又损失效率;而使用HybridSource后,多个源在作业图与DataStreamAPI 视角下就是一个单一 source。
使用它需要引入flink-connector-base依赖(通常作为具体连接器的传递依赖一并获得)。
切换位置的两种设定方式:
方式一:构图时固定起点。适用于各源覆盖范围预先可知的场景——比如文件读到预定切换时间点后,继续从 Kafka 读取:
long switchTimestamp = ...; // derive from file input paths FileSource<String> fileSource = FileSource.forRecordStreamFormat(new TextLineInputFormat(), Path.fromLocalFile(testDir)).build(); KafkaSource<String> kafkaSource = KafkaSource.<String>builder() .setStartingOffsets(OffsetsInitializer.timestamp(switchTimestamp + 1)) .build(); HybridSource<String> hybridSource = HybridSource.builder(fileSource) .addSource(kafkaSource) .build();方式二:切换时刻动态定位。适用于文件源积压很大、处理时长可能超过下一源数据保留期(retention)的场景——切换必须发生在"当前时间 - X"。此时需要在切换时再确定下一源的起点,通过实现SourceFactory接收上一个文件 enumerator 的结束位置,延迟构造KafkaSource:
FileSource<String> fileSource = CustomFileSource.readTillOneDayFromLatest(); HybridSource<String> hybridSource = HybridSource.<String, CustomFileSplitEnumerator>builder(fileSource) .addSource( switchContext -> { CustomFileSplitEnumerator previousEnumerator = switchContext.getPreviousEnumerator(); // how to get timestamp depends on specific enumerator long switchTimestamp = previousEnumerator.getEndTimestamp(); KafkaSource<String> kafkaSource = KafkaSource.<String>builder() .setStartingOffsets(OffsetsInitializer.timestamp(switchTimestamp + 1)) .build(); return kafkaSource; }, Boundedness.CONTINUOUS_UNBOUNDED);注意该方式要求 enumerator 支持获取结束时间戳,当前可能需要对源做定制(FileSource对动态结束位置的支持跟踪于 FLINK-23633);同一源码树中的 HybridSourceSplitEnumerator.java 展示了各子源 split 的枚举与交接实现。
5.3 FileSystem 连接器
FileSystem 连接器(source/sink)的完整配置(包括FileSink的行编码、滚动策略、分桶(bucketing)与分区提交等)请参阅 filesystem.md 及 格式化器目录(CSV、JSON、Avro、Parquet、Hadoop、文本文件等)。
六、端到端容错保证:选择连接器前必读
不同连接器参与 Flink checkpoint/快照机制的程度不同,直接决定端到端投递语义。Fault Tolerance Guarantees 文档 给出了官方结论:只有当 source 参与快照机制时,Flink 才能保证用户状态更新的 exactly-once;同样,要获得端到端 exactly-once 投递,sink 也必须参与 checkpoint。
内置/官方 source 的状态更新保证:
| Source | 保证 | 备注 |
|---|---|---|
| Apache Kafka | exactly once | 使用与版本匹配的 Kafka 连接器 |
| AWS Kinesis Streams | exactly once | |
| RabbitMQ | at most once (v0.10) / exactly once (v1.0) | |
| Google PubSub | at least once | |
| Collections | exactly once | |
| Files | exactly once | |
| Sockets | at most once |
各官方 sink 的投递保证(假设状态更新为 exactly-once):
| Sink | 保证 | 备注 |
|---|---|---|
| Elasticsearch | at least once | |
| Opensearch | at least once | |
| Kafka producer | at least once / exactly once | v0.11+ 事务生产者可实现 exactly once |
| Cassandra sink | at least once / exactly once | 仅幂等更新时可 exactly once |
| Amazon DynamoDB | at least once | |
| Amazon Kinesis Data Streams | at least once | |
| Amazon Kinesis Data Firehose | at least once | |
| File sinks | exactly once | |
| Socket sinks | at least once | |
| Standard output | at least once | |
| Redis sink | at least once |
每个连接器的精确语义细节需查阅对应文档。据此可以得出实践准则:追求端到端 exactly-once 时,优先选择参与两阶段提交(如 Kafka 事务生产者、FileSink)的连接器,并保持状态存储的 exactly-once 配置。
七、Apache Bahir 扩展连接器
除 Flink 项目自身外,更多流式连接器通过Apache Bahir社区发布,包括:
| 连接器 | 类型 |
|---|---|
| Apache ActiveMQ | source / sink |
| Apache Flume | sink |
| Redis | sink |
| Akka | sink |
| Netty | source |
这些连接器以独立的 Flink 扩展形式提供,使用时同样需为作业引入对应依赖并部署相应的第三方服务。
八、不依赖连接器的接入方式:Async I/O 数据富化
使用连接器并不是把数据送入/送出 Flink 的唯一途径。一种常见模式是在Map或FlatMap中查询外部数据库或 Web 服务以富化(enrich)主数据流,例如为订单流补充用户画像、为日志流补充 IP 归属地。这种"每条记录一次外部调用"的模式如果写成同步阻塞调用,会严重拖慢吞吐。
为此 Flink 提供了 Async I/O API(AsyncFunction+DataStream.asyncWaitOperator),允许在等待外部请求返回期间继续处理其他记录,让富化类作业既高效又健壮。设计富化管道时,Async I/O 往往是比引入重量级连接器更轻量、更精准的答案。
九、选择建议与下一步
根据上述内容可以形成清晰的选型路径:
- 本地联调/演示:优先使用预定义源(集合、socket、文件)与 DataGen,零外部依赖;
- 生产实时接入:按消息/存储系统选择官方连接器(Kafka、Kinesis、Pulsar、RabbitMQ、JDBC 等),并核对 guarantees.md 中的语义保证;
- 历史数据 + 实时流:用 Hybrid Source 将 FileSource 与 KafkaSource 串成单一输入流;
- 字段级富化:不引入连接器,直接用 Async I/O 在算子内完成外部查询;
- 社区生态补充:Bahir 覆盖 ActiveMQ、Redis 等未进入 Flink 主项目的系统。
深入阅读可继续探索 DataStream 编程指南(预定义源汇的完整方法语义与示例)、DataGen、FileSystem、Hybrid Source 与 guarantees,并结合 StreamExecutionEnvironment.java、DataGeneratorSource.java、HybridSource.java 阅读底层实现,做到"知其然亦知其所以然"。
【免费下载链接】flink项目地址: https://gitcode.com/gh_mirrors/fli/flink
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考