You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.07.21 23:34:58