Flink批处理模式下Kafka偏移量未正确提交问题求助
问题描述
我正在开发一个每日将Kafka数据导出至S3的数据管道,每日数据量极少(不足100万条记录,单条大小1KB),计划每日运行一次管道,从上次提交的偏移量消费至最新偏移量,写入S3的Parquet文件。但在Flink批处理模式(RuntimeExecutionMode.BATCH)下运行bounded DataStream任务时,出现偏移量提交异常:有时仅单个分区偏移量被提交,有时全部分区都未提交,但日志显示所有分区偏移量均已成功提交。
环境信息
Kafka Topic Partitions: 3 Kafka Topic Replication: 1 Java: 11 Flink: 1.17
相关日志
org.apache.kafka.clients.consumer.internals.ConsumerCoordinator [] - [Consumer clientId=flink-batch-test-2-1, groupId=flink-batch-test-2] Sending asynchronous auto-commit of offsets {my-topic-name-1=OffsetAndMetadata{offset=7907, leaderEpoch=6, metadata=''}} org.apache.kafka.clients.consumer.internals.ConsumerCoordinator [] - [Consumer clientId=flink-batch-test-2-0, groupId=flink-batch-test-2] Sending asynchronous auto-commit of offsets {my-topic-name-0=OffsetAndMetadata{offset=45198, leaderEpoch=4, metadata=''}} org.apache.kafka.clients.consumer.internals.ConsumerCoordinator [] - [Consumer clientId=flink-batch-test-2-2, groupId=flink-batch-test-2] Sending asynchronous auto-commit of offsets {my-topic-name-2=OffsetAndMetadata{offset=7791, leaderEpoch=2, metadata=''}} org.apache.kafka.clients.consumer.internals.Fetcher [] - [Consumer clientId=flink-batch-test-2-0, groupId=flink-batch-test-2] Added READ_UNCOMMITTED fetch request for partition my-topic-name-0 at position FetchPosition{offset=45198, offsetEpoch=Optional[4], currentLeader=LeaderAndEpoch{leader=Optional[prefix2.mycluster.com:9092 (id: 1003 rack: null)], epoch=4}} to node prefix2.mycluster.com:9092 (id: 1003 rack: null) org.apache.kafka.clients.consumer.internals.Fetcher [] - [Consumer clientId=flink-batch-test-2-1, groupId=flink-batch-test-2] Added READ_UNCOMMITTED fetch request for partition my-topic-name-1 at position FetchPosition{offset=7907, offsetEpoch=Optional[6], currentLeader=LeaderAndEpoch{leader=Optional[prefix3.mycluster.com:9092 (id: 1004 rack: null)], epoch=6}} to node prefix3.mycluster.com:9092 (id: 1004 rack: null) org.apache.kafka.clients.consumer.internals.Fetcher [] - [Consumer clientId=flink-batch-test-2-2, groupId=flink-batch-test-2] Added READ_UNCOMMITTED fetch request for partition my-topic-name-2 at position FetchPosition{offset=7791, offsetEpoch=Optional[2], currentLeader=LeaderAndEpoch{leader=Optional[prefix1.mycluster.com:9092 (id: 1002 rack: null)], epoch=2}} to node prefix1.mycluster.com:9092 (id: 1002 rack: null) org.apache.kafka.clients.consumer.internals.ConsumerCoordinator [] - [Consumer clientId=flink-batch-test-2-2, groupId=flink-batch-test-2] Committed offset 7791 for partition my-topic-name-2 org.apache.kafka.clients.consumer.internals.ConsumerCoordinator [] - [Consumer clientId=flink-batch-test-2-1, groupId=flink-batch-test-2] Committed offset 7907 for partition my-topic-name-1 org.apache.kafka.clients.consumer.internals.ConsumerCoordinator [] - [Consumer clientId=flink-batch-test-2-0, groupId=flink-batch-test-2] Committed offset 45198 for partition my-topic-name-0
代码实现
package org.example; import org.slf4j.Logger; import java.util.Properties; import org.slf4j.LoggerFactory; import org.apache.kafka.clients.consumer.ConsumerConfig; import org.apache.kafka.clients.consumer.OffsetResetStrategy; import org.apache.kafka.common.serialization.StringDeserializer; import org.apache.flink.api.common.RuntimeExecutionMode; import org.apache.flink.api.common.functions.MapFunction; import org.apache.flink.api.common.eventtime.WatermarkStrategy; import org.apache.flink.streaming.api.datastream.DataStream; import org.apache.flink.streaming.api.functions.sink.PrintSinkFunction; import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment; import org.apache.flink.connector.kafka.source.KafkaSource; import org.apache.flink.connector.kafka.source.enumerator.initializer.OffsetsInitializer; import org.apache.flink.connector.kafka.source.reader.deserializer.KafkaRecordDeserializationSchema; public class Main { public static void main(String[] args) throws Exception { final Logger logger = LoggerFactory.getLogger(Main.class); final StreamExecutionEnvironment environment = StreamExecutionEnvironment.getExecutionEnvironment(); environment.setRuntimeMode(RuntimeExecutionMode.BATCH); System.out.println(environment); DataStream<String> srcStream = environment.fromSource(getKafkaSource(), WatermarkStrategy.noWatermarks(), "Kafka Source") .name("kafka source") .uid("kafka source") .setParallelism(1); srcStream.map(new MapFunction<String, String>() { @Override public String map(String data) throws Exception { return data; } }).name("Transformation"); srcStream.addSink(new PrintSinkFunction<>()).name("Print Sink"); environment.execute("test"); environment.close(); } public static KafkaSource<String> getKafkaSource(){ Properties prop = new Properties(); prop.setProperty(ConsumerConfig.ENABLE_AUTO_COMMIT_CONFIG, "true"); return KafkaSource.<String>builder() .setBootstrapServers("prefix1.mycluster.com:9092,prefix2.mycluster.com:9092,prefix3.mycluster.com:9092,prefix4.mycluster.com:9092") .setTopics("my-topic-name") .setGroupId("flink-batch-test-2") .setDeserializer(KafkaRecordDeserializationSchema.valueOnly(StringDeserializer.class)) .setStartingOffsets(OffsetsInitializer.committedOffsets(OffsetResetStrategy.EARLIEST)) .setProperties(prop) .setBounded(OffsetsInitializer.latest()) .build(); } }
解决方案
1. 禁用Kafka自动提交,改用Flink管理偏移量
批处理模式下,Kafka异步自动提交机制与Flink批处理生命周期不兼容。Flink任务结束时会直接关闭资源,可能导致Kafka的异步提交请求未完成就被中断,这就是日志显示提交成功但实际未生效的核心原因。
修改getKafkaSource()方法中的配置:
Properties prop = new Properties(); // 禁用自动提交,交给Flink管理 prop.setProperty(ConsumerConfig.ENABLE_AUTO_COMMIT_CONFIG, "false");
2. 调整Kafka Source并行度匹配分区数
当前代码中Source并行度设为1,单实例处理3个分区会增加偏移量提交的不确定性。将并行度改为与Topic分区数一致(3),让每个分区由独立的Source实例处理:
DataStream<String> srcStream = environment.fromSource(getKafkaSource(), WatermarkStrategy.noWatermarks(), "Kafka Source") .name("kafka source") .uid("kafka source") .setParallelism(3); // 匹配Topic分区数
3. 移除手动关闭环境的代码
Flink任务在environment.execute()完成后会自动关闭环境,手动调用environment.close()可能提前终止资源清理流程,影响偏移量提交逻辑,直接删除该行代码。
修改后的核心代码片段
Kafka Source配置:
public static KafkaSource<String> getKafkaSource(){ Properties prop = new Properties(); // 禁用Kafka自动提交 prop.setProperty(ConsumerConfig.ENABLE_AUTO_COMMIT_CONFIG, "false"); return KafkaSource.<String>builder() .setBootstrapServers("prefix1.mycluster.com:9092,prefix2.mycluster.com:9092,prefix3.mycluster.com:9092,prefix4.mycluster.com:9092") .setTopics("my-topic-name") .setGroupId("flink-batch-test-2") .setDeserializer(KafkaRecordDeserializationSchema.valueOnly(StringDeserializer.class)) .setStartingOffsets(OffsetsInitializer.committedOffsets(OffsetResetStrategy.EARLIEST)) .setProperties(prop) .setBounded(OffsetsInitializer.latest()) .build(); }
主函数调整:
public static void main(String[] args) throws Exception { final Logger logger = LoggerFactory.getLogger(Main.class); final StreamExecutionEnvironment environment = StreamExecutionEnvironment.getExecutionEnvironment(); environment.setRuntimeMode(RuntimeExecutionMode.BATCH); DataStream<String> srcStream = environment.fromSource(getKafkaSource(), WatermarkStrategy.noWatermarks(), "Kafka Source") .name("kafka source") .uid("kafka source") .setParallelism(3); // 匹配分区数 srcStream.map(new MapFunction<String, String>() { @Override public String map(String data) throws Exception { return data; } }).name("Transformation"); srcStream.addSink(new PrintSinkFunction<>()).name("Print Sink"); // 执行任务,无需手动关闭环境 environment.execute("test"); }
验证方法
任务完成后,使用Kafka命令行工具检查偏移量提交状态:
kafka-consumer-groups.sh --bootstrap-server prefix1.mycluster.com:9092 --describe --group flink-batch-test-2
确认所有分区的CURRENT-OFFSET与日志中显示的提交偏移量一致即可。
内容的提问来源于stack exchange,提问作者Mrugesh Kadia
相关产品推荐
相关产品推荐

