ARTICLE DETAIL

资讯详情

深耕郑州网站建设与运营推广的一线实战洞察。

Apache Flink DataStream 连接器全景指南:预定义源汇、官方连接器与接入方式详解

Apache Flink DataStream 连接器全景指南:预定义源汇、官方连接器与接入方式详解 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 SourcesFlink 内建的基础数据源开箱即用覆盖了文件、目录、Socket、集合与迭代器四类场景。它们全部定义在StreamExecutionEnvironment上对应实现位于 StreamExecutionEnvironment.java其中readTextFile/readFile系列在 L1635-L1925fromCollection/fromElements/fromSequence系列在 L1371-L1620socketTextStream在 L1927-L2010addSource在 L2145 之后。2.1 文件类数据源方法说明readTextFile(path)按行读取符合TextInputFormat规范的文本文件逐行返回StringreadFile(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将元素写入 SocketaddSink调用自定义 sink functionFlink 自带的连接器如 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.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]GeneratorFunctionLong, String generatorFunction index - Number: index; long numberOfRecords 1000; DataGeneratorSourceString source new DataGeneratorSource(generatorFunction, numberOfRecords, Types.STRING); DataStreamSourceString stream env.fromSource(source, WatermarkStrategy.noWatermarks(), Generator Source);输出顺序与并行度相关每个子序列内部按序产出如果并行度限制为 1则整体呈现从Number: 0到Number: 999的严格有序输出。限速Rate LimitingDataGeneratorSource内置限速能力。下面的代码让所有源子任务合计每秒不超过 100 条GeneratorFunctionLong, Long generatorFunction index - index; double recordsPerSecond 100; DataGeneratorSourceString 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 FileSourceString fileSource FileSource.forRecordStreamFormat(new TextLineInputFormat(), Path.fromLocalFile(testDir)).build(); KafkaSourceString kafkaSource KafkaSource.Stringbuilder() .setStartingOffsets(OffsetsInitializer.timestamp(switchTimestamp 1)) .build(); HybridSourceString hybridSource HybridSource.builder(fileSource) .addSource(kafkaSource) .build();方式二切换时刻动态定位。适用于文件源积压很大、处理时长可能超过下一源数据保留期retention的场景——切换必须发生在当前时间 - X。此时需要在切换时再确定下一源的起点通过实现SourceFactory接收上一个文件 enumerator 的结束位置延迟构造KafkaSourceFileSourceString fileSource CustomFileSource.readTillOneDayFromLatest(); HybridSourceString hybridSource HybridSource.String, CustomFileSplitEnumeratorbuilder(fileSource) .addSource( switchContext - { CustomFileSplitEnumerator previousEnumerator switchContext.getPreviousEnumerator(); // how to get timestamp depends on specific enumerator long switchTimestamp previousEnumerator.getEndTimestamp(); KafkaSourceString kafkaSource KafkaSource.Stringbuilder() .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 onceRabbitMQat most once (v0.10) / exactly once (v1.0)Google PubSubat least onceCollectionsexactly onceFilesexactly onceSocketsat most once各官方 sink 的投递保证假设状态更新为 exactly-onceSink保证备注Elasticsearchat least onceOpensearchat least onceKafka producerat least once / exactly oncev0.11 事务生产者可实现 exactly onceCassandra sinkat least once / exactly once仅幂等更新时可 exactly onceAmazon DynamoDBat least onceAmazon Kinesis Data Streamsat least onceAmazon Kinesis Data Firehoseat least onceFile sinksexactly onceSocket sinksat least onceStandard outputat least onceRedis sinkat least once每个连接器的精确语义细节需查阅对应文档。据此可以得出实践准则追求端到端 exactly-once 时优先选择参与两阶段提交如 Kafka 事务生产者、FileSink的连接器并保持状态存储的 exactly-once 配置。七、Apache Bahir 扩展连接器除 Flink 项目自身外更多流式连接器通过Apache Bahir社区发布包括连接器类型Apache ActiveMQsource / sinkApache FlumesinkRedissinkAkkasinkNettysource这些连接器以独立的 Flink 扩展形式提供使用时同样需为作业引入对应依赖并部署相应的第三方服务。八、不依赖连接器的接入方式Async I/O 数据富化使用连接器并不是把数据送入/送出 Flink 的唯一途径。一种常见模式是在Map或FlatMap中查询外部数据库或 Web 服务以富化enrich主数据流例如为订单流补充用户画像、为日志流补充 IP 归属地。这种每条记录一次外部调用的模式如果写成同步阻塞调用会严重拖慢吞吐。为此 Flink 提供了 Async I/O APIAsyncFunctionDataStream.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),仅供参考
返回列表