如何避免主Topic重复消费已被DLQ处理的消息?
Kafka主Topic消息推DLQ后重启重复消费问题解决
问题描述
我维护的Kafka消费者类包含主Topic监听器与DLQ监听器,主Topic监听器处理消息失败时,消息会通过Bean工厂推送到DLQ Topic且能被成功处理,但重启消费者应用后,已被DLQ处理的消息会被主Topic监听器重复消费。
相关代码
KafkaConsumer.java
public class KafkaConsumer { //MAIN TOPIC LISTENER @KafkaListener(id = "main-topic", topics = "main-topic", groupId = "main", containerFactory = "kafkaListenerContainerFactory", clientIdPrefix = "main-topic") public void mainListener(ConsumerRecord<String, String> consumerRecord, Acknowledgement ack) { // 将consumerRecord值转换为Object类并保存到DB dbService.saveTodb(consumerRecord.value(), new ObjectMapper()); ack.acknowledge(); } //DLQ LISTENER @KafkaListener(id = "DLQ-topic", topics = "DLQ-topic", groupId = "main", clientIdPrefix = "DLQ", autostartup= "false") public void dlqListener(ConsumerRecord<String, String> consumerRecord, Acknowledgement ack) { dbService.saveTodb(consumerRecord.value(), new ObjectMapper()); ack.acknowledge(); } }
注:原代码中DLQ监听器方法名与主Topic监听器重复,已修正为dlqListener避免歧义
KafkaBeanFactory.java
public class KafkaBeanFactory{ @Bean(name = "kafkaListenerContainerFactory") public ConcurrentKafkaListenerContainerFactory<Object, Object> kafkaListenerContainerFactory( ConcurrentKafkaListenerContainerFactoryConfigurer configurer, ConsumerFactory<Object, Object> kafkaConsumerFactory, KafkaTemplate<Object, Object> template) { ConcurrentKafkaListenerContainerFactory<Object, Object> factory = new ConcurrentKafkaListenerContainerFactory<>(); configurer.configure(factory, kafkaConsumerFactory); var recoverer = new DeadLetterPublishingRecoverer(template, (record, ex) -> new TopicPartition("DLQ-topic", record.partition())); var errorHandler = new DefaultErrorHandler(recoverer, new FixedBackOff(3, 20000)); errorHandler.addRetryableExceptions(JsonProcessingException.class, DBException.class); errorHandler.setAckAfterHandle(true); factory.setCommonErrorHandler(errorHandler); return factory; } }
application.yaml
kafka: bootstrap-servers: localhost:9092 # 示例值 client-id: main&DLQ properties: security: protocol: SASL_SSL sasl: mechanism: PLAIN jaas: config: org.apache.kafka.common.security.plain.PlainLoginModule required username="$ConnectionString" password="<string>"; security: protocol: SASL_SSL consumer: key-deserializer: org.apache.kafka.common.serialization.StringDeserializer value-deserializer: org.apache.kafka.common.serialization.StringDeserializer groupId: main enable-auto-commit: false auto.offset.reset: earliest listener: ack-mode: MANUAL_IMMEDIATE
问题原因
核心原因是主Topic中处理失败并推送到DLQ的消息,其偏移量未被正确提交。由于启用了MANUAL_IMMEDIATE手动确认模式,监听器仅在处理成功时调用ack.acknowledge(),处理失败时该逻辑不会执行;而当前配置的DefaultErrorHandler未明确指定在恢复(推送DLQ)后提交偏移量,导致消费组偏移量停留在失败消息的位置,重启后会重新拉取该消息。
解决方案
1. 配置错误处理器在恢复后提交偏移量
修改KafkaBeanFactory中的DefaultErrorHandler配置,添加setCommitRecovered(true),确保消息推送到DLQ后,主Topic的偏移量被提交:
var errorHandler = new DefaultErrorHandler(recoverer, new FixedBackOff(3, 20000)); errorHandler.addRetryableExceptions(JsonProcessingException.class, DBException.class); errorHandler.setAckAfterHandle(true); errorHandler.setCommitRecovered(true); // 新增:恢复后提交偏移量 factory.setCommonErrorHandler(errorHandler);
2. 优化DLQ监听器消费组(可选)
将DLQ监听器的groupId改为独立值(如main-dlq),避免与主Topic消费组混淆,便于后续监控和维护:
@KafkaListener(id = "DLQ-topic", topics = "DLQ-topic", groupId = "main-dlq", clientIdPrefix = "DLQ", autostartup= "false")
3. 验证偏移量提交状态
使用Kafka命令行工具检查消费组偏移量,确认主Topic的偏移量已前进到失败消息之后:
kafka-consumer-groups.sh --bootstrap-server localhost:9092 --describe --group main
验证流程
- 部署修改后的应用,模拟主Topic消息处理失败场景
- 确认消息被推送到DLQ并成功处理
- 重启应用,检查主Topic是否不再重复消费已处理的消息
- 通过命令行工具验证消费组偏移量已正确更新
内容的提问来源于stack exchange,提问作者KGT
相关产品推荐
相关产品推荐

