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

Kafka消费者异常重试10次,添加MDC日志后恢复正常的原因解析

问题:Kafka监听器重复接收单条消息10次,添加MDC后恢复正常的原因

场景描述

向Kafka主题发送单条消息后,监听器会重复接收并处理该消息共10次(每次在前一次处理完成后触发)。添加一行MDC日志代码后,消费者恢复正常,仅处理一次消息。

消费者核心代码

@KafkaListener(topics = "xxx", groupId = "xxx")
public void onMessageUser(ConsumerRecord<?, ?> record, Acknowledgment ack) {
    // 业务处理代码
    // some code
    ack.acknowledge();
}

Kafka相关配置

spring.kafka.producer.retries=5
spring.kafka.producer.acks=1
spring.kafka.consumer.max-poll-records=50
spring.kafka.consumer.poll-timeout=5000
spring.kafka.consumer.enable-auto-commit=false
spring.kafka.listener.ack-mode=MANUAL

添加后恢复正常的MDC代码:

org.slf4j.MDC.put("requestId", xxx);

根因分析

感谢@Artem Bilan和@Gary Russell的提示,已定位问题根源:

1. 触发异常的代码

业务处理片段中的代码存在空指针风险:

Map<String,Object> prop = new ConcurrentHashMap<>();   
prop.put(CommonConstants.Request.REQUEST_ID,MDC.get(CommonConstants.Request.REQUEST_ID));

当MDC.get(CommonConstants.Request.REQUEST_ID)返回null时,ConcurrentHashMap的put方法会抛出空指针异常(因为ConcurrentHashMap不允许key或value为null),异常栈如下:

java.lang.NullPointerException: null    
at java.base/java.util.concurrent.ConcurrentHashMap.putVal(ConcurrentHashMap.java:1011)     
at java.base/java.util.concurrent.ConcurrentHashMap.put(ConcurrentHashMap.java:1006)

2. 异常未被发现的原因

  • 该业务代码未添加try-catch块捕获异常;
  • 项目中的Spring Web全局异常处理器仅处理Web请求链路的异常,不覆盖Kafka消息处理流程,导致异常未被记录到错误日志,长期未被察觉。

3. 重复消费的触发逻辑

未捕获的异常会触发Spring Kafka的默认重试策略,默认重试次数为9次,加上首次处理,总共会执行10次消息处理。

4. 添加MDC后恢复正常的原因

添加MDC.put("requestId", xxx)后,MDC.get()能获取到有效值,不再触发空指针异常,消息处理正常完成后执行ack.acknowledge()提交偏移量,因此消费者仅处理一次消息。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.23 16:52:35