使用固定端口的EmbeddedKafkaBroker集成测试问题咨询
在Spring Boot 3.3.2项目中,未使用Spring-Kafka模板,直接依赖Apache kafka-clients包,需要用EmbeddedKafkaBroker做集成测试。相关代码如下:
服务与配置类
@Service public class MyPublishService { @Autowired MyEventsProducer eventProducer; // 基于KafkaProducer和AdminClient的封装类 @PostConstruct void init() { eventProducer.createTopic("myTopic"); // 调用AdminClient创建主题 } public void publish(String topic, Object payload) { // 向Kafka发送消息 eventProducer.send(topic, payload); } } @Configuration class MyKafkaConfig { String server; // 集成测试时期望为localhost:9092 @Bean public MyEventsProducer myEventsProducer(){ if(localEnv()) { // 判断是否为本地Kafka环境 server = "localhost:9092"; } return new MyEventsProducer(server); } }
集成测试代码
@ExtendWith(SpringExtension.class) @SpringBootTest( classes = Application.class, webEnvironment = SpringBootTest.WebEnvironment.RANDOM_PORT ) @ActiveProfiles("test") @AutoConfigureMockMvc @EmbeddedKafka(partitions = 1, brokerProperties = { "listeners=PLAINTEXT://localhost:9092", "port=9092" }) public class MyPublishServiceTest{ @Autowired MyPublishService myPublishService; @Autowired EmbeddedKafkaBroker embeddedKafkaBroker; @Test public void publishMessageTest() { System.out.println(embeddedKafkaBroker.getBrokersAsString()); myPublishService.publishCollaboration(collaboration); } }
待解决问题
- 为何
embeddedKafkaBroker.getBrokersAsString()始终返回随机端口?已尝试在application-test.yaml中配置spring.kafka.bootstrap-servers: localhost:9092但无效,如何获取固定的host:port? - 在
MyPublishService的@PostConstruct中创建主题,如何确保Embedded Kafka在执行该代码前启动,避免出现TimeOutException?
(注:使用Spring-kafka-test 3.1.1、org.testcontainers:kafka 1.19.3)
问题1:获取固定端口的EmbeddedKafkaBroker地址
你配置的@EmbeddedKafka参数虽然指定了port=9092,但getBrokersAsString()返回随机端口的核心原因是:EmbeddedKafkaBroker默认会忽略硬编码的端口配置,除非显式关闭端口随机分配逻辑。
具体修复步骤
- 在
@EmbeddedKafka注解中添加controlledShutdown = true,确保固定端口配置生效 - 手动覆盖
MyKafkaConfig中的server地址,让生产者使用EmbeddedKafka的实际地址
修改后的测试类示例:
@EmbeddedKafka( partitions = 1, controlledShutdown = true, brokerProperties = { "listeners=PLAINTEXT://localhost:9092", "port=9092" } ) public class MyPublishServiceTest{ @Autowired EmbeddedKafkaBroker embeddedKafkaBroker; @Autowired MyKafkaConfig myKafkaConfig; @BeforeEach void setUp() { // 手动覆盖配置中的server地址,确保生产者连接到EmbeddedKafka myKafkaConfig.server = embeddedKafkaBroker.getBrokersAsString(); // 如果MyEventsProducer是单例,可重新初始化实例绑定新地址 // myEventsProducer = new MyEventsProducer(embeddedKafkaBroker.getBrokersAsString()); } // ... 测试方法 }
另外需检查localEnv()方法的判断逻辑,确保test环境下该条件为true,否则server地址不会被设置为localhost:9092。
问题2:确保EmbeddedKafka在@PostConstruct前启动
@PostConstruct会在Spring Bean初始化完成后立即执行,而EmbeddedKafkaBroker的启动时机可能晚于MyPublishService的初始化,导致主题创建超时。提供三种可行方案:
方案1:用@DependsOn强制依赖顺序
修改MyPublishService的注解,添加@DependsOn("embeddedKafkaBroker"),强制Spring先初始化EmbeddedKafkaBroker,再初始化MyPublishService:
@Service @DependsOn("embeddedKafkaBroker") public class MyPublishService { // ... 原有代码 }
方案2:延迟主题创建到首次发送消息时
去掉@PostConstruct中的主题创建逻辑,改为在publish方法中先检查主题是否存在,不存在再创建:
public void publish(String topic, Object payload) { if (!eventProducer.topicExists(topic)) { eventProducer.createTopic(topic); } eventProducer.send(topic, payload); }
方案3:在测试类中提前启动EmbeddedKafka
在测试类中手动启动EmbeddedKafkaBroker,确保在Spring上下文初始化前完成启动:
@ExtendWith(SpringExtension.class) @SpringBootTest(...) @ActiveProfiles("test") @AutoConfigureMockMvc @EmbeddedKafka(...) public class MyPublishServiceTest{ @Autowired static EmbeddedKafkaBroker embeddedKafkaBroker; @BeforeAll static void startKafka() { embeddedKafkaBroker.start(); } // ... 测试方法 }
注意:此方案需要EmbeddedKafkaBroker被静态注入,且@BeforeAll方法为静态方法。
内容的提问来源于stack exchange,提问作者user2206366

