Apache Flink DataStream 连接器全景指南:预定义源汇、官方连接器与接入方式详解
2026/9/23 14:47:00 网站建设 项目流程

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 作业的数据输入与输出,按来源可以分为三类:

  1. 预定义源与汇(Predefined Sources and Sinks):内置于 Flink 运行时,始终可用,无需额外依赖;
  2. Flink 项目官方连接器(Flink Project Connectors):由 Apache Flink 项目维护,随源码发布,用于对接各类第三方系统;
  3. 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。

两个必须注意的语义陷阱

  1. PROCESS_CONTINUOUSLY模式下,文件一旦被修改,其内容会被整体重新处理——即使只是在文件末尾追加数据,也会导致全部内容被重复处理,从而破坏 exactly-once 语义;
  2. 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自定义文件输出,支持自定义对象到字节的转换
writeToSocketSerializationSchema将元素写入 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 Kafkasource / sink分布式消息队列
Apache Cassandrasource / sink分布式 NoSQL 数据库
Amazon DynamoDBsinkAWS 托管 NoSQL 数据库
Amazon Kinesis Data Streamssource / sinkAWS 实时数据流服务
Amazon Kinesis Data FirehosesinkAWS 流数据投递服务
DataGensource内置测试数据生成器
Elasticsearchsink分布式搜索与分析引擎
Opensearchsink开源搜索与分析引擎
FileSystemsource / sink本地/分布式文件系统
RabbitMQsource / sinkAMQP 消息队列
Google PubSubsource / sinkGoogle 云消息服务
Hybrid Sourcesource多源顺序切换组合源
Apache Pulsarsource云原生消息流平台
JDBCsink各类关系型数据库
MongoDBsource / sink文档型 NoSQL 数据库

其中 DataGen、FileSystem、Hybrid Source 的详细文档位于本仓库 connectors/datastream 目录下(datagen.mdfilesystem.mdhybridsource.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 Kafkaexactly once使用与版本匹配的 Kafka 连接器
AWS Kinesis Streamsexactly once
RabbitMQat most once (v0.10) / exactly once (v1.0)
Google PubSubat least once
Collectionsexactly once
Filesexactly once
Socketsat most once

各官方 sink 的投递保证(假设状态更新为 exactly-once):

Sink保证备注
Elasticsearchat least once
Opensearchat least once
Kafka producerat least once / exactly oncev0.11+ 事务生产者可实现 exactly once
Cassandra sinkat least once / exactly once仅幂等更新时可 exactly once
Amazon DynamoDBat least once
Amazon Kinesis Data Streamsat least once
Amazon Kinesis Data Firehoseat least once
File sinksexactly once
Socket sinksat least once
Standard outputat least once
Redis sinkat least once

每个连接器的精确语义细节需查阅对应文档。据此可以得出实践准则:追求端到端 exactly-once 时,优先选择参与两阶段提交(如 Kafka 事务生产者、FileSink)的连接器,并保持状态存储的 exactly-once 配置。

七、Apache Bahir 扩展连接器

除 Flink 项目自身外,更多流式连接器通过Apache Bahir社区发布,包括:

连接器类型
Apache ActiveMQsource / sink
Apache Flumesink
Redissink
Akkasink
Nettysource

这些连接器以独立的 Flink 扩展形式提供,使用时同样需为作业引入对应依赖并部署相应的第三方服务。

八、不依赖连接器的接入方式:Async I/O 数据富化

使用连接器并不是把数据送入/送出 Flink 的唯一途径。一种常见模式是在MapFlatMap查询外部数据库或 Web 服务以富化(enrich)主数据流,例如为订单流补充用户画像、为日志流补充 IP 归属地。这种"每条记录一次外部调用"的模式如果写成同步阻塞调用,会严重拖慢吞吐。

为此 Flink 提供了 Async I/O API(AsyncFunction+DataStream.asyncWaitOperator),允许在等待外部请求返回期间继续处理其他记录,让富化类作业既高效又健壮。设计富化管道时,Async I/O 往往是比引入重量级连接器更轻量、更精准的答案。

九、选择建议与下一步

根据上述内容可以形成清晰的选型路径:

  1. 本地联调/演示:优先使用预定义源(集合、socket、文件)与 DataGen,零外部依赖;
  2. 生产实时接入:按消息/存储系统选择官方连接器(Kafka、Kinesis、Pulsar、RabbitMQ、JDBC 等),并核对 guarantees.md 中的语义保证;
  3. 历史数据 + 实时流:用 Hybrid Source 将 FileSource 与 KafkaSource 串成单一输入流;
  4. 字段级富化:不引入连接器,直接用 Async I/O 在算子内完成外部查询;
  5. 社区生态补充: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),仅供参考

需要专业的网站建设服务?

联系我们获取免费的网站建设咨询和方案报价,让我们帮助您实现业务目标

立即咨询