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
相关产品推荐
相关产品推荐

