Spring Kafka批量监听器自定义主题非阻塞重试方案咨询
方案选型结论
针对高版本Spring Kafka(2.8.x及以上,含当前主流3.x版本)的批量消费非阻塞重试场景,选型结论如下:
- 不推荐使用
RecoveringBatchErrorHandler:该组件实现的是阻塞式重试,重试等待期间会占用消费者线程,对应分区无法拉取新消息,吞吐量低,不符合非阻塞指数退避的要求。 - 优先使用Spring Kafka原生非阻塞重试能力(底层基于
RetryTopicConfigurer实现):不需要为每个重试主题、死信主题单独编写@KafkaListener方法,框架会自动生成对应主题的转发、延迟消费逻辑;完全支持手动指定预创建的主题名称,可关闭框架自动建主题能力,完全匹配你的约束条件。 - 你预先创建的2个重试主题+1个死信主题的结构,和这套机制的逻辑完全适配,不需要调整主题命名或数量。
实现代码示例
基础配置类
配置批量消费容器工厂,绑定非阻塞重试规则,关闭自动建主题能力:
@Configuration @EnableKafka public class KafkaBatchConfig { @Bean public ConsumerFactory<String, String> consumerFactory() { Map<String, Object> props = new HashMap<>(); props.put(ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG, "你的Kafka集群地址"); props.put(ConsumerConfig.GROUP_ID_CONFIG, "product-consumer-group"); props.put(ConsumerConfig.KEY_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class); props.put(ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class); props.put(ConsumerConfig.MAX_POLL_RECORDS_CONFIG, 500); props.put(ConsumerConfig.ENABLE_AUTO_COMMIT_CONFIG, false); return new DefaultKafkaConsumerFactory<>(props); } @Bean public ProducerFactory<String, String> producerFactory() { Map<String, Object> props = new HashMap<>(); props.put(ProducerConfig.BOOTSTRAP_SERVERS_CONFIG, "你的Kafka集群地址"); props.put(ProducerConfig.KEY_SERIALIZER_CLASS_CONFIG, StringSerializer.class); props.put(ProducerConfig.VALUE_SERIALIZER_CLASS_CONFIG, StringSerializer.class); return new DefaultKafkaProducerFactory<>(props); } @Bean public KafkaTemplate<String, String> kafkaTemplate(ProducerFactory<String, String> producerFactory) { return new KafkaTemplate<>(producerFactory); } @Bean public KafkaListenerContainerFactory<ConcurrentMessageListenerContainer<String, String>> kafkaBatchListenerFactory( ConsumerFactory<String, String> consumerFactory, KafkaTemplate<String, String> kafkaTemplate ) { ConcurrentKafkaListenerContainerFactory<String, String> factory = new ConcurrentKafkaListenerContainerFactory<>(); factory.setConsumerFactory(consumerFactory); factory.setBatchListener(true); factory.getContainerProperties().setAckMode(ContainerProperties.AckMode.MANUAL_IMMEDIATE); // 配置非阻塞重试规则 RetryTopicConfigurationBuilder retryConfigBuilder = RetryTopicConfigurationBuilder .newInstance() // 完全关闭框架自动创建主题的能力,使用预创建的主题 .autoCreateTopics(false, false) // 指数退避规则:首次重试间隔1s,间隔倍数2.0,最大重试间隔10s .exponentialBackoff(1000, 2.0, 10000) // 最大尝试次数=1次主消费+2次重试,共3次,耗尽后进入死信 .maxAttempts(3) // 自定义重试主题命名,直接映射预创建的ProductTopic.Retry-1、ProductTopic.Retry-2 .retryTopicNaming((mainTopic, attempt, context) -> mainTopic + ".Retry-" + attempt) // 自定义死信主题命名,映射预创建的ProductTopic.Retry-DLT .dltNaming(mainTopic -> mainTopic + ".Retry-DLT") .withListenerFactory(factory) .withKafkaTemplate(kafkaTemplate); factory.setRetryTopicConfiguration(retryConfigBuilder.build()); return factory; } }
业务监听器
保留你原有核心业务逻辑,不需要为重试/死信主题单独编写监听器:
@Slf4j @Component public class ProductTopicListener { @KafkaListener( topics = "ProductTopic", containerFactory = "kafkaBatchListenerFactory" ) public void onBatch(List<Message<String>> messages, Acknowledgment acknowledgment) { // 原有DB批量插入逻辑 consume(messages); // 业务执行成功再提交偏移量 acknowledgment.acknowledge(); } // 可选:自定义死信处理逻辑,不需要单独配置死信监听器 @DltHandler public void handleDeadLetterMessages(List<Message<String>> dltMessages) { log.error("消息两次重试全部失败进入死信,消息数量:{}", dltMessages.size()); // 可扩展死信落库、告警等逻辑 } }
最佳实践说明
- 非阻塞重试执行逻辑:主主题消费抛出异常时,框架会将失败的整批消息转发到对应次数的重试主题,转发成功后提交主主题偏移量,不会阻塞主消费者线程;重试主题的监听器会按照配置的退避延迟时间拉取消息消费,消费失败则转发到下一级重试主题,两次重试都失败后转发至死信主题。
- 配置
autoCreateTopics(false, false)后,框架不会向Broker发送任何创建主题的请求,完全使用你提前建好的主题,满足环境限制要求。 - 重试主题的监听器默认和主主题使用相同的消费配置,如果需要给重试/死信主题单独配置并发数、拉取批次大小、消费隔离级别,可以通过
RetryTopicConfigurationBuilder的retryTopicConfigurers、dltConfigurer方法单独指定,不需要手写监听器逻辑。 - 偏移量提交保持原有手动提交逻辑即可:业务执行成功手动提交偏移量,消费失败时框架完成重试/死信转发后会自动提交对应偏移量,不会出现消息丢失、重复消费卡死的问题。
- 2.9版本之后的Spring Kafka已经原生支持批量模式下的非阻塞重试,消息转发、批次处理逻辑都经过官方验证,稳定性远高于自定义错误处理器的实现,不需要自己实现消息转发、退避调度的逻辑。
内容的提问来源于stack exchange,提问作者hideburn
相关产品推荐
相关产品推荐

