Spring Kafka非阻塞重试:如何获取重试次数并统计相关指标
Kafka非阻塞重试场景下的重试次数获取与消费统计方案
问题背景
现有Spring Boot应用无生产者,包含2个监听器,基于org.springframework.kafka.retrytopic配置了Kafka非阻塞重试机制,核心配置Bean如下:
@Bean RetryTopicConfiguration kafkaRetryTopicConfiguration(final ConcurrentKafkaListenerContainerFactory<String, String> kafkaListenerContainerFactory, final KafkaTemplate<String, String> kafkaTemplate) { return RetryTopicConfigurationBuilder .newInstance() .doNotAutoCreateRetryTopics() .autoStartDltHandler(false) .useSingleTopicForFixedDelays() .includeTopics(kafkaTopics.getConsumerTopics()) .maxAttempts(maxAttempts) .fixedBackOff(backOffPeriod.toMillis()) .retryTopicSuffix( "-" + kafkaProperties.getConsumer().getGroupId() + kafkaProperties.getProperties().get("retry.suffix")) .dltSuffix("-" + kafkaProperties.getConsumer().getGroupId() + kafkaProperties.getProperties().get("error.suffix")) .listenerFactory(kafkaListenerContainerFactory) .create(kafkaTemplate); }
当前配置通过单个通用@KafkaListener消费主主题、重试主题消息,自动完成重试流程。需实现两个核心目标:
- 监听器内获取消息的重试次数(
maxAttempts设为3,需区分首次重试、二次重试) - 统计以下消费指标:
- 主主题成功消费消息数
- 重试主题成功消费消息数
- 死信队列(DLT)消息数(已通过拦截器实现)
- 首次重试后处理的消息数
- 二次重试后处理的消息数
解决方案
一、获取重试次数:利用RetryTopic内置消息头
Spring Kafka的RetryTopic机制会自动在重试消息中注入kafka_retry_attempt消息头,记录当前重试次数(计数从1开始,主主题消费时该头不存在或值为0)。在监听器中直接从ConsumerRecord提取该值即可:
@KafkaListener(topics = "#{'${kafka.consumer.topics}'.split(',')}", groupId = "${kafka.consumer.group-id}") public void consume(ConsumerRecord<String, String> record) { // 初始化重试次数为0(主主题消费场景) int retryAttempt = 0; Header retryHeader = record.headers().lastHeader("kafka_retry_attempt"); if (retryHeader != null) { retryAttempt = Integer.parseInt(new String(retryHeader.value())); } // 传入重试次数执行业务逻辑 processMessage(record, retryAttempt); }
二、消费统计:结合重试次数与主题类型实现
通过自定义线程安全的统计类,在监听器中根据重试次数和消费场景更新统计值:
1. 统计类定义
@Component public class KafkaConsumeMetrics { // 主主题成功消费数 private final AtomicInteger mainTopicSuccessCount = new AtomicInteger(0); // 重试主题成功消费总数 private final AtomicInteger retryTopicSuccessCount = new AtomicInteger(0); // 首次重试成功处理数 private final AtomicInteger firstRetrySuccessCount = new AtomicInteger(0); // 二次重试成功处理数 private final AtomicInteger secondRetrySuccessCount = new AtomicInteger(0); // DLT消息数(复用已有拦截器逻辑) private final AtomicInteger dltCount = new AtomicInteger(0); // 自增方法 public void incrementMainTopicSuccess() { mainTopicSuccessCount.incrementAndGet(); } public void incrementRetryTopicSuccess() { retryTopicSuccessCount.incrementAndGet(); } public void incrementFirstRetrySuccess() { firstRetrySuccessCount.incrementAndGet(); } public void incrementSecondRetrySuccess() { secondRetrySuccessCount.incrementAndGet(); } public void incrementDltCount() { dltCount.incrementAndGet(); } // 省略各统计值的getter方法 }
2. 监听器中更新统计
@Autowired private KafkaConsumeMetrics consumeMetrics; @KafkaListener(topics = "#{'${kafka.consumer.topics}'.split(',')}", groupId = "${kafka.consumer.group-id}") public void consume(ConsumerRecord<String, String> record) { int retryAttempt = 0; Header retryHeader = record.headers().lastHeader("kafka_retry_attempt"); if (retryHeader != null) { retryAttempt = Integer.parseInt(new String(retryHeader.value())); } try { // 执行业务处理逻辑 businessProcess(record.value()); // 根据重试次数更新统计指标 if (retryAttempt == 0) { consumeMetrics.incrementMainTopicSuccess(); } else { consumeMetrics.incrementRetryTopicSuccess(); if (retryAttempt == 1) { consumeMetrics.incrementFirstRetrySuccess(); } else if (retryAttempt == 2) { consumeMetrics.incrementSecondRetrySuccess(); } } } catch (Exception e) { // 异常会由RetryTopic机制自动转发到重试主题或DLT,无需手动处理 throw new RuntimeException("消费失败,触发重试流程", e); } }
三、补充说明
- 若需要更精准区分主题类型,也可通过
record.topic()判断主题后缀,但利用kafka_retry_attempt消息头更直接可靠,无需依赖主题命名规则。 - DLT消息统计可直接在已实现的拦截器中调用
consumeMetrics.incrementDltCount()完成计数。 kafka_retry_attempt是RetryTopic机制的内置消息头,无需手动添加,只要配置了RetryTopic就会自动注入到重试消息中。
内容的提问来源于stack exchange,提问作者sammy20
相关产品推荐
相关产品推荐

