如何使用Embedded Kafka实现Spring Boot Kafka生产者无Mock真实类测试
无Mock的Spring Kafka生产者集成测试实现方案
报错原因说明
你遇到的Broker不可用报错,核心原因是Embedded Kafka启动时会使用随机端口,而你的Spring Boot应用默认的spring.kafka.bootstrap-servers配置还是指向固定地址(比如默认的localhost:9092),没有和Embedded Kafka的实际监听地址对齐,导致生产者连接失败。
实现步骤
我们用Spring官方提供的@EmbeddedKafka注解(比EmbeddedKafkaRule更适配Spring Boot测试场景)完成全真实类的集成测试,不需要Mock任何组件,流程如下:
- 确保依赖正确引入,
spring-kafka-test的版本要和项目中spring-kafka的版本完全一致 - 测试类添加
@SpringBootTest和@EmbeddedKafka注解,启动嵌入式Kafka服务 - 用
@DynamicPropertySource动态将嵌入式Kafka的地址注入Spring配置,覆盖默认的Broker地址 - 注入待测试的
MessageProducer和你的业务@KafkaListener消费者,用同步工具(比如CountDownLatch)验证消息发送和消费的完整流程
完整测试代码示例
首先修正你示例代码中KafkaTemplate泛型的多余引号问题,业务代码如下
@Component public class MessageProducer { private KafkaTemplate<String, String> kafkaTemplate; @Autowired public MessageProducer(KafkaTemplate<String, String> kafkaTemplate) { this.kafkaTemplate = kafkaTemplate; } public void sendMessage(String message, String topicName) { kafkaTemplate.send(topicName, message); } }
测试类代码(JUnit 5 + Spring Boot 2.4+)
import org.junit.jupiter.api.Test; import org.springframework.beans.factory.annotation.Autowired; import org.springframework.boot.test.context.SpringBootTest; import org.springframework.kafka.annotation.KafkaListener; import org.springframework.kafka.test.context.EmbeddedKafka; import org.springframework.kafka.test.EmbeddedKafkaBroker; import org.springframework.test.context.DynamicPropertyRegistry; import org.springframework.test.context.DynamicPropertySource; import java.util.concurrent.CountDownLatch; import java.util.concurrent.TimeUnit; import static org.junit.jupiter.api.Assertions.*; @SpringBootTest // 启动1个Broker节点,自动创建test-topic主题,分区数为1 @EmbeddedKafka(partitions = 1, topics = { "test-topic" }) public class MessageProducerTest { @Autowired private MessageProducer messageProducer; @Autowired private EmbeddedKafkaBroker embeddedKafkaBroker; // 如果你已经有业务@KafkaListener,直接注入你的业务消费者类即可,此处为示例测试消费者 @org.springframework.stereotype.Component public static class TestConsumer { // 用于同步异步消费流程 public CountDownLatch latch = new CountDownLatch(1); public String receivedMsg; @KafkaListener(topics = "test-topic", groupId = "test-consumer-group") public void listen(String message) { receivedMsg = message; latch.countDown(); } } @Autowired private TestConsumer testConsumer; // 动态注入Kafka配置,覆盖默认的Broker地址 @DynamicPropertySource static void setKafkaProperties(DynamicPropertyRegistry registry) { registry.add("spring.kafka.bootstrap-servers", embeddedKafkaBroker::getBrokersAsString); // 消费者配置:首次消费从最早的消息开始读,避免消息发送早于消费者启动导致丢失 registry.add("spring.kafka.consumer.auto-offset-reset", () -> "earliest"); } @Test void testSendMessageSuccess() throws InterruptedException { String testMsg = "test embedded kafka msg"; String testTopic = "test-topic"; // 调用生产者发送消息 messageProducer.sendMessage(testMsg, testTopic); // 等待消费完成,最多等待5秒 boolean isReceived = testConsumer.latch.await(5, TimeUnit.SECONDS); // 验证结果 assertTrue(isReceived, "消息未在规定时间内被消费者接收"); assertEquals(testMsg, testConsumer.receivedMsg); } }
如果你使用JUnit 4,需要调整为EmbeddedKafkaRule的写法:
import org.junit.ClassRule; import org.junit.BeforeClass; import org.junit.Test; import org.springframework.kafka.test.rule.EmbeddedKafkaRule; // 其他导入和上面一致 @SpringBootTest public class MessageProducerTest { @ClassRule public static EmbeddedKafkaRule embeddedKafkaRule = new EmbeddedKafkaRule(1, true, "test-topic"); @BeforeClass public static void setKafkaConfig() { // 手动将Embedded Kafka地址设置到系统属性,覆盖Spring配置 System.setProperty("spring.kafka.bootstrap-servers", embeddedKafkaRule.getEmbeddedKafka().getBrokersAsString()); System.setProperty("spring.kafka.consumer.auto-offset-reset", "earliest"); } // 剩下的测试逻辑和JUnit5版本完全一致 }
注意事项
- 测试类不要硬编码Kafka的Broker地址,全部从EmbeddedKafka的实例中动态获取
- 消费是异步流程,必须用
CountDownLatch/CompletableFuture等同步工具等待消费完成,不能直接断言 - 如果不需要加载全量Spring上下文,可以在
@SpringBootTest中指定classes参数,只加载需要的Bean(生产者、消费者、Kafka自动配置类),加快测试速度
内容的提问来源于stack exchange,提问作者ppb
相关产品推荐
相关产品推荐

