寻求支持Spark Runner与Kafka批量偏移提交的Apache Beam适配版本
提问(翻译自英文)
我当前使用Apache Beam 2.29.0版本开发项目,希望升级到支持**批量处理结束时提交Kafka偏移量(至少一次语义)**的版本。尝试升级到2.64.0版本后,Spark Runner无法完全兼容,出现如下异常:
onWindowExpiration is not supported: org.apache.beam.sdk.transforms.GroupIntoBatches$GroupIntoBatchesDoFn
我的代码实现如下:
Pipeline p = Pipeline.create(options); p.apply(KafkaIO.<String, String>read() .withBootstrapServers("localhost:9092") .withTopic("data_input") .withConsumerConfigUpdates(ImmutableMap.of("enable.auto.commit", "false")) // 禁用自动提交 .commitOffsetsInFinalize() // 批量作业中使用手动提交 ) .apply(ParDo.of(new DoFn<KV<String, String>, Void>() { @ProcessElement public void processElement(ProcessContext c) { // 处理每个元素 } })) .apply(new CheckpointFn()); // 自定义检查点转换 p.run().waitUntilFinish();
解决方案
适配版本推荐
优先选择Apache Beam 2.48.0 ~ 2.59.0区间内的版本,理由如下:
- 该区间已稳定支持
commitOffsetsInFinalize()批量提交Kafka偏移量的特性(该特性在Beam 2.34.0左右正式落地) - 未引入导致Spark Runner兼容问题的
GroupIntoBatches窗口过期回调逻辑(该逻辑在Beam 2.60.0后新增,Spark Runner未及时适配)
其中2.59.0是最优选择,它是Spark Runner兼容性断裂前的最后几个稳定版本,同时包含批量偏移量提交的完整支持。
异常原因分析
你遇到的onWindowExpiration is not supported错误,源于Beam 2.60.0后对GroupIntoBatches变换的内部实现调整:新增了窗口过期回调逻辑,但Spark Runner的适配进度未跟上,导致该DoFn无法在Spark环境中正常运行。而2.59.0及更早版本的GroupIntoBatches无此逻辑,可与Spark Runner完美兼容。
代码优化提示
- 移除自定义的
CheckpointFn中与偏移量提交相关的逻辑,避免和commitOffsetsInFinalize()的批量提交逻辑冲突——批量模式下,commitOffsetsInFinalize()会在管道完全结束时统一提交所有消费的偏移量。 - 保持
enable.auto.commit为false,防止Kafka自动提交与Beam手动提交的逻辑冲突。
内容的提问来源于stack exchange,提问作者Fabio
相关产品推荐
相关产品推荐

