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

从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版本的以下行为:

  1. 消息处理失败时不进行任何重试
  2. 仅特定异常触发自定义恢复逻辑
  3. 非特定异常仅记录错误日志,不执行额外处理

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.02 08:10:23