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
相关产品推荐
相关产品推荐

