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

KafkaListener集成测试问题:验证消息未被监听器消费

Kafka监听器集成测试断言失败排查

各位好(希望传奇人物Gary Russell能看到)

我正在为Kafka监听器编写集成测试,目标是验证生产者发送的消息能被监听器正常消费,但目前断言检查失败。更关键的是,我在测试代码里重写的onMessage方法中设置了断点,却从未触发。

测试套件代码

@EmbeddedKafka(topics = {TOPIC})
@RunWith(SpringRunner.class)
@SpringBootTest(
    classes = {KafkaIntegrationTestConfiguration.class}
)
@DirtiesContext(classMode = DirtiesContext.ClassMode.AFTER_EACH_TEST_METHOD)
public class KafkaIntegrationTest {
  @Inject
  private KafkaListenerEndpointRegistry kafkaListenerEndpointRegistry;

  @Inject
  private KafkaTemplate<String, Event> producer;

  @SpyBean
  private CustomKafkaConsumer consumer;

  @Test
  public void testKafkaListener() throws Exception {

    // initialize listener container
    ConcurrentMessageListenerContainer<?, ?> container =
      (ConcurrentMessageListenerContainer<?, ?>) kafkaListenerEndpointRegistry.getListenerContainer("custom-consumer");
    container.stop();
    @SuppressWarnings("unchecked")
    AcknowledgingConsumerAwareMessageListener<String, Event> messageListener =
        (AcknowledgingConsumerAwareMessageListener<String, Event>) container
        .getContainerProperties().getMessageListener();
    // create CountDownLatch
    CountDownLatch latch = new CountDownLatch(1);

    // wrap KafkaListener
    container.getContainerProperties()
      .setMessageListener(new AcknowledgingConsumerAwareMessageListener<String, Event>() {

        @Override
        public void onMessage(ConsumerRecord<String, Event> data, Acknowledgment acknowledgment,
                              Consumer<?, ?> consumer) {
          try {
            messageListener.onMessage(data, acknowledgment, consumer);
          } finally {
            latch.countDown();
          }
        }

      });

    container.start();

    // send event to kafka
    sendToKafkaTemplate();

    assertThat(latch.await(10, TimeUnit.SECONDS)).isTrue();
  }
}

生产者配置代码

@EnableKafka
@EmbeddedKafka(topics = {TOPIC})
@Import({ KafkaAutoConfiguration.class })
@ComponentScan("..kafkaConfiguration")
@ComponentScan("..kafka")
public class KafkaIntegrationTestConfiguration {
  @Value("${kafka.mainClusterSchemaRegistry}")
  private String schemaRegistry;

  @Bean
  private KafkaTemplate<String, Event> customKafkaTemplate(EmbeddedKafkaBroker ekb) {

    Map<String, Object> config = KafkaTestUtils.producerProps(ekb);

    config.put(ProducerConfig.KEY_SERIALIZER_CLASS_CONFIG, StringSerializer.class.getName());
    config.put(ProducerConfig.VALUE_SERIALIZER_CLASS_CONFIG, KafkaAvroSerializer.class.getName());
    config.put(AbstractKafkaSchemaSerDeConfig.SCHEMA_REGISTRY_URL_CONFIG, schemaRegistry);

    DefaultKafkaProducerFactory<String, Event> factory = new DefaultKafkaProducerFactory<>(config);

    KafkaTemplate<String, Event> kafkaTemplate = new KafkaTemplate<>(factory);
    kafkaTemplate.setDefaultTopic(TOPIC);
    return kafkaTemplate;
  }
}

Kafka监听器容器工厂配置

@Bean(name = LISTENER_CONTAINER_FACTORY_NAME)
public ConcurrentKafkaListenerContainerFactory<String, Event>
    kafkaListenerContainerFactory(SslConfig sslConfig, KafkaConfig kafkaConfig) {
  final ConcurrentKafkaListenerContainerFactory<String, Event> factory =
      new ConcurrentKafkaListenerContainerFactory<>();

  final ConsumerFactory<String, Event> consumerFactor =
      consumerFactory(sslConfig, kafkaConfig);
  factory.setConsumerFactory(consumerFactor);

  factory.getContainerProperties().setAckMode(ContainerProperties.AckMode.MANUAL_IMMEDIATE);
  // need to update to use new Error Handler
  factory.setErrorHandler(buildListenerErrorHandler());

  return factory;
}


@KafkaListener(
    id = "sync-consumer",
    topics = TOPIC,
    groupId = GROUP_ID,
    containerFactory = LISTENER_CONTAINER_FACTORY_NAME
)

排查方向建议

  • 容器ID不匹配:测试代码中获取容器用的是custom-consumer,但实际@KafkaListener的id是sync-consumer,这会导致获取到的容器为空,后续包装逻辑完全无效,先修正ID匹配问题。
  • 消费者配置未对接EmbeddedKafka:检查consumerFactory的配置,是否正确使用了EmbeddedKafka的broker地址?如果消费者仍指向生产集群,自然无法消费测试环境的消息。
  • 序列化/反序列化不一致:生产者用KafkaAvroSerializer,消费者对应的反序列化器是否配置正确?错误的反序列化器会导致消息无法解析,可能直接被错误处理器拦截,不会进入onMessage。
  • 消息发送未确认:在sendToKafkaTemplate()中添加发送结果阻塞等待,比如producer.send(...).get(),确保消息确实写入EmbeddedKafka。
  • 监听器就绪延迟:容器重启后,等待消费者完成订阅再发送消息,可在container.start()后添加短暂延迟,或用EmbeddedKafka工具类等待消费者就绪。
  • 错误处理器日志排查:自定义错误处理器buildListenerErrorHandler()可能吞掉了异常,添加日志输出,排查是否有反序列化或其他消费错误。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.30 05:40:38