@RetryableTopic搭配多KafkaListenerContainerFactory失效排查
问题分析与解决方案
错误根源
你遇到的FilteringBatchMessageListenerAdapter无法转换为MessageListener的错误,核心原因是:尽管你将容器工厂的batchListener设为false,但工厂层面配置的recordFilterStrategy与@RetryableTopic的容器创建逻辑发生冲突,导致容器错误生成了批量类型的消息监听适配器,而单条消息监听器无法兼容该类型。
修复步骤
1. 移除工厂层面的过滤配置,改用方法级过滤
将过滤逻辑从容器工厂移到监听器方法上,使用@Filter注解实现,避免触发批量适配器的生成:
@Slf4j @RequiredArgsConstructor @Component public class RetryConsumer { LocalDateTime now; @RetryableTopic( attempts = "5", autoCreateTopics = "false", backoff = @Backoff(30_000), fixedDelayTopicStrategy = FixedDelayStrategy.SINGLE_TOPIC, listenerContainerFactory = "retryKafkaListenerContainerFactory" ) @KafkaListener(id = "retry-consumer", topics = "topic1", containerFactory = "retryKafkaListenerContainerFactory") @Filter("messageFilter") // 新增Filter注解绑定过滤方法 public void handleMessage ( String message, @Header(KafkaHeaders.RECEIVED_TOPIC) String topic ) { // 原有业务逻辑 log.info("Received message: {} from topic: {}", message, topic); if (now == null) { now = LocalDateTime.now(); System.out.println("attempt at:" + now.toString()); } else { System.out.println("minutes between attempts:" + Duration.between(LocalDateTime.now(), now).toSeconds()); } throw new RuntimeException("Test exception"); } // 新增过滤方法,实现原有过滤逻辑 public boolean messageFilter(ConsumerRecord<String, String> record) { if (Boolean.TRUE.equals(kafkaProperties.getConsumer().getLogPayload())) { log.info(record.toString()); } return false; } @DltHandler public void handleDlt ( String message, @Header(KafkaHeaders.RECEIVED_TOPIC) String topic ) { log.info("Message: {} handled by dlq topic: {}", message, topic); } }
2. 调整容器工厂配置
移除setRecordFilterStrategy的调用,同时建议注入已有的ConsumerFactory而非手动创建,避免重复配置引发的问题:
@Bean @Qualifier("retryKafkaListenerContainerFactory") public ConcurrentKafkaListenerContainerFactory<String, String> retryKafkaListenerContainerFactory( MessageConverter messageConverter, ConsumerFactory<String, String> consumerFactory // 注入已有ConsumerFactory ) { ConcurrentKafkaListenerContainerFactory<String, String> factory = new ConcurrentKafkaListenerContainerFactory<>(); factory.setBatchListener(false); factory.setMessageConverter(messageConverter); factory.setCommonErrorHandler(defaultErrorHandler()); factory.setConsumerFactory(consumerFactory); // 移除setRecordFilterStrategy相关配置 return factory; }
3. 检查MessageConverter类型
确保你注入的MessageConverter是单条消息类型(如StringJsonMessageConverter),而非批量消息转换器(BatchMessagingMessageConverter)。如果是批量转换器,替换为单条类型的转换器。
补充说明
容器工厂层面的recordFilterStrategy在与@RetryableTopic结合使用时,容易触发内部逻辑对适配器类型的误判,改用方法级的@Filter注解既能实现相同的过滤效果,又能避免类型不兼容的问题。
内容的提问来源于stack exchange,提问作者BadChanneler
相关产品推荐
相关产品推荐

