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

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消费主主题、重试主题消息,自动完成重试流程。需实现两个核心目标:

  1. 监听器内获取消息的重试次数(maxAttempts设为3,需区分首次重试、二次重试)
  2. 统计以下消费指标:
    • 主主题成功消费消息数
    • 重试主题成功消费消息数
    • 死信队列(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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.14 13:10:30