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

