defaulterrorhandler是spring kafka中处理消费异常的核心组件,支持重试、跳过、提交偏移量及投递死信(dlt);需配置backoff控制重试节奏、指定recoverer(如deadletterpublishingrecoverer)、明确fatal异常并确保dlt topic可写且有监听。

Spring Kafka 中的 DefaultErrorHandler 是处理消费异常的核心组件,它替代了已弃用的 SeekToCurrentErrorHandler,支持重试、跳过、提交偏移量以及投递死信(DLT)等策略。配置的关键在于:**明确重试行为、指定恢复器(Recoverer)、控制哪些异常不重试,并确保 DLT Topic 可写且监听到位**。
重试次数与间隔配置
通过 BackOff 控制重试节奏,避免密集失败打满资源:
-
固定间隔重试:使用
FixedBackOff(3000L, 3L)表示每次间隔 3 秒,最多重试 3 次(含首次失败共 4 次尝试) -
指数退避重试:用
ExponentialBackOff(1000L, 2.0, 60000L)实现首延 1s、倍数 2、上限 60s 的退避策略 - 重试次数为 0 时,表示不重试,直接进入恢复逻辑(如发往 DLT)
绑定死信队列(DLT)的两种方式
必须配合 ConsumerRecordRecoverer 才能触发 DLT 投递:
-
默认自动 DLT Topic:使用
DeadLetterPublishingRecoverer,会自动将消息发往{original-topic}-DLT(如order-topic-DLT),前提是该 Topic 已存在或启用autoCreateTopics=true -
自定义 DLT Topic 名称:构造
DeadLetterPublishingRecoverer时传入TopicPartition或Function,实现按业务/错误类型路由到不同 DLT Topic - 若仅需记录日志或人工干预,也可用
LoggingConsumerRecordRecoverer替代,不发 DLT
控制哪些异常不重试(关键细节)
某些异常重试无意义,DefaultErrorHandler 默认将其归为 fatal 类型,跳过重试直接恢复:
- 默认 fatal 异常包括:
DeserializationException、MessageConversionException、ClassCastException等 - 若想让反序列化失败也重试(例如网络抖动导致 Schema Registry 临时不可用),需显式移除:
errorHandler.removeClassification(DeserializationException.class) - 可添加自定义不可重试异常:
errorHandler.addNotRetryableException(BusinessValidationException.class)
完整配置示例(Bean 方式)
在容器中声明 ConcurrentKafkaListenerContainerFactory 并注入定制化 DefaultErrorHandler:
@Bean
public ConcurrentKafkaListenerContainerFactory<string string> kafkaListenerContainerFactory(
ConsumerFactory<string string> consumerFactory) {
ConcurrentKafkaListenerContainerFactory<string string> factory =
new ConcurrentKafkaListenerContainerFactory();
factory.setConsumerFactory(consumerFactory);
factory.getContainerProperties().setAckMode(ContainerProperties.AckMode.MANUAL_IMMEDIATE);
// 配置重试与恢复
BackOff backOff = new FixedBackOff(2000L, 2L); // 2秒间隔,最多再试2次(共3次)
DeadLetterPublishingRecoverer recoverer =
new DeadLetterPublishingRecoverer(kafkaTemplate); // 自动推送到 -DLT Topic
DefaultErrorHandler errorHandler = new DefaultErrorHandler(recoverer, backOff);
errorHandler.addNotRetryableException(IllegalArgumentException.class); // 明确不重试某类业务异常
factory.setCommonErrorHandler(errorHandler);
return factory;
}</string></string></string>
注意:kafkaTemplate 必须已配置好 producer,且具备向 DLT Topic 写入权限;同时需有对应 @KafkaListener 监听 DLT Topic 做后续人工或自动补偿处理。
Java免费学习笔记:立即使用
解锁 Java 大师之旅:从入门到精通的终极指南











