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(); } } }
问题原因分析
- 重复配置引发冲突:同时使用
@RetryableTopic注解和独立的RetryTopicConfigurationBean,会导致重试配置叠加,生成多套重试主题,引发消息重复消费。 - 重试次数理解偏差:
attempts = "3"指的是总执行次数(首次消费+重试次数),所以实际重试次数是2次;但你看到的4次调用,是重复配置导致的额外消费触发。 - 异常捕获范围不足:
consumer()仅捕获JsonProcessingException,故意抛出的ArithmeticException未被捕获,会触发Spring Kafka的额外失败处理逻辑,加剧重复问题。 - 消费者组配置错误:
idIsGroup=false会让消费者组ID与监听ID解绑,可能导致同一消息被多实例重复处理。
修复方案
- 移除重复配置:删除独立的
RetryTopicConfigurationBean,@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
相关产品推荐
相关产品推荐

