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

JUnit5中不依赖Spring Boot使用@EmbeddedKafka的可行性及替代方案咨询

在JUnit5中不依赖Spring Boot使用内存Kafka的方案

一、仅基于Spring Core使用@EmbeddedKafka

完全可以不依赖Spring Boot,仅通过Spring Core配合JUnit5使用@EmbeddedKafka注解。该注解属于spring-kafka-test模块,是Spring Kafka生态的一部分,和Spring Boot没有强绑定关系。

只需引入必要依赖(spring-core、spring-context、spring-kafka、spring-kafka-test及JUnit5相关依赖),通过JUnit5的Spring扩展加载Spring上下文即可,无需使用@SpringBootTest。

示例代码:

import org.junit.jupiter.api.Test;
import org.junit.jupiter.api.extension.ExtendWith;
import org.springframework.context.annotation.Bean;
import org.springframework.context.annotation.Configuration;
import org.springframework.kafka.core.KafkaTemplate;
import org.springframework.kafka.core.ProducerFactory;
import org.springframework.kafka.test.context.EmbeddedKafka;
import org.springframework.kafka.test.utils.KafkaTestUtils;
import org.springframework.test.context.ContextConfiguration;
import org.springframework.test.context.junit.jupiter.SpringExtension;

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

@ExtendWith(SpringExtension.class)
@ContextConfiguration(classes = EmbeddedKafkaWithoutSpringBootTest.Config.class)
@EmbeddedKafka(partitions = 1, topics = {"test-topic"})
public class EmbeddedKafkaWithoutSpringBootTest {

    @Configuration
    static class Config {
        @Bean
        public ProducerFactory<String, String> producerFactory() {
            return KafkaTestUtils.producerFactory(EmbeddedKafkaBroker.getSpringEmbeddedKafkaBroker());
        }

        @Bean
        public KafkaTemplate<String, String> kafkaTemplate(ProducerFactory<String, String> producerFactory) {
            return new KafkaTemplate<>(producerFactory);
        }
    }

    @Test
    void testKafkaSendAndReceive(KafkaTemplate<String, String> kafkaTemplate) {
        // 发送消息
        kafkaTemplate.send("test-topic", "key", "test-message");

        // 接收消息验证
        String receivedMessage = KafkaTestUtils.getSingleRecord(
                KafkaTestUtils.consumerFactory(EmbeddedKafkaBroker.getSpringEmbeddedKafkaBroker()),
                "test-topic").value();
        assertEquals("test-message", receivedMessage);
    }
}

二、原生方式启动内存Kafka(无Spring容器)

如果不想依赖Spring上下文,可直接使用spring-kafka-test中的EmbeddedKafkaBroker类手动启动和管理内存Kafka,类似JUnit4的RawKafka实现方式:

import org.apache.kafka.clients.consumer.Consumer;
import org.apache.kafka.clients.consumer.ConsumerRecord;
import org.apache.kafka.clients.producer.KafkaProducer;
import org.apache.kafka.clients.producer.Producer;
import org.apache.kafka.clients.producer.ProducerRecord;
import org.junit.jupiter.api.AfterEach;
import org.junit.jupiter.api.BeforeEach;
import org.junit.jupiter.api.Test;
import org.springframework.kafka.test.EmbeddedKafkaBroker;
import org.springframework.kafka.test.utils.KafkaTestUtils;

import java.util.Map;

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

public class RawEmbeddedKafkaJUnit5Test {

    private EmbeddedKafkaBroker embeddedKafka;
    private Producer<String, String> producer;
    private Consumer<String, String> consumer;

    @BeforeEach
    void setUp() {
        // 启动嵌入式Kafka
        embeddedKafka = new EmbeddedKafkaBroker(1, false, "test-topic");
        embeddedKafka.start();

        // 创建生产者和消费者
        Map<String, Object> producerProps = KafkaTestUtils.producerProps(embeddedKafka);
        producer = new KafkaProducer<>(producerProps);

        Map<String, Object> consumerProps = KafkaTestUtils.consumerProps("test-group", "true", embeddedKafka);
        consumer = KafkaTestUtils.createConsumer(consumerProps);
        consumer.subscribe(java.util.Collections.singletonList("test-topic"));
    }

    @Test
    void testSendReceive() {
        // 发送消息
        producer.send(new ProducerRecord<>("test-topic", "key", "raw-test-message")).join();

        // 接收并验证
        ConsumerRecord<String, String> record = KafkaTestUtils.getSingleRecord(consumer, "test-topic");
        assertEquals("raw-test-message", record.value());
    }

    @AfterEach
    void tearDown() {
        // 关闭资源
        producer.close();
        consumer.close();
        embeddedKafka.stop();
    }
}

三、其他内存Kafka可选方案

1. Testcontainers Kafka

如果需要更贴近生产环境的测试(比如模拟真实Kafka集群行为),可以使用Testcontainers启动Kafka容器,无需依赖嵌入式Kafka:

import org.junit.jupiter.api.Test;
import org.testcontainers.containers.KafkaContainer;
import org.testcontainers.junit.jupiter.Container;
import org.testcontainers.junit.jupiter.Testcontainers;
import org.testcontainers.utility.DockerImageName;

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

@Testcontainers
public class TestcontainersKafkaTest {

    @Container
    private static final KafkaContainer kafkaContainer = new KafkaContainer(DockerImageName.parse("confluentinc/cp-kafka:latest"));

    @Test
    void testKafkaContainer() {
        // 获取Kafka bootstrap地址
        String bootstrapServers = kafkaContainer.getBootstrapServers();
        assertNotNull(bootstrapServers);
        // 此处可基于该地址创建生产者/消费者进行测试
    }
}

2. 第三方内存Kafka库

比如kafka-unit这类第三方库,但这类库大多维护不活跃,推荐优先使用官方的EmbeddedKafkaBroker或Testcontainers方案。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.12 18:46:07