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

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时消费者无法触发?有没有更优的测试方案?


问题原因及解决方案

一、消费者未触发的核心原因

  1. Mock对象未注入Spring容器
    @Mock创建的HelloKafkaService只是本地Mock实例,并未替换Spring容器中的真实服务实例。实际运行的HelloKafkaListener依赖的是容器内的真实服务,自然不会调用本地Mock,导致verify断言失败。

  2. 缺少Schema Registry容器
    代码中配置了Schema Registry地址http://localhost:8081,但仅启动了Kafka容器,没有启动Schema Registry容器。消费者无法反序列化Avro消息,默认配置下会静默失败,导致消息无法被处理。

  3. 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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.30 07:08:14