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

@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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.23 10:52:03