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

Spring Kafka 2.X迁移至3.0:RecoveryCallback适配方案咨询

Spring Kafka 3.0 实现重试与恢复回调(替代2.X版本方案)

Spring Kafka 3.0移除了ConcurrentKafkaListenerContainerFactory直接设置RetryTemplate和RecoveryCallback的API,改用DefaultErrorHandler统一处理重试及恢复逻辑,以下是和你旧代码功能一致的实现方案:

@Bean
public ConcurrentKafkaListenerContainerFactory<?, ?> retryConcurrentKafkaListenerContainerFactory() {
    ConcurrentKafkaListenerContainerFactory<Object, Object> factory = new ConcurrentKafkaListenerContainerFactory<>();
    factory.setConsumerFactory(consumerFactory());

    // 配置重试规则:RetryTemplate定义重试次数、退避策略等
    RetryTemplate retryTemplate = new RetryTemplate();
    // 示例:最大重试3次
    SimpleRetryPolicy retryPolicy = new SimpleRetryPolicy();
    retryPolicy.setMaxAttempts(3);
    retryTemplate.setRetryPolicy(retryPolicy);
    // 示例:每次重试间隔1秒
    FixedBackOffPolicy backOffPolicy = new FixedBackOffPolicy();
    backOffPolicy.setBackOffPeriod(1000);
    retryTemplate.setBackOffPolicy(backOffPolicy);

    // 创建DefaultErrorHandler,整合重试规则和恢复回调
    DefaultErrorHandler errorHandler = new DefaultErrorHandler(
            // 重试耗尽后的恢复逻辑
            (failedRecord, exception) -> {
                ConsumerRecord<?, ?> record = (ConsumerRecord<?, ?>) failedRecord;
                System.out.println("Recovery callback invoked for record: " + record.value());

                // 获取手动提交的Acknowledgment(如果开启了手动提交模式)
                Acknowledgment acknowledgment = record.headers().lastHeader(KafkaHeaders.ACKNOWLEDGMENT) != null
                        ? (Acknowledgment) record.headers().lastHeader(KafkaHeaders.ACKNOWLEDGMENT).value()
                        : null;
                if (acknowledgment != null) {
                    acknowledgment.acknowledge();
                }
            },
            // 将RetryTemplate的策略适配到ErrorHandler中
            new RetryPolicyBackOffManager(retryTemplate)
    );

    // 为容器工厂设置错误处理器
    factory.setCommonErrorHandler(errorHandler);

    // 其他自定义配置(如手动提交模式等,按需调整)
    // factory.getContainerProperties().setAckMode(ContainerProperties.AckMode.MANUAL);

    return factory;
}

关键变化说明

  • 不再使用RetryOperationsInterceptor,转而通过DefaultErrorHandler作为统一错误处理入口,整合重试和恢复逻辑。
  • 恢复回调直接作为DefaultErrorHandler的构造参数传入,failedRecord就是触发失败的消息记录。
  • 获取Acknowledgment的方式改为从消息头KafkaHeaders.ACKNOWLEDGMENT中提取,替代旧版本从重试上下文获取的逻辑。
  • RetryPolicyBackOffManager负责把你原来RetryTemplate里的重试、退避规则映射到DefaultErrorHandler,保证和旧版本行为一致。

如果需要针对特定异常设置不同重试策略,可替换SimpleRetryPolicy为ExceptionClassifierRetryPolicy实现更灵活的规则。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.16 07:53:19