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

使用固定端口的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);
    }
}

待解决问题

  1. 为何embeddedKafkaBroker.getBrokersAsString()始终返回随机端口?已尝试在application-test.yaml中配置spring.kafka.bootstrap-servers: localhost:9092但无效,如何获取固定的host:port?
  2. 在MyPublishService的@PostConstruct中创建主题,如何确保Embedded Kafka在执行该代码前启动,避免出现TimeOutException?

(注:使用Spring-kafka-test 3.1.1、org.testcontainers:kafka 1.19.3)


解决方案

问题1:获取固定端口的EmbeddedKafkaBroker地址

你配置的@EmbeddedKafka参数虽然指定了port=9092,但getBrokersAsString()返回随机端口的核心原因是:EmbeddedKafkaBroker默认会忽略硬编码的端口配置,除非显式关闭端口随机分配逻辑。

具体修复步骤

  1. 在@EmbeddedKafka注解中添加controlledShutdown = true,确保固定端口配置生效
  2. 手动覆盖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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.13 22:47:06