从spring-kafka 2.7.14迁移到3.0.1的CommonErrorHandler配置咨询
Spring Kafka 3.0.1 适配配置(复现2.7.14原有行为)
核心思路:用DefaultErrorHandler整合所有逻辑
Spring Kafka 3.x中DefaultErrorHandler是CommonErrorHandler的默认实现,可一站式替代旧版RetryTemplate+RecoveryCallback+ErrorHandler的组合,完全匹配你需要的原有行为。
完整配置示例
import org.apache.kafka.clients.consumer.ConsumerRecord; import org.springframework.context.annotation.Bean; import org.springframework.context.annotation.Configuration; import org.springframework.kafka.config.ConcurrentKafkaListenerContainerFactory; import org.springframework.kafka.core.ConsumerFactory; import org.springframework.kafka.core.KafkaTemplate; import org.springframework.kafka.listener.DefaultErrorHandler; import org.springframework.kafka.listener.ConsumerRecordRecoverer; import org.springframework.util.backoff.FixedBackOff; import org.slf4j.Logger; import org.slf4j.LoggerFactory; @Configuration public class KafkaConfig { private static final Logger log = LoggerFactory.getLogger(KafkaConfig.class); @Bean public ConcurrentKafkaListenerContainerFactory<String, Object> kafkaListenerContainerFactory( ConsumerFactory<String, Object> consumerFactory, KafkaTemplate<String, Object> kafkaTemplate) { ConcurrentKafkaListenerContainerFactory<String, Object> factory = new ConcurrentKafkaListenerContainerFactory<>(); factory.setConsumerFactory(consumerFactory); // 1. 替代原RecoveryCallback:自定义恢复逻辑 ConsumerRecordRecoverer customRecoverer = (record, exception) -> { // 这里实现你原RecoveryCallback的业务逻辑,比如发死信队列、告警等 log.error("消息处理失败,执行恢复逻辑:topic={}, offset={}", record.topic(), record.offset(), exception); // 示例:发送到死信队列 // kafkaTemplate.send("dead-letter-topic", record.key(), record.value()); }; // 2. 实现永不重试:FixedBackOff(0L, 0)表示无等待、无重试,直接走后续逻辑 FixedBackOff neverRetryBackOff = new FixedBackOff(0L, 0); // 3. 初始化DefaultErrorHandler,整合恢复器与重试规则 DefaultErrorHandler errorHandler = new DefaultErrorHandler(customRecoverer, neverRetryBackOff); // 4. 配置异常过滤:仅特定异常触发恢复,其他异常仅记日志 // 这里替换成你需要的特定异常类,比如YourBusinessException errorHandler.setSkipRecoveryFor(exception -> !(exception instanceof YourSpecificException)); // 绑定错误处理器到容器工厂 factory.setCommonErrorHandler(errorHandler); return factory; } // 自定义特定异常示例(替换成你实际业务中的异常类) static class YourSpecificException extends RuntimeException { public YourSpecificException(String message) { super(message); } } }
关键配置说明
- 永不重试实现:
FixedBackOff(0L, 0)完全替代旧版factory.setRetryTemplate(neverRetry...)的逻辑,失败后直接进入后续处理。 - 恢复逻辑替代:
ConsumerRecordRecoverer直接对应原RecoveryCallback的业务逻辑,可灵活扩展死信、告警等操作。 - 异常过滤控制:
setSkipRecoveryFor通过Lambda判断异常类型,仅让指定异常进入恢复流程,其他异常自动记录错误日志后跳过恢复,完全匹配你的需求。
行为匹配验证
该配置1:1复现Spring Kafka 2.7.14版本的以下行为:
- 消息处理失败时不进行任何重试
- 仅特定异常触发自定义恢复逻辑
- 非特定异常仅记录错误日志,不执行额外处理
内容的提问来源于stack exchange,提问作者user2717708
相关产品推荐
相关产品推荐

