如何在JUnit测试中于Spring Boot应用启动后启动嵌入式Kafka集群
JUnit测试延迟启动嵌入式Kafka实现方案
核心原理
@EmbeddedKafka注解默认会在Spring应用上下文初始化前启动Kafka集群,要调整启动时机,只要放弃注解的自动托管逻辑,手动控制EmbeddedKafkaBroker实例的生命周期即可。
实现步骤
- 移除测试类上的
@EmbeddedKafka注解,改为手动创建EmbeddedKafkaBroker实例,不由Spring自动触发启动 - 禁用默认Kafka健康检查和调整客户端重试参数,避免应用启动阶段因Kafka未就绪直接报错:
# 测试配置文件或@TestPropertySource中添加 management.health.kafka.enabled=false # 配置生产者/消费者重试参数,按需调整时长 spring.kafka.producer.retries=20 spring.kafka.producer.properties.retry.backoff.ms=1000 spring.kafka.consumer.auto-offset-reset=earliest
- 在测试逻辑中先执行等待逻辑,再手动启动Kafka并将服务地址注册到Spring环境
完整代码示例(JUnit 5)
import org.springframework.kafka.test.EmbeddedKafkaBroker; import org.springframework.kafka.core.KafkaTemplate; import org.springframework.context.ConfigurableApplicationContext; import org.springframework.boot.test.context.SpringBootTest; import org.springframework.test.context.TestPropertySource; import org.junit.jupiter.api.Test; import org.springframework.beans.factory.annotation.Autowired; @SpringBootTest @TestPropertySource(properties = { "management.health.kafka.enabled=false", "spring.kafka.bootstrap-servers=${spring.embedded.kafka.brokers}" }) public class KafkaDelayStartTest { // 手动创建嵌入式Kafka实例,参数依次为:broker数量、是否受控关闭、分区数、预创建的Topic名 private final EmbeddedKafkaBroker embeddedKafka = new EmbeddedKafkaBroker(1, false, 1, "test_topic"); @Autowired private ConfigurableApplicationContext context; @Autowired private KafkaTemplate<String, String> kafkaTemplate; @Test public void testAppHandleWhenKafkaNotReady() throws InterruptedException { // 模拟应用启动后延迟启动Kafka的等待时长,示例为30秒,按需调整 Thread.sleep(30000); // 手动启动嵌入式Kafka集群 embeddedKafka.start(); // 将Kafka服务地址注册到Spring环境中,让Kafka客户端可以获取到连接地址 context.getEnvironment().getSystemProperties().put( "spring.embedded.kafka.brokers", embeddedKafka.getBrokersAsString() ); // 等待Kafka集群完全就绪,按需调整等待时长 Thread.sleep(5000); // 此处编写测试逻辑,例如发送消息、验证消费逻辑、异常处理逻辑等 kafkaTemplate.send("test_topic", "test_msg"); // 测试完成后手动关闭Kafka集群 embeddedKafka.stop(); } }
注意事项
- 如果应用中存在启动时自动运行的Kafka消费监听,Kafka启动完成后可以调用
KafkaListenerEndpointRegistry的start()方法手动重启监听容器,加快客户端重连效率 - 多测试用例复用延迟启动逻辑时,可以将Kafka启停逻辑封装到
@BeforeEach/@AfterEach方法或者自定义测试扩展中,减少重复代码
内容的提问来源于stack exchange,提问作者TinekeW
相关产品推荐
相关产品推荐

