如何使用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
相关产品推荐
相关产品推荐

