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

如何使用Mockito测试KafkaTemplate 解决测试中实例为空、无法获取返回值问题

测试时kafkaTemplate为null是因为测试运行过程中没有完成该对象的注入,要么未启动Spring测试上下文完成自动装配,要么未手动注入mock实例到被测类中,可通过以下两种方案实现测试:

方案1:单元测试(使用Mockito mock KafkaTemplate,执行速度快,无需依赖Kafka环境)

如果原类使用字段注入(即你当前@Autowired打在字段上的写法),可以通过反射工具注入依赖,也可以优化为构造注入更便于测试。

不修改原业务代码的测试写法

// JUnit5 版本
import org.junit.jupiter.api.BeforeEach;
import org.junit.jupiter.api.Test;
import org.junit.jupiter.api.extension.ExtendWith;
import org.mockito.Mock;
import org.mockito.junit.jupiter.MockitoExtension;
import org.springframework.kafka.core.KafkaTemplate;
import org.springframework.kafka.support.SendResult;
import org.springframework.test.util.ReflectionTestUtils;
import org.springframework.util.concurrent.ListenableFuture;
import org.springframework.util.concurrent.SettableListenableFuture;

import static org.mockito.ArgumentMatchers.any;
import static org.mockito.ArgumentMatchers.eq;
import static org.mockito.Mockito.*;

@ExtendWith(MockitoExtension.class)
public class KafkaProducerServiceTest {
    @Mock
    private KafkaTemplate<String, ShopResponse> kafkaTemplate;

    private KafkaProducerService kafkaProducerService;
    private final String testTopic = "test-order-topic";

    @BeforeEach
    void setUp() {
        kafkaProducerService = new KafkaProducerService();
        // 手动注入字段依赖
        ReflectionTestUtils.setField(kafkaProducerService, "kafkaTemplate", kafkaTemplate);
        ReflectionTestUtils.setField(kafkaProducerService, "kafkaOrderResponseTopic", testTopic);
    }

    @Test
    void testSendMessage() {
        // 构造mock的返回future
        SettableListenableFuture<SendResult<String, ShopResponse>> mockFuture = new SettableListenableFuture<>();
        SendResult<String, ShopResponse> mockResult = mock(SendResult.class);
        mockFuture.set(mockResult);
        //  stub send方法的返回值
        when(kafkaTemplate.send(eq(testTopic), any(ShopResponse.class))).thenReturn(mockFuture);

        // 调用被测方法
        OrderResponseHeader testHeader = new OrderResponseHeader();
        // 这里可按需给testHeader设置测试字段
        kafkaProducerService.sendOrderResponseToKafka(testHeader);

        // 验证send方法被正确调用,同时你可以直接使用之前构造的mockFuture获取返回值
        verify(kafkaTemplate, times(1)).send(eq(testTopic), any(ShopResponse.class));
        // 如果你想直接拿到业务方法内的future,建议修改业务方法将future作为返回值返回
    }
}

优化业务代码为构造注入(更推荐)

修改后无需使用反射工具,代码可维护性更高:

// 优化后的业务类
public class KafkaProducerService  {
    private final KafkaTemplate<String, ShopResponse> kafkaTemplate;
    private final String kafkaOrderResponseTopic;

    public KafkaProducerService(KafkaTemplate<String, ShopResponse> kafkaTemplate,
                               @Value("${kafka.order.response.topic}") String kafkaOrderResponseTopic) {
        this.kafkaTemplate = kafkaTemplate;
        this.kafkaOrderResponseTopic = kafkaOrderResponseTopic;
    }

    // 建议增加返回值,方便调用方获取发送结果
    public ListenableFuture<SendResult<String, ShopResponse>> sendOrderResponseToKafka(OrderResponseHeader orderResponseHeader) {
        ShopResponse shopResponse = createShopResponse(orderResponseHeader);
        return kafkaTemplate.send(kafkaOrderResponseTopic, shopResponse);
    }
    // 省略createShopResponse实现
}

方案2:集成测试(使用嵌入式Kafka验证真实发送流程)

如果需要验证真实的消息发送、序列化逻辑,可以使用Spring提供的嵌入式Kafka做集成测试,此时Spring上下文会自动注入所有依赖,不会出现kafkaTemplate为null的问题。

依赖引入(Maven)

<dependency>
    <groupId>org.springframework.kafka</groupId>
    <artifactId>spring-kafka-test</artifactId>
    <scope>test</scope>
</dependency>

测试代码示例

import org.junit.jupiter.api.Test;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.beans.factory.annotation.Value;
import org.springframework.boot.test.context.SpringBootTest;
import org.springframework.kafka.test.context.EmbeddedKafka;
import org.springframework.kafka.support.SendResult;
import org.springframework.util.concurrent.ListenableFuture;

import static org.junit.jupiter.api.Assertions.assertEquals;

@SpringBootTest
@EmbeddedKafka(partitions = 1, brokerProperties = {"listeners=PLAINTEXT://localhost:9092"})
public class KafkaProducerIntegrationTest {
    @Autowired
    private KafkaProducerService kafkaProducerService;

    @Value("${kafka.order.response.topic}")
    private String testTopic;

    @Test
    void testRealSend() throws Exception {
        OrderResponseHeader testHeader = new OrderResponseHeader();
        // 调用方法获取future
        ListenableFuture<SendResult<String, ShopResponse>> future = kafkaProducerService.sendOrderResponseToKafka(testHeader);
        // 等待发送完成获取结果
        SendResult<String, ShopResponse> sendResult = future.get();
        // 验证发送结果
        assertEquals(testTopic, sendResult.getRecordMetadata().topic());
    }
}

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.09.28 08:57:02