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
相关产品推荐
相关产品推荐

