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

如何使用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任何组件,流程如下:

  1. 确保依赖正确引入,spring-kafka-test的版本要和项目中spring-kafka的版本完全一致
  2. 测试类添加@SpringBootTest和@EmbeddedKafka注解,启动嵌入式Kafka服务
  3. 用@DynamicPropertySource动态将嵌入式Kafka的地址注入Spring配置,覆盖默认的Broker地址
  4. 注入待测试的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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.09.23 16:15:01