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

Flink Kafka Connector配置EXACTLY_ONCE仍产生重复消息求助

问题描述

版本信息

  • flink.version 1.15.2
  • scala.binary.version 2.12
  • java.version 1.11

代码

public static void main(String[] args) throws Exception { 
    StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
    Properties prop = new Properties();
    prop.put("commit.offsets.on.checkpoint", "true"); 
    KafkaSource<String> source = KafkaSource.<String>builder()
        .setBootstrapServers(kafkaBroker)
        .setTopics("inputTopic")
        .setGroupId("my-group"+ System.currentTimeMillis())
        .setStartingOffsets(OffsetsInitializer.committedOffsets(OffsetResetStrategy.EARLIEST))
        .setValueOnlyDeserializer(new SimpleStringSchema())
        .setProperties(prop)
        .build(); 
    DataStream<String> sourceStream = env.fromSource(
        source,
        WatermarkStrategy.noWatermarks(), "Kafka Source");  
    Properties sinkProps = new Properties(); 
    sinkProps.put(ProducerConfig.TRANSACTION_TIMEOUT_CONFIG, 6000); 
    KafkaSink<String> sink = KafkaSink.<String>builder()
        .setBootstrapServers(kafkaBroker)
        .setKafkaProducerConfig(sinkProps)
        .setRecordSerializer(new OffsetSerializer("outputTopic"))
        .setDeliverGuarantee(DeliveryGuarantee.EXACTLY_ONCE)
        .setTransactionalIdPrefix("trx-"+System.currentTimeMillis())
        .build();  
    sourceStream
        .keyBy(new KeySelector<String, String>(){ 
            @Override
            public String getKey(String value) throws Exception {
                System.out.println("Key >>" + value); 
                return value;
            } 
        })
        .map(new MapFunction<String, String>() {
            @Override public String map(String value) throws Exception {  
                System.out.println("offset >>" + value);
                if(value.equalsIgnoreCase("4"))
                {  
                    System.out.println("Custom unhandled exception");
                    throw new Exception("Custom unhandled exception");
                } 
                return value; 
            }
        }) 
        .sinkTo(sink);
    env.enableCheckpointing(500);
    env.getCheckpointConfig().setCheckpointTimeout(10000);
    env.getCheckpointConfig().setCheckpointingMode(CheckpointingMode.EXACTLY_ONCE);
    env.getCheckpointConfig().setMinPauseBetweenCheckpoints(10L);
    env.getCheckpointConfig().setTolerableCheckpointFailureNumber(1);
    env.getCheckpointConfig().setExternalizedCheckpointCleanup(
        ExternalizedCheckpointCleanup.NO_EXTERNALIZED_CHECKPOINTS);
    env.setRestartStrategy(RestartStrategies.fixedDelayRestart(1
        ,org.apache.flink.api.common.time.Time.of(1,TimeUnit.SECONDS)));
    env.execute("tester");
}  

已尝试操作

  • 检查Kafka集群事务是否启用
  • 修改TRANSACTION_TIMEOUT_CONFIG
  • 每次运行修改消费者组ID
  • 移除DataStream的keyBy操作
  • 添加以下配置:
    env.getCheckpointConfig().setMaxConcurrentCheckpoints(1);
    env.getCheckpointConfig().enableUnalignedCheckpoints();
    
  • 使用不同类型的ExternalizedCheckpointCleanup

诉求

期望每条消费的消息仅向Kafka生产一条消息,但发生异常时,之前已成功处理的消息会被重新处理,寻求解决方案。


解决方案

1. 异常本地化处理,避免任务重启

直接抛出未捕获异常会触发Flink重启机制,重启后任务会从上一次成功的检查点位置重新消费数据,这就导致已处理的消息被重复处理。可以在map算子内捕获异常,将异常消息路由到死信队列,而非中断整个任务:

@Override public String map(String value) throws Exception {  
    System.out.println("offset >>" + value);
    if(value.equalsIgnoreCase("4"))
    {  
        System.out.println("Custom unhandled exception");
        // 将异常消息发送到死信队列单独处理
        deadLetterSink.send(value);
        // 返回null或特殊标识,跳过后续生产逻辑
        return null;
    } 
    return value; 
}

2. 修正事务超时与检查点的时间匹配

当前TRANSACTION_TIMEOUT_CONFIG=6000(6秒)小于检查点超时时间(10秒),这会导致Kafka事务在检查点完成前超时,触发事务回滚,重启后重复发送消息。需调整事务超时时间大于检查点间隔+检查点超时时间,同时确保Kafka集群的transaction.max.timeout.ms(默认900秒)大于该值:

sinkProps.put(ProducerConfig.TRANSACTION_TIMEOUT_CONFIG, 15000);

3. 固定消费者组ID与事务ID前缀

每次运行动态生成这两个ID,会导致Flink无法复用之前的检查点和已提交的偏移量,重启后只能从最早偏移量开始消费。固定ID后,任务重启可以直接从最近的检查点恢复状态:

.setGroupId("my-group-fixed")
// ...
.setTransactionalIdPrefix("trx-fixed")

4. 启用外部化检查点保留

将外部化检查点清理策略设置为RETAIN_ON_CANCELLATION,这样任务失败或重启时,检查点会被保留,恢复时直接从检查点继续处理,避免从头消费:

env.getCheckpointConfig().setExternalizedCheckpointCleanup(
    ExternalizedCheckpointCleanup.RETAIN_ON_CANCELLATION);

内容的提问来源于stack exchange,提问作者Kanad Mehta

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.18 02:35:34