升级Spring Boot/Spring Kafka后自定义DeadLetterPublishingRecoverer未触发的问题排查
升级Spring Boot/Spring Kafka后自定义DeadLetterPublishingRecoverer未触发的问题排查
你好,我之前在升级Spring Boot 3和Spring Kafka 3.x的时候也碰到过类似的问题,DefaultErrorHandler里的Recoverer死活不触发,后来一步步排查才找到原因,给你梳理几个关键排查点和解决方案:
一、先确认容器工厂是否正确绑定了自定义的DefaultErrorHandler
很多时候问题出在这里:你虽然定义了DefaultErrorHandler的Bean,但容器工厂并没有使用它。Spring Boot的ConcurrentKafkaListenerContainerFactoryConfigurer会自动配置一个默认的ErrorHandler,如果你没有手动覆盖,自定义的Bean就不会生效。
你需要在容器工厂的配置方法里,明确设置CommonErrorHandler:
@Bean public ConcurrentKafkaListenerContainerFactory<Object, Object> kafkaListenerContainerFactory( ConcurrentKafkaListenerContainerFactoryConfigurer configurer, ConsumerFactory<Object, Object> consumerFactory, KafkaTemplate<Object, Object> kafkaTemplate, DefaultErrorHandler defaultErrorHandler) { ConcurrentKafkaListenerContainerFactory<Object, Object> factory = new ConcurrentKafkaListenerContainerFactory<>(); configurer.configure(factory, consumerFactory); // 关键步骤:把自定义的DefaultErrorHandler绑定到容器工厂 factory.setCommonErrorHandler(defaultErrorHandler); return factory; }
二、检查BackOff配置是否符合预期
DefaultErrorHandler的重试逻辑和旧的SeekToCurrentErrorHandler有细微差别:
FixedBackOff的第二个参数是重试次数,而不是总尝试次数(总尝试次数=重试次数+1)- 如果你的BackOff设置的重试次数太多,可能需要等所有重试完成后才会触发Recoverer;如果希望失败后直接进入死信,需要把重试次数设为0(即总尝试次数为1)
比如,要实现“第一次失败就触发Recoverer”的配置:
@Bean public DefaultErrorHandler defaultErrorHandler(ConsumerRecordRecoverer deadLetterRecoverer) { // FixedBackOff(间隔时间, 重试次数),这里重试次数设为0,失败后直接走Recoverer FixedBackOff backOff = new FixedBackOff(0L, 0); DefaultErrorHandler errorHandler = new DefaultErrorHandler(deadLetterRecoverer, backOff); // 可选:如果某些异常不需要重试(比如反序列化异常,重试也没用),直接标记为不可重试 errorHandler.addNotRetryableException(DeserializationException.class); return errorHandler; }
三、验证DeadLetterPublishingRecoverer的主题解析逻辑是否正确
如果你的Recoverer是DeadLetterPublishingRecoverer,那要确保BiFunction能正确生成死信主题的TopicPartition:
- 比如你是不是写错了死信主题的命名规则?
- 有没有权限往死信主题发送消息?
- 主题是否存在(如果是自动创建主题,要确保Kafka集群开启了自动创建)
示例正确的Recoverer配置:
@Bean public ConsumerRecordRecoverer deadLetterRecoverer(KafkaTemplate<Object, Object> kafkaTemplate) { return new DeadLetterPublishingRecoverer(kafkaTemplate, (record, ex) -> { // 这里根据业务需求生成死信主题,比如原主题后缀加"-dlq" String dlqTopic = record.topic() + "-dlq"; return new TopicPartition(dlqTopic, record.partition()); }); }
四、查看日志排查重试流程
开启Spring Kafka的DEBUG日志,观察以下关键日志:
- 有没有出现
Retrying delivery for record(s):如果有,说明重试逻辑在执行,需要等重试次数耗尽才会触发Recoverer - 有没有出现
Recovering record after X attempts:如果有,说明已经进入Recoverer逻辑,此时要检查Recoverer内部是否有异常(比如发送死信失败被吞了) - 如果没有任何重试或恢复的日志,那大概率是容器工厂没绑定到自定义的ErrorHandler
五、注意异常类型的处理
DefaultErrorHandler默认会重试所有RuntimeException,但如果你的异常属于以下情况,可能不会触发重试:
- 被标记为不可重试的异常(通过
addNotRetryableException添加的) - 属于
FatalException子类的异常(比如序列化相关异常)
如果你的业务异常不需要重试,直接进入死信,记得把它加入不可重试列表:
errorHandler.addNotRetryableException(YourBusinessException.class);
备注:内容来源于stack exchange,提问作者corstad
相关产品推荐
相关产品推荐

