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
相关产品推荐
相关产品推荐

