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

如何解决@KafkaListener与Mockito测试Kafka Consumer时Mock不生效问题?

问题根源

你猜的没错——测试里用@InjectMocks创建的KafkaConsumerExchangeRate是测试类自己的实例,而真正接收Kafka消息的是Spring容器初始化的另一个consumer实例,它用的是真实的ExchangeRepository,不是你mock的那个,所以mock的repository根本没被调用,自然触发verify失败。

解决步骤

1. 用@MockBean替换@Mock+@InjectMocks

@MockBean会自动把Spring容器里的真实ExchangeRepository替换成mock实例,这样容器里的consumer就会使用这个mock,不需要再手动创建consumer实例。

修改测试类的依赖注入部分:

@SpringBootTest
@DirtiesContext
@EmbeddedKafka(partitions = 1, brokerProperties = {"listeners=PLAINTEXT://localhost:9092", "port=9092"})
class KafkaConsumerExchangeRateTest {

    @Autowired
    private KafkaTemplate<String, ListExchangeRatesDto> template;

    // 用@MockBean替换@Mock和@InjectMocks
    @MockBean
    private ExchangeRepository repository;

    // ... 测试方法逻辑不变(后续可优化等待逻辑)
}

2. 对齐Kafka连接配置

主目录的consumer配置是localhost:29092,但测试用的EmbeddedKafka是9092,需要在测试目录的application.properties里覆盖配置:

spring.kafka.bootstrap-servers=localhost:9092
# 保留原有Producer配置

确保consumer连接到测试用的嵌入式Kafka,而非本地真实Kafka服务。

3. 替换不可靠的Thread.sleep(推荐)

用CountDownLatch等待消息处理完成,比固定时长的sleep更可靠:

@Test
void givenEmbeddedKafkaBroker_whenSendingWithSimpleProducer_thenMessageReceived() throws InterruptedException {
    // 初始化CountDownLatch,计数1表示等待1次消息处理完成
    CountDownLatch latch = new CountDownLatch(1);

    // given部分逻辑不变...

    // 在mock的save方法里触发latch计数减1
    doAnswer(invocation -> {
        latch.countDown();
        return null;
    }).when(repository).save(any());

    // 发送消息
    template.send("exchange-rate", listOfDtos);
    // 等待最多5秒,直到latch计数为0
    latch.await(5, TimeUnit.SECONDS);

    // 验证save方法调用
    verify(repository, times(1)).save(any());
}
修改后的完整测试类示例
@SpringBootTest
@DirtiesContext
@EmbeddedKafka(partitions = 1, brokerProperties = {"listeners=PLAINTEXT://localhost:9092", "port=9092"})
class KafkaConsumerExchangeRateTest {

    @Autowired
    private KafkaTemplate<String, ListExchangeRatesDto> template;

    @MockBean
    private ExchangeRepository repository;

    @Test
    void givenEmbeddedKafkaBroker_whenSendingWithSimpleProducer_thenMessageReceived() throws InterruptedException {
        //given
        ExchangeRate exchangeRate = ExchangeRate.builder()
                .buyingRate(BigDecimal.ONE)
                .sellingRate(BigDecimal.TEN)
                .updatedAt(LocalDateTime.now())
                .currency1("EUR")
                .currency2("USD")
                .build();

        ExchangeOutgoingRateDto dto = ExchangeOutgoingRateDto.builder()
                .buyingRate(BigDecimal.ONE)
                .sellingRate(BigDecimal.TEN)
                .updateAt(LocalDateTime.now().toString())
                .currency1("EUR")
                .currency2("USD")
                .build();

        ListExchangeRatesDto listOfDtos = ListExchangeRatesDto.builder()
                .exchangesRates(List.of(dto))
                .build();

        //and
        doReturn(Optional.of(exchangeRate)).when(repository)
                .findExchangeRateByCurrency1AndCurrency2(anyString(), anyString());

        CountDownLatch latch = new CountDownLatch(1);
        doAnswer(invocation -> {
            latch.countDown();
            return null;
        }).when(repository).save(any());

        //when
        template.send("exchange-rate", listOfDtos);
        latch.await(5, TimeUnit.SECONDS);

        //then
        verify(repository, times(1)).save(any());
    }
}

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.29 09:47:06