SpringBoot中用Testcontainers+Mockito测试Kafka消费者未触发问题排查
搭建Kafka消费者时,使用Testcontainers和Mockito做集成测试,生产者发送消息正常,但消费者始终未触发,抛出如下错误:
Wanted but not invoked:
helloKafkaService.handleMessage(
);
-> at com.my.kafka.HelloKafkaListenerTestIT.helloKafkaListenerTest
Actually, there were zero interactions with this mock.
相关代码
服务接口HelloKafkaService.java
@Service public interface HelloKafkaService { public void handleMessage(ExampleEvent exampleEvent); }
消费者HelloKafkaListener.java
@Component public class HelloKafkaListener { private static final Logger log = LoggerFactory.getLogger(HelloKafkaListener.class); private final HelloKafkaService helloKafkaService; public HelloKafkaListener(HelloKafkaService helloKafkaService) { this.helloKafkaService = helloKafkaService; } @KafkaListener( topics = "my-topic", groupId = "my-topic:HelloKafkaListener") public void process(ExampleEvent event) { this.helloKafkaService.handleMessage(event); log.info("Processing event: " + event.getExampleField()); } }
测试类HelloKafkaListenerTestIT.java
@SpringBootTest @Testcontainers public class HelloKafkaListenerTestIT { @Container public static KafkaContainer kafkaContainer = new KafkaContainer(DockerImageName.parse("confluentinc/cp-kafka:7.0.5")); @Mock private HelloKafkaService helloKafkaService; @Autowired private KafkaTemplate<String, ExampleEvent> kafkaTemplate; @DynamicPropertySource static void kafkaProperties(DynamicPropertyRegistry registry) { registry.add(ProducerConfig.BOOTSTRAP_SERVERS_CONFIG, kafkaContainer::getBootstrapServers); registry.add(ProducerConfig.CLIENT_ID_CONFIG, () -> "test-id"); registry.add(ProducerConfig.KEY_SERIALIZER_CLASS_CONFIG, () -> StringSerializer.class.getName()); registry.add(ProducerConfig.VALUE_SERIALIZER_CLASS_CONFIG, () -> KafkaAvroSerializer.class.getName()); } @BeforeEach public void setUp() { Serializer keySerializer; Serializer valueSerializer; var avroConfig = Map.of( KafkaAvroSerializerConfig.SCHEMA_REGISTRY_URL_CONFIG, "http://localhost:8081", KafkaAvroSerializerConfig.AUTO_REGISTER_SCHEMAS, true, KafkaAvroSerializerConfig.USE_LATEST_VERSION, true ); keySerializer = new StringSerializer(); valueSerializer = new SpecificAvroSerializer<ExampleEvent>(); valueSerializer.configure(avroConfig, false); Map<String, Object> config = new HashMap<>(); config.put(ProducerConfig.BOOTSTRAP_SERVERS_CONFIG, kafkaContainer.getBootstrapServers()); config.put(ProducerConfig.KEY_SERIALIZER_CLASS_CONFIG, StringSerializer.class); config.put(ProducerConfig.VALUE_SERIALIZER_CLASS_CONFIG, KafkaAvroSerializer.class); config.put(ProducerConfig.CLIENT_ID_CONFIG, "test-id"); config.put("schema.registry.url", "http://localhost:8081"); DefaultKafkaProducerFactory<String, ExampleEvent> producerFactory = new DefaultKafkaProducerFactory<>(config, keySerializer, valueSerializer); this.kafkaTemplate = new KafkaTemplate<>(producerFactory); } static { kafkaContainer.start(); } @Test public void helloKafkaListenerTest() { ArgumentCaptor<ExampleEvent> captor = ArgumentCaptor.forClass(ExampleEvent.class); ExampleEvent exampleEvent = new ExampleEvent("Hello World!"); // Confirmed this.kafkaTemplate.send works well this.kafkaTemplate.send("my-topic", "key", exampleEvent); // Throw the error here: verify(helloKafkaService, timeout(5000)).handleMessage(captor.capture()); assertNotNull(captor.getValue()); assertEquals("Hello World!", captor.getValue().getExampleField()); } @AfterAll static void tearDown() { kafkaContainer.stop(); } }
Avro事件定义event.avsc
{ "namespace": "my.namespace", "type": "record", "name": "ExampleEvent", "doc": "A sample event", "fields": [ { "name": "exampleField", "type": "string" } ] }
疑问
为什么使用@Mock和verify时消费者无法触发?有没有更优的测试方案?
一、消费者未触发的核心原因
Mock对象未注入Spring容器
@Mock创建的HelloKafkaService只是本地Mock实例,并未替换Spring容器中的真实服务实例。实际运行的HelloKafkaListener依赖的是容器内的真实服务,自然不会调用本地Mock,导致verify断言失败。缺少Schema Registry容器
代码中配置了Schema Registry地址http://localhost:8081,但仅启动了Kafka容器,没有启动Schema Registry容器。消费者无法反序列化Avro消息,默认配置下会静默失败,导致消息无法被处理。KafkaTemplate初始化冲突
先通过@Autowired注入KafkaTemplate,又在setUp方法中手动创建新实例覆盖原有对象,容易引发配置不一致问题。
二、优化后的测试方案
1. 修复Mock注入问题
使用@MockBean替换@Mock,@MockBean会自动将Mock对象注册到Spring容器,替换原有的HelloKafkaService实例,确保HelloKafkaListener使用该Mock对象。
2. 添加Schema Registry容器
引入Testcontainers的Schema Registry容器,与Kafka容器联动,保证Avro消息的序列化/反序列化正常。
3. 简化Kafka配置
通过@DynamicPropertySource统一配置生产者和消费者的Schema Registry地址,直接使用Spring自动配置的KafkaTemplate实例。
优化后的测试类代码
@SpringBootTest @Testcontainers public class HelloKafkaListenerTestIT { @Container public static KafkaContainer kafkaContainer = new KafkaContainer(DockerImageName.parse("confluentinc/cp-kafka:7.0.5")); @Container public static SchemaRegistryContainer schemaRegistryContainer = new SchemaRegistryContainer(DockerImageName.parse("confluentinc/cp-schema-registry:7.0.5")) .withKafka(kafkaContainer); @MockBean private HelloKafkaService helloKafkaService; @Autowired private KafkaTemplate<String, ExampleEvent> kafkaTemplate; @DynamicPropertySource static void kafkaProperties(DynamicPropertyRegistry registry) { // 生产者配置 registry.add(ProducerConfig.BOOTSTRAP_SERVERS_CONFIG, kafkaContainer::getBootstrapServers); registry.add(ProducerConfig.KEY_SERIALIZER_CLASS_CONFIG, () -> StringSerializer.class.getName()); registry.add(ProducerConfig.VALUE_SERIALIZER_CLASS_CONFIG, () -> KafkaAvroSerializer.class.getName()); registry.add(KafkaAvroSerializerConfig.SCHEMA_REGISTRY_URL_CONFIG, schemaRegistryContainer::getSchemaRegistryUrl); // 消费者配置 registry.add(ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG, kafkaContainer::getBootstrapServers); registry.add(ConsumerConfig.KEY_DESERIALIZER_CLASS_CONFIG, () -> StringDeserializer.class.getName()); registry.add(ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG, () -> KafkaAvroDeserializer.class.getName()); registry.add(KafkaAvroDeserializerConfig.SCHEMA_REGISTRY_URL_CONFIG, schemaRegistryContainer::getSchemaRegistryUrl); registry.add(KafkaAvroDeserializerConfig.SPECIFIC_AVRO_READER_CONFIG, () -> "true"); } @Test public void helloKafkaListenerTest() throws InterruptedException { ArgumentCaptor<ExampleEvent> captor = ArgumentCaptor.forClass(ExampleEvent.class); ExampleEvent exampleEvent = new ExampleEvent("Hello World!"); // 发送消息并等待确认 kafkaTemplate.send("my-topic", "key", exampleEvent).get(3, TimeUnit.SECONDS); // 验证Mock方法被调用 verify(helloKafkaService, timeout(5000).times(1)).handleMessage(captor.capture()); // 断言消息内容正确 assertNotNull(captor.getValue()); assertEquals("Hello World!", captor.getValue().getExampleField()); } }
额外优化建议
- 移除静态代码块中的
kafkaContainer.start()和@AfterAll中的stop(),Testcontainers会自动管理容器的生命周期。 - 发送消息时调用
get()方法等待发送完成,避免消息未写入Kafka就开始验证。 - 配置消费者
SPECIFIC_AVRO_READER_CONFIG=true,确保反序列化为具体的ExampleEvent类型。
内容的提问来源于stack exchange,提问作者Zichen Ma

