多彩编程 多彩编程MZPH · CODE BLOG
ARTICLE DETAIL

文章详情

深耕前端与后端开发技术的一线实战笔记与踩坑复盘。

Apache Beam KafkaIO 实战指南:Kafka 读写配置、动态读取与精确一次写入原理解析

Apache Beam KafkaIO 实战指南:Kafka 读写配置、动态读取与精确一次写入原理解析 【免费下载链接】beamApache Beam is a unified programming model for Batch and Streaming data processing.项目地址https://gitcode.com/gh_mirrors/beam18/beam点击查看免费下载Apache Beam 的 KafkaIO 是一组用于从 Apache Kafka 读取消息、向 Kafka 写入消息的 I/O transforms。本指南以仓库中 KafkaIO 模块 及其核心实现 KafkaIO.java 的 JavaDoc 为主体完整覆盖依赖配置、读取/写入 API 的每一项参数、事件时间与水印策略、动态读取、Schema Registry 集成以及 Exactly-once 写入等实战细节并辅以源码级实现证据。读完本文你将能够独立完成一个基于 Beam Kafka 的流式管道从配置 Maven 依赖开始正确写出可运行的多主题读取、动态 topic 发现、Avro 反序列化与精确一次写入管道并理解底层 runner 如何通过 checkpoint 与 offset 管理保证数据语义。一、依赖配置beam-sdks-java-io-kafka 与 kafka-clients使用 KafkaIO 前必须在项目中加入对beam-sdks-java-io-kafka的依赖。需要特别强调的是KafkaIO 在运行时支持多种版本的 Kafka client它不会传递性地拉取某个特定版本的kafka-clients因此你需要自己显式声明一个兼容的kafka-clients作为运行时依赖。通常当前版本及较新版本的 Kafka 都受支持完整的支持版本清单以 KafkaIO 的 JavaDoc 为准。Maven 依赖声明如下版本号需替换为你实际使用的 Beam 与 Kafka 版本dependency groupIdorg.apache.beam/groupId artifactIdbeam-sdks-java-io-kafka/artifactId version.../version /dependency dependency groupIdorg.apache.kafka/groupId artifactIdkafka-clients/artifactId versiona_recent_version/version scoperuntime/scope /dependency从源码来看KafkaIO 与 Kafka client 的所有交互都通过kafka-clients完成模块内置的 ConsumerSpEL.java 使用 SpEL 反射方式调用不同版本 client 的 API如offsetsForTimes这正是它能在不强制绑定 client 版本的前提下兼容多版本的原因。关于版本兼容性KafkaIO.java 的 JavaDoc 明确指出运行时支持kafka-clients0.10.1 及以上版本0.9.x 至 0.10.0.0 这些更老的版本虽然仍可用但已弃用、可能在不久的将来被移除。请务必保证应用引入的 client 版本与 Kafka 集群版本兼容——如果版本不兼容Kafka client 通常会在初始化时报出清晰的错误信息。当检测到 client 版本过老时KafkaIO 还会在运行时打出 WARN 日志提示升级见expand()中对ConsumerSpEL.hasOffsetsForTimes()的判断。仓库中还保留了 kafka-01103、kafka-100、kafka-251 等多个版本目录以及 kafka-integration-test.gradle用于针对不同 Kafka 版本做集成验证。二、从 Kafka 读取KafkaIO.read() 基础用法KafkaIO 的读侧是一个UnboundedSource返回的是一束无界的PCollectionKafkaRecordK, V。KafkaRecord除了携带 key 和 value 之外还包含 topic、partition、offset、timestamp、timestampType、headers 等 Kafka 元数据。虽然大多数应用只消费单个 topic但 source 可以配置为消费多个 topic甚至可以精确到一组TopicPartition。配置一个 Kafka source 至少有四个必填项bootstrapServersKafka broker 地址列表一个或多个 topicwithTopic()或withTopics(ListString)key 反序列化器withKeyDeserializer(...)value 反序列化器withValueDeserializer(...)。最小可运行示例来自 KafkaIO.java 的 JavaDocpipeline .apply(KafkaIO.Long, Stringread() .withBootstrapServers(broker_1:9092,broker_2:9092) .withTopic(my_topic) // 多主题请用 withTopics(ListString) .withKeyDeserializer(LongDeserializer.class) .withValueDeserializer(StringDeserializer.class) // 以上四项为必填返回 PCollectionKafkaRecordLong, String // —— 以下均为可选配置 —— // 通过 ConsumerConfig 进一步定制读取用的 KafkaConsumer例如设置 group.id .withConsumerConfigUpdates(ImmutableMap.of(group.id, my_beam_app_1)) // 基于 Kafka LogAppendTime 设置事件时间与水印 // 默认是 withProcessingTime()对使用 CreateTime 时间戳的 topic 用 withCreateTime() .withLogAppendTime() // 只读取已提交committed的消息 .withReadCommitted() // 管道消费到的 offset 在 finalize 阶段回传给 Kafka .commitOffsetsInFinalize() // 提供可序列化函数决定运行时是否停止读取某个 TopicPartition // 注意只有基于 ReadFromKafkaDoFn 的实现会响应此信号 .withCheckStopReadingFn(new SerializedFunctionTopicPartition, Boolean() {}) // 将解析失败的坏消息送到备选 sinkErrorHandler 模式 .withBadRecordErrorHandler(errorHandler) // 若不需要 Kafka 元数据可以丢弃它 .withoutMetadata() // PCollectionKVLong, String ) .apply(Values.Stringcreate()) // PCollectionString ...2.1 Topic 的三种指定方式三选一源码中的withTopics、withTopicPartitions、withTopicPattern三个方法相互排斥同时设置会触发checkState异常withTopic(String)/withTopics(ListString)读取这些 topic 的所有 partitionwithTopicPartitions(ListTopicPartition)只读取指定的 partition 子集withTopicPattern(String)按正则匹配的 topic 的所有 partition。expand()方法中同样做了强制校验非动态读取时withTopic/withTopics/withTopicPartitions/withTopicPattern必须设置其一否则抛异常提示。2.2 反序列化器与 Coder 推断Kafka 提供了常见类型的反序列化器位于org.apache.kafka.common.serialization如LongDeserializer、StringDeserializer。除了反序列化器Beam runner 在必要时还需要Coder来物化 key 和 value 对象。大多数情况下Coder 可以由反序列化器类型自动推断出来无需显式指定。但在推断失败时可以连同反序列化器一起显式指定 CoderwithKeyDeserializerAndCoder(Class, Coder)/withValueDeserializerAndCoder(Class, Coder)针对DeserializerProvider变体还有withKeyDeserializerProviderAndCoder(...)/withValueDeserializerProviderAndCoder(...)。注意Kafka 消息是通过 key/value反序列化器解释的。若 key 或 value 反序列化器未设置expand()会直接报错 withKeyDeserializer() is required。从源码Read.Builder.resolveCoder(...)看Coder 推断逻辑是反射读取反序列化器deserialize方法的返回类型byte[]对应NullableCoderByteArrayCoder、Integer对应NullableCoderVarIntCoder、Long对应NullableCoderVarLongCoder其他类型则抛出异常提示无法推断。2.3 便利入口与常见辅助方法KafkaIO.readBytes()直接返回一个 key/value 均为byte[]的Read内部已自动设置ByteArrayDeserializerKafkaIO.read()创建未初始化的Read默认值为空 topics、maxNumRecords Long.MAX_VALUE、关闭 offset 提交、关闭动态读取、时间戳策略为ProcessingTime见 KafkaIO.java 的 read() 实现withoutMetadata()丢弃 Kafka 元数据输出PCollectionKVK, V。三、分区分配与 Checkpointing初始 offset 的确定顺序Kafka 的 partition 会在各 splitworker之间均匀分布。Checkpointing 得到完整支持每个 split 都可以从之前的 checkpoint 恢复具体能力受 runner 限制split 与 checkpoint 的细节参见KafkaUnboundedSource.split(int, PipelineOptions)。当管道首次启动、或没有任何 checkpoint时source 默认从latestoffset 开始消费。你可以通过Read.withConsumerConfigUpdates(Map)修改 ConsumerConfig 覆盖该行为例如改为 earliest也可以开启 Kafka 的 offset auto commit 以从上次提交点恢复。综合来看KafkaIO.read设置初始 offset 遵循以下顺序JavaDoc 原文runner 提供的KafkaCheckpointMarkcheckpoint 优先当ConsumerConfig.ENABLE_AUTO_COMMIT_CONFIG true时使用 Kafka 中存储的 consumer offset默认从latestoffset 开始。从 KafkaIOUtils.java 可以看到默认的 Consumer 配置key/value 反序列化器默认ByteArrayDeserializer、receive.buffer.bytes 512 * 1024注释说明这是针对 KAFKA-3135 的大接收缓冲能数倍提升吞吐、auto.offset.reset latest、enable.auto.commit false。另一个值得注意的实现细节是KafkaIO.read()后台实际运行着两个 consumer——主 consumer 读取数据副 offset consumer 用于估算 backlog拉取各 partition 的最新 offset。默认情况下 offset consumer 继承主 consumer 配置并带一个自动生成的 group.id这在需要额外配置的安全 Kafka 集群上可能不工作因此提供了withOffsetConsumerConfigOverrides(Map)为其补充配置当看到 exception while fetching latest offset for partition {} 的 WARN 日志时就该检查这项配置。此外Kafka 的 seek-to-initial-offset 是阻塞操作对某些旧版 client 可能永久阻塞该问题由 KIP-266 通过default.api.timeout.ms解决。KafkaIO.read 自身实现了超时逻辑避免永久阻塞同时会识别并尊重你在 consumer config 中传入的default.api.timeout.ms。四、按时间戳定位withStartReadTime 与 withStopReadTimewithStartReadTime(Instant)使用时间戳来设置起始 offsetwithStopReadTime(Instant)设置结束 offset。两者都仅受 Kafka Client 0.10.1.0 及以上版本支持且要求消息格式版本在 0.10.0 之后消息自带时间戳。关于withStartReadTime有两个要点它的优先级高于ConsumerConfig.AUTO_OFFSET_RESET_CONFIG和任何自动提交的 offset在以下两种情况下会硬失败① 某些 partition 中不存在时间戳大于等于目标时间戳的消息② 某 partition 的消息格式版本早于 0.10.0即消息没有时间戳。当 runner 不支持offsetsForTimes而你又设置了 start/stop time 时expand()会抛出参数校验异常并提示通过 Maven 属性kafka.clients.version将 client 升级到 0.10.1.0 以上。另外若设置了withStopReadTime基于 SDF 的读实现会自动把 source 标记为有界withBounded()。五、动态读取运行时发现 Topic / Partition对于给定的 bootstrap_serverKafkaIO 还能动态发现并读取可用的TopicPartition并在不需要时停止读取。动态读取由两个组件协作完成WatchForKafkaTopicPartitions负责周期性扫描并发射新出现的TopicPartitionReadFromKafkaDoFn负责消费每个KafkaSourceDescriptor。动态读取主要解决两类场景某个 topic 或 partition 被添加/删除某个 topic 或 partition 被添加、删除后又重新添加。在提供了checkStopReadingFn的前提下动态读取还能额外处理某个 topic 或 partition 被停止某个 topic 或 partition 被添加、停止后又重新添加。启用方式JavaDoc 示例pipeline .apply(KafkaIO.Long, Stringread() // 每 1 小时扫描一次可用的 TopicPartition 并发射新增的 .withDynamicRead(Duration.standardHours(1)) .withCheckStopReadingFn(new SerializedFunctionTopicPartition, Boolean() {}) .withBootstrapServers(broker_1:9092,broker_2:9092) .withKeyDeserializer(LongDeserializer.class) .withValueDeserializer(StringDeserializer.class) ) .apply(Values.Stringcreate()) // PCollectionString ...使用注意withDynamicRead(Duration)会开启动态读取并设置扫描周期不传则默认周期为 1 小时。开启动态读取后管道构建时不再强制要求提供 topic/partition但运行时要求开启beam_fn_api实验特性——expand()中会校验ExperimentalOptions.hasExperiment(pipelineOptions, beam_fn_api)否则报错 Kafka Dynamic Read requires enabling experiment beam_fn_api。JavaDoc 还详细描述了两种已知的竞态条件建议在遇到重复或漏读时排查stop/remove 后重新添加导致无法再次发射例如WatchForKafkaTopicPartitions每 1 小时更新跟踪集合10:00 跟踪到 partition A 并开始读取10:30 发现其被 stop/remove 而停止读取并返回ProcessContinuation.stop()10:45 用户想重新读取 A但 11:00 扫描时它只知道 A 仍是活跃 partition不会再发射 Astop/remove 与重新添加交错导致重复记录ReadFromKafkaDoFn还在处理被停止的 partition A 时WatchForKafkaTopicPartitions又发射了重新添加的 A导致同一 partition 被并发处理、输出重复记录。根本原因是两个 DoFn 对 stop/remove 信号的反应无法保证在同一时刻发生。六、ReadSourceDescriptors在运行时确定读取描述KafkaIO.readSourceDescriptors()内部即ReadSourceDescriptors是另一条读取路径它以一个PCollectionKafkaSourceDescriptor为输入、输出PCollectionKafkaRecord核心实现基于Splittable DoFnSDF。它与KafkaIO.Read最大的区别在于不需要在管道构建期提供 topic 列表、topic pattern、topic partitions、start/stop 时间等 source 描述这些描述可以在运行时才由管道动态填充。典型场景是先从 BigQuery 表查询出要消费的 Kafka topic 列表再交给ReadSourceDescriptors去读取。基础用法构造期只配置反序列化器pipeline .apply(Create.of(KafkaSourceDescriptor.of(new TopicPartition(topic, 1))) .apply(KafkaIO.readAll() .withBootstrapServers(broker_1:9092,broker_2:9092) .withKeyDeserializer(LongDeserializer.class) .withValueDeserializer(StringDeserializer.class));bootstrapServers也可以由KafkaSourceDescriptor自身携带pipeline .apply(Create.of( KafkaSourceDescriptor.of( new TopicPartition(topic, 1), null, null, ImmutableList.of(broker_1:9092, broker_2:9092)) .apply(KafkaIO.readAll() .withKeyDeserializer(LongDeserializer.class) .withValueDeserializer(StringDeserializer.class));注意若KafkaSourceDescriptor不提供 bootstrap servers 且withBootstrapServers未配置expand()会打印 WARN 日志提示必须在运行时通过 descriptor 填充否则管道会失败。ReadSourceDescriptors的消费侧配置与KafkaIO.Read基本一致consumer config、consumer factory fn、offset consumer config、key/value coder、key/value deserializer provider、isCommitOffsetEnabled对应commitOffsetsInFinalize另有几个处理记录相关的专属配置commitOffsets()处理完记录后提交 offset。注意如果 consumer config 里同时设置了enable.auto.committruecommitOffsets()会被忽略源码expand()中会打 INFO 日志提示二者只需其一withExtractOutputTimestampFn(fn)为每条KafkaRecord计算输出时间戳并控制水印推进内置三种类型withProcessingTime()、withCreateTime()、withLogAppendTime()对应实现见内部类ExtractOutputTimestampFns——useCreateTime/useLogAppendTime均校验记录的 timestampType不匹配会抛IllegalArgumentException水印估计器三选一withWallTimeWatermarkEstimator()、withMonotonicallyIncreasingWatermarkEstimator()默认、withManualWatermarkEstimator()。组合示例pipeline .apply(Create.of( KafkaSourceDescriptor.of( new TopicPartition(topic, 1), null, null, ImmutableList.of(broker_1:9092, broker_2:9092)) .apply(KafkaIO.readAll() .withKeyDeserializer(LongDeserializer.class) .withValueDeserializer(StringDeserializer.class) .withProcessingTime() .commitOffsets());从源码的展开逻辑expand()看当commitOffsets()开启时变换会展开为ParDo(ReadFromKafkaDoFn) → Reshuffle → Map(输出 KafkaRecord) → KafkaCommitOffset的复合结构该展开在 Dataflow 上跑 cross-language 时不受支持。七、Confluent Schema Registry按 Schema 反序列化 Avro如果你希望按 Confluent Schema Registry 中的 schema 反序列化 key/valueKafkaIO 可以从指定的 Schema Registry URL 拉取 schema 并用于反序列化且会依据对应的Deserializer自动推断 Coder。对于 Avro schema输出将是PCollectionKafkaRecord其中 key/value 类型为org.apache.avro.generic.GenericRecord此时无需指定 key/value 反序列化器和 coder因为会默认设置为KafkaAvroDeserializer与AvroCoder。示例value 使用 Schema Registry 中存储的 Avro schemakey 为 LongPCollectionKafkaRecordLong, GenericRecord input pipeline .apply(KafkaIO.Long, GenericRecordread() .withBootstrapServers(broker_1:9092,broker_2:9092) .withTopic(my_topic) .withKeyDeserializer(LongDeserializer.class) // 指定 Schema Registry URL 与 value subject .withValueDeserializer( ConfluentSchemaRegistryDeserializerProvider.of(http://localhost:8081, my_topic-value)) ...还可以向 schema registry client 传递认证等属性ImmutableMapString, Object csrConfig ImmutableMap.String, Objectbuilder() .put(AbstractKafkaAvroSerDeConfig.BASIC_AUTH_CREDENTIALS_SOURCE,USER_INFO) .put(AbstractKafkaAvroSerDeConfig.USER_INFO_CONFIG,username:password) .build(); PCollectionKafkaRecordLong, GenericRecord input pipeline .apply(KafkaIO.Long, GenericRecordread() .withBootstrapServers(broker_1:9092,broker_2:9092) .withTopic(my_topic) .withKeyDeserializer(LongDeserializer.class) .withValueDeserializer( ConfluentSchemaRegistryDeserializerProvider.of(https://localhost:8081, my_topic-value, null, csrConfig)) ...对应的实现类是 ConfluentSchemaRegistryDeserializerProvider.java其测试用例见 ConfluentSchemaRegistryDeserializerProviderTest.java。八、写入 KafkaKafkaIO.write() 系列KafkaIO 的 sink 支持把 key-value 对写入 Kafka topic也可以只写 value或直接写原生ProducerRecord。配置一个 Kafka sink 至少需要bootstrapServers、目标 topic、key 和 value 序列化器。8.1 标准 KV 写入PCollectionKVLong, String kvColl ...; kvColl.apply(KafkaIO.Long, Stringwrite() .withBootstrapServers(broker_1:9092,broker_2:9092) .withTopic(results) .withKeySerializer(LongSerializer.class) .withValueSerializer(StringSerializer.class) // 通过 ProducerConfig 进一步定制 KafkaProducer例如开启压缩 .withProducerConfigUpdates(ImmutableMap.of(compression.type, gzip)) // 为发布的记录设置时间戳使用元素时间戳 .withInputTimestamp() // 或自定义时间戳函数 .withPublishTimestampFunction((elem, elemTs) - ...) // 序列化失败的记录可送入错误处理器 .withBadRecordErrorHandler(errorHandler) // 可选开启精确一次写入需受支持的 runner .withEOS(20, eos-sink-group-id); );8.2 只写 value空 keyPCollectionString strings ...; strings.apply(KafkaIO.Void, Stringwrite() .withBootstrapServers(broker_1:9092,broker_2:9092) .withTopic(results) .withValueSerializer(StringSerializer.class) // 只需要 value 序列化器 .values() );8.3 写入原生 ProducerRecord如需写入ProducerRecord可逐条指定 topic、partition、时间戳等使用KafkaIO.writeRecords()PCollectionProducerRecordLong, String records ...; records.apply(KafkaIO.Long, StringwriteRecords() .withBootstrapServers(broker_1:9092,broker_2:9092) .withTopic(results) .withKeySerializer(LongSerializer.class) .withValueSerializer(StringSerializer.class) );WriteRecords.withTopic(...)设置的只是默认 topicProducerRecord可以逐条覆盖。若写入时设置了 keyKafka 会用 key 的 hash 决定 partition。write()内部实际就是包装了一个writeRecords()transform见 KafkaIO.write() 实现。8.4 写入 Avro value要产出 Avro value可使用io.confluent.kafka.serializers.KafkaAvroSerializer。为了让该类的泛型类型在KafkaIO.write()/withValueSerializer()中正常工作需要通过 (Class) 强转擦除泛型KafkaIO.Long, Stringwrite() ... .withValueSerializer((Class)KafkaAvroSerializer.class) .withProducerConfigUpdates( Map with schema registry configuration details ) ...8.5 发布时间戳的两个注意事项withInputTimestamp()等价于withPublishTimestampFunction((e, ts) - ts)即用管道元素时间戳作为发布到 Kafka 的消息时间戳Kafka 的保留策略基于消息时间戳如果管道在处理过去时间的数据而发布的时间戳早于 Kafka 集群的log.retention.hours消息可能在被发布后立即被 Kafka 删除。8.6 精确一次写入Exactly-once SinkwithEOS(int numShards, String sinkGroupId)为写入 Kafka 提供精确一次语义使应用在 Beam 管道内部的精确一次之上获得端到端精确一次保证即使管道执行中发生重试worker 重启、数据重分配/自动扩容等写入 sink 的记录也只在 Kafka 上提交一次。实现原理Beam runner 通常只为管道结果提供精确一次语义而对用户代码在 transform 中产生的外部副作用如写 Kafka不提供保证——若不加处理这类写入可能发生多次。开启 EOS 后sink 将兼容 runner 的 checkpoint 语义与 Kafka0.11的事务机制绑定记录先进入事务与 checkpoint 一起原子提交从而保证只写一次。该能力依赖 runner 的 checkpoint 语义并非所有 runner 兼容实现会在初始化时检查不允许的 runner 会直接抛异常。JavaDoc 明确说明 Dataflow、Flink、Spark runner 兼容。参数说明numShards设置 sink 并行度。状态元数据以sinkGroupId为前缀分布在 Kafka 上的这么多虚拟 partition 中经验法则是设为接近 Kafka topic 的 partition 数sinkGroupId用于在 Kafka 上存储少量状态元数据的 group id类似KafkaConsumer的 consumer group id。每个 job 应使用唯一的 group id这样 job 重启/更新时能保留状态以维持精确一次语义。该状态与 sink 事务在 Kafka 上原子提交对应KafkaProducer.sendOffsetsToTransaction(Map, String)。sink 在初始化时会做多项 sanity check避免误用看似由同一 job 写入的状态。性能提示JavaDoc 原文精确一次 sink 涉及两次记录 shuffle且记录要经过两轮序列化-反序列化。如果数据量大、序列化成本高CPU 开销会比较明显可以通过在写入 Kafka sink 前先把数据序列化为 byte array 来降低 CPU 成本。sink的默认 producer 配置与默认 consumer 配置类似地由常量维护withProducerConfigUpdates(Map)只更新相同 key 的值、不会整体覆盖默认配置withProducerFactoryFn(...)可自定义 producer 工厂主要用于测试默认为KafkaProducer。九、事件时间戳与水印TimestampPolicy默认情况下KafkaIO reader 会把记录时间戳事件时间设为处理时间source 水印推进为当前墙钟时间。如果 topic 启用了 Kafka 服务端摄入时间戳LogAppendTime可以调用withLogAppendTime()启用也可以实现TimestampPolicyFactory提供完全自定义的策略并通过withTimestampPolicyFactory(...)接入。三种内置策略ProcessingTimePolicy默认每条记录的时间戳为它在 reader 中成为当前记录的时刻水印始终推进到当前时间。若 Kafka 服务端时间LogAppendTime可用JavaDoc 建议优先使用withLogAppendTime()而非本策略。LogAppendTimePolicywithLogAppendTime()为每条记录赋 Kafka 的 log append time服务端摄入时间。每个 Kafka partition 的水印是最后一条已读记录的时间戳当 partition 空闲时水印推进到墙钟时间之后数秒。要求每条从 Kafka 消费的记录的时间戳类型都是LOG_APPEND_TIME——Kafka 中该模式需在 topic 上启用启用后所有后续记录的时间戳都是 log append time若个别记录的时间戳类型不是LOG_APPEND_TIME则取其时间戳为前一条记录的时间戳或最新水印中较大者。整个 source 的水印是各 partition 水印中最小的那个——如果某个 reader 因 partition 间记录分布不均而落后它会拖住整个 source 的水印。CreateTimePolicywithCreateTime(Duration maxDelay)基于记录的CREATE_TIME时间戳。如果某条记录的时间戳类型不是CREATE_TIME则报错。partition 内时间戳要求大致单调递增允许的乱序上限为maxDelay例如 1 分钟。任一时刻水印为Min(now, 当前已见最大事件时间) - maxDelay且水印不会设置为未来时间、上限为now - maxDelaypartition 空闲时水印推进到now - maxDelay。maxDelay的含义是partition 中任意一条记录之后的时间戳预期不早于当前记录时间戳 - maxDelay。自定义策略TimestampPolicyFactory.createTimestampPolicy(TopicPartition, Optional)会在 reader 启动时针对每个 partition 被调用一次返回该 partition 的时间戳/水印策略。注withTimestampFn2/withWatermarkFn2/withTimestampFn/withWatermarkFn四个方法自 2.4 起已弃用应统一改用withTimestampPolicyFactory(...)。十、读侧的其它重要开关withReadCommitted()在 consumer 配置中设置isolation.levelread_committed确保不读取未提交消息。Kafka 0.11 引入事务写要求端到端精确一次语义的应用应只读已提交消息详见KafkaConsumerJavaDoc。commitOffsetsInFinalize()把已 finalize 的 checkpoint 对应 offset 提交回 Kafka对应CheckpointMark.finalizeCheckpoint()。它有助于管道重启时减少记录缺口或重复处理但不提供硬性的处理保证调用finalizeCheckpoint()后提交可能有短暂延迟reader 可能正阻塞在读取上。它与 Kafka consumer 的AUTO_COMMIT配置相互独立——通常启用其中一个即可不要同时启用。注意启用本方法要求 consumer 配置里设置group.id否则expand()直接报错若同时设置了enable.auto.committrue会打出 WARN 提示二者只需其一。withMaxNumRecords(long)/withMaxReadTime(Duration)与Read.Unbounded类似主要用于测试和演示应用可将无界读取转为有界读取。withBadRecordErrorHandler(ErrorHandler)把解析失败的坏消息发送到备选 sink见ErrorHandler文档。注意基于旧版 UnboundedSource 的实现不支持坏记录错误处理器会打 WARN 提示改用 SDF 实现坏记录路由仅在 SDF 实现ReadSourceDescriptors中生效。withConsumerFactoryFn(...)自定义创建Consumer的工厂函数默认KafkaConsumer可用于对接其它 Kafka consumer 实现/版本。withCheckStopReadingFn(...)运行时判定是否停止读取某TopicPartition的可序列化函数见第五节动态读取。十一、SDF 与 UnboundedSourceKafkaIO 的两种读实现从 KafkaIO.java 的 expand() 可以看到读侧在展开时会在两种实现间选择ReadFromKafkaViaUnboundedLEGACY基于UnboundedSource的旧实现通过org.apache.beam.sdk.io.Read.from(...)包装若设置了maxNumRecords或maxReadTime会转为有界读取。坏记录错误处理器不受支持ReadFromKafkaViaSDFSDF基于 Splittable DoFn 的新实现内部复用ReadSourceDescriptors支持动态读取、坏记录错误处理器、stop 时间戳有界化等新特性。选择规则源码可验证当管道设置了beam_fn_api_use_deprecated_read、use_deprecated_read、use_unbounded_sdf_wrapper等实验开关或 runner 兼容性判定只支持 LEGACY 实现、或 runner 更倾向旧实现时走 LEGACY否则默认走 SDF。其中runnerPrefersLegacyRead(...)的逻辑是Dataflow runner 与开启beam_fn_api的管道倾向 SDF其余非 portable runner 出于性能考量默认用旧实现除非显式开启use_sdf_read。兼容性判定逻辑集中在 KafkaIOReadImplementationCompatibility.java配套测试见 KafkaIOReadImplementationCompatibilityTest.java。十二、Protobuf 测试资源重建模块的测试资源目录中包含一个 protobuf descriptor setproto_byte_utils.pb用于测试自定义序列化。当需要重新生成该 descriptor set 时可按 README 中的说明使用protoc命令protoc \ -Isdks/java/io/kafka/src/test/resources/ \ --descriptor_set_outsdks/java/io/kafka/src/test/resources/proto_byte/file_descriptor/proto_byte_utils.pb \ sdks/java/io/kafka/src/test/resources/proto_byte/proto_byte_utils.proto命令中的路径以仓库根目录为基准。十三、测试与验证如何在仓库中进一步探索本文涉及的诸多行为都有对应的单元/集成测试可佐证便于读者深入理解KafkaIOTest.java主测试套件覆盖读/写主流程KafkaIOIT.java面向真实 Kafka 集群的集成测试KafkaReadSchemaTransformProviderTest.java 与 KafkaWriteSchemaTransformProviderTest.javaSchema Transform 的读/写提供者测试KafkaDlqTest.java坏记录/死信队列DLQ行为验证WatchForKafkaTopicPartitionsTest.java动态 topic/partition 发现测试KafkaRecordCoderTest.java、TopicPartitionCoderTest.java、ProducerRecordCoderTest.java各类 Coder 的编码/解码测试。总结Apache Beam KafkaIO 以beam-sdks-java-io-kafka 自带的kafka-clients运行时依赖为底座提供了从最简单主题读取到动态 partition 发现、从默认处理时间戳到LogAppendTime/CreateTime/自定义时间戳策略、从普通 KV 写入到基于 Kafka 事务的精确一次写入的完整能力矩阵。理解其双 consumer 架构主消费 offset 估算、三选一的 topic 指定方式、checkpoint 优先的初始 offset 确定顺序以及SDF 与 UnboundedSource 两种读实现的选择逻辑是正确配置生产级管道、排查重复/漏读与时序问题的关键。更多权威细节建议直接查阅 KafkaIO.java 的完整 JavaDoc——它本身就是该模块最完整的使用手册。赞分享【免费下载链接】beamApache Beam is a unified programming model for Batch and Streaming data processing.项目地址https://gitcode.com/gh_mirrors/beam18/beam点击查看免费下载相关推荐Apache Beam KafkaIO 实战指南Java 读写 Apache Kafka 的完整配置与源码级原理Apache Beam KafkaIO 实战指南Java 读写 Apache Kafka 的完整配置与源码级原理 Apache Beam 的 KafkaIO批处理流处理大数据Apache Beam KafkaIO 实战指南用 Java SDK 读写 Kafka 消息Apache Beam KafkaIO 实战指南用 Java SDK 读写 Kafka 消息 Apache Beam 的 KafkaIO 是统一批流编程模型中大数据批处理流处理数据工程Apache Beam 中 KafkaIO 连接器从 Kafka Topic 读取与写入的完整实践指南Apache Beam 中 KafkaIO 连接器从 Kafka Topic 读取与写入的完整实践指南 导读 Apache Beam 为 Apache Kaf大数据批处理流处理数据工程上一篇MongoDB数据验证规则版本管理Robo 3T追踪约束变更的完整指南下一篇kohya_ss模型压缩技术SVDF分解降低模型体积70%的终极指南创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
返回列表