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

如何避免主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

验证流程

  1. 部署修改后的应用,模拟主Topic消息处理失败场景
  2. 确认消息被推送到DLQ并成功处理
  3. 重启应用,检查主Topic是否不再重复消费已处理的消息
  4. 通过命令行工具验证消费组偏移量已正确更新

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.20 06:30:51