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

Kafka @RetryableTopic重试后@DltHandler多次触发问题排查

Kafka @RetryableTopic 重试次数与DLT执行异常问题

我在测试Kafka的@RetryableTopic特性,预期逻辑是消息重试3次后推送至SQS,但实际运行出现两个异常:

  • consumer()方法被重试线程调用了4次
  • @DltHandler标注的processMessage()方法被调用了2次

按Spring Retry机制,@DltHandler应该只在所有重试耗尽后触发一次,目前通过线程名称跟踪consumer()的调用,但找不到配置遗漏的点。

相关代码

@Bean(ProductServiceConstants.PRODUCT_KAFKA_DLT_PRODUCER_FACTORY)
public KafkaTemplate<String, String> kafkaTemplateForDlt() {
    return new KafkaTemplate<>(producerFactory());
}

@Bean
public RetryTopicConfiguration myRetryTopic(@Qualifier(ProductServiceConstants.PRODUCT_KAFKA_DLT_PRODUCER_FACTORY)KafkaTemplate<String, String> template) {
    return RetryTopicConfigurationBuilder
            .newInstance()
            .create(template);
}

@Slf4j
@Component
public class ProductEventConsumer {
@Autowired
private ProductServiceImpl productServiceImpl;

@Autowired ObjectMapper objectMapper;

@Value("${aws.sqsDLQ}")
private String productdlq;

@Autowired
private SqsTemplate sqsTemplate;

@RetryableTopic(
          backoff = @Backoff(delayExpression = "10000", multiplierExpression = "0"), 
          attempts = "3", 
          kafkaTemplate = ProductServiceConstants.PRODUCT_KAFKA_DLT_PRODUCER_FACTORY,
          include = {SocketTimeoutException.class,ArithmeticException.class})
@KafkaListener(id=ProductServiceConstants.PRODUCT_KAFKA_CONSUMER_ID, idIsGroup=false,
        topics="#{'${spring.kafka.product-topic}'}",containerFactory=ProductServiceConstants.PRODUCT_KAFKA_CONSUMER_FACTORY)
public void consumer(ConsumerRecord<String,String> consumerRecord, Acknowledgment ack) {
    try{
        log.info("START:Received request via kafka:{} thread:{}",consumerRecord.value()
                ,Thread.currentThread().getName());
        int result = 10 / 0;
        ack.acknowledge();
    }catch( JsonProcessingException e) {
        log.error("END:Exception occured while saving item:{}",e.getMessage());
    }
}

@DltHandler
public void processMessage(ConsumerRecord<String,String> consumerRecord, Acknowledgment ack) {
    try{
        log.error("START:Pushing message to SQS DLQ:{}",consumerRecord.key());
        sqsTemplate.send(sqsSendOptions -> sqsSendOptions.queue(productdlq).payload(consumerRecord.value()));
    }catch(Exception e) {
        log.error("END:Failure while pushing msg to sqs dlq:{} key:{}",e.getMessage(),consumerRecord.key());
    }
    finally {
        ack.acknowledge();
    }
}
}

问题原因分析

  1. 重复配置引发冲突:同时使用@RetryableTopic注解和独立的RetryTopicConfiguration Bean,会导致重试配置叠加,生成多套重试主题,引发消息重复消费。
  2. 重试次数理解偏差:attempts = "3"指的是总执行次数(首次消费+重试次数),所以实际重试次数是2次;但你看到的4次调用,是重复配置导致的额外消费触发。
  3. 异常捕获范围不足:consumer()仅捕获JsonProcessingException,故意抛出的ArithmeticException未被捕获,会触发Spring Kafka的额外失败处理逻辑,加剧重复问题。
  4. 消费者组配置错误:idIsGroup=false会让消费者组ID与监听ID解绑,可能导致同一消息被多实例重复处理。

修复方案

  • 移除重复配置:删除独立的RetryTopicConfiguration Bean,@RetryableTopic注解已足够完成重试规则配置。
  • 修正重试次数:如果需要3次重试(不含首次消费),将attempts设为"4"(首次+3次重试)。
  • 统一异常处理:要么在consumer()中捕获所有目标异常,要么直接抛出异常让@RetryableTopic处理,避免未捕获异常引发额外逻辑。
  • 修正消费者组设置:将idIsGroup设为默认值true,确保同一组内消费者不会重复消费。

修正后核心代码示例

// 移除多余的RetryTopicConfiguration Bean

@Slf4j
@Component
public class ProductEventConsumer {
    // ... 保留原有注入代码

    @RetryableTopic(
            backoff = @Backoff(delayExpression = "10000", multiplierExpression = "0"),
            attempts = "4", // 首次消费 + 3次重试
            kafkaTemplate = ProductServiceConstants.PRODUCT_KAFKA_DLT_PRODUCER_FACTORY,
            include = {SocketTimeoutException.class, ArithmeticException.class})
    @KafkaListener(id = ProductServiceConstants.PRODUCT_KAFKA_CONSUMER_ID,
            topics = "#{'${spring.kafka.product-topic}'}",
            containerFactory = ProductServiceConstants.PRODUCT_KAFKA_CONSUMER_FACTORY)
    public void consumer(ConsumerRecord<String, String> consumerRecord, Acknowledgment ack) {
        log.info("START:Received request via kafka:{} thread:{}", consumerRecord.value(), Thread.currentThread().getName());
        int result = 10 / 0; // 故意抛出异常
        ack.acknowledge();
    }

    @DltHandler
    public void processMessage(ConsumerRecord<String, String> consumerRecord, Acknowledgment ack) {
        try {
            log.error("START:Pushing message to SQS DLQ:{}", consumerRecord.key());
            sqsTemplate.send(sqsSendOptions -> sqsSendOptions.queue(productdlq).payload(consumerRecord.value()));
        } catch (Exception e) {
            log.error("END:Failure while pushing msg to sqs dlq:{} key:{}", e.getMessage(), consumerRecord.key());
        } finally {
            ack.acknowledge();
        }
    }
}

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.27 16:14:52