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

Spring Boot中Kafka Testcontainers动态覆盖Bootstrap Servers配置问题

解决Kafka Testcontainers集成测试中Bootstrap Servers配置不生效的问题

问题根源分析

  1. Kafka容器的Advertised Listeners配置错误:默认的confluentinc/cp-kafka:latest镜像会把KAFKA_ADVERTISED_LISTENERS设为容器内部的PLAINTEXT://localhost:29092,但Testcontainers映射的是随机外部端口,客户端拿到的元数据是内部端口,导致无法连接到实际启动的Kafka服务。
  2. 手动修改KafkaProperties时机过晚:测试方法中修改kafkaProperties.setBootstrapServers时,Spring已经完成了Kafka生产者/消费者客户端的初始化,此时修改配置不会生效。
  3. 测试主题未自动创建:如果未开启主题自动创建,发送消息时Kafka不会自动生成test-topic,最终导致超时错误。

具体修复步骤

1. 修正Kafka容器的对外暴露地址配置

修改KafkaContainer初始化代码,让容器对外暴露Testcontainers映射的随机外部端口,确保客户端能拿到正确的连接地址:

@Container
private static final KafkaContainer kafkaContainer = new KafkaContainer(DockerImageName.parse("confluentinc/cp-kafka:latest"))
        .withEmbeddedZookeeper()
        // 使用Testcontainers Kafka内置方法配置外部监听地址
        .withListenerAddresses(KafkaContainer.ExternalListener.plaintext());

如果使用旧版本Testcontainers,也可以手动设置环境变量:

@Container
private static final KafkaContainer kafkaContainer = new KafkaContainer(DockerImageName.parse("confluentinc/cp-kafka:latest"))
        .withEmbeddedZookeeper()
        .withEnv("KAFKA_ADVERTISED_LISTENERS", 
                String.format("PLAINTEXT://localhost:%d,PLAINTEXT_INTERNAL://localhost:29092", 
                        kafkaContainer.getMappedPort(9092)))
        .withEnv("KAFKA_LISTENERS", "PLAINTEXT://0.0.0.0:9092,PLAINTEXT_INTERNAL://0.0.0.0:29092");

2. 移除无效的手动配置代码

删除测试方法中的kafkaProperties.setBootstrapServers(List.of(kafkaContainer.getBootstrapServers()));,@DynamicPropertySource已经在Spring上下文初始化前完成了配置覆盖,这个手动修改完全无效。

3. 确保测试主题被创建

可以选择两种方式:

  • 在application-test.yml中开启自动创建主题:
spring:
  kafka:
    admin:
      auto-create: true
    consumer:
      properties:
        allow.auto.create.topics: true
    producer:
      properties:
        allow.auto.create.topics: true
  • 或者在测试类中通过KafkaAdmin手动创建主题:
@Autowired
private KafkaAdmin kafkaAdmin;

@BeforeEach
void setup() {
    NewTopic topic = TopicBuilder.name(TEST_TOPIC).partitions(1).replicas(1).build();
    kafkaAdmin.createOrModifyTopics(topic);
}

4. 确保所有Kafka组件的Bootstrap配置被覆盖

扩展@DynamicPropertySource,覆盖生产者、消费者、Admin所有相关的Bootstrap配置:

@DynamicPropertySource
static void kafkaProperties(DynamicPropertyRegistry registry) {
    String bootstrapServers = kafkaContainer.getBootstrapServers();
    registry.add("spring.kafka.bootstrap-servers", () -> bootstrapServers);
    registry.add("spring.kafka.producer.bootstrap-servers", () -> bootstrapServers);
    registry.add("spring.kafka.consumer.bootstrap-servers", () -> bootstrapServers);
    registry.add("spring.kafka.admin.bootstrap-servers", () -> bootstrapServers);
}

5. 优化测试等待逻辑(可选)

使用同步发送确认消息发送成功,避免因异步发送导致的等待超时:

// 同步发送消息,设置5秒超时
kafkaProducer.sendAvroMessage(actual).get(5, TimeUnit.SECONDS);

完整修复后的测试类示例

@Slf4j
@Testcontainers
@ActiveProfiles({"test"})
@SpringBootTest(classes = KafkaApplication.class)
class ProducerIT {

    private final static String TEST_TOPIC = "test-topic";

    private TestAvroMessage expectedPayload;

    private final CountDownLatch latch = new CountDownLatch(1);

    @Autowired
    private KafkaService kafkaService;

    @Autowired
    private Producer kafkaProducer;

    @Autowired
    private KafkaAdmin kafkaAdmin;

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

    @DynamicPropertySource
    static void kafkaProperties(DynamicPropertyRegistry registry) {
        String bootstrapServers = kafkaContainer.getBootstrapServers();
        registry.add("spring.kafka.bootstrap-servers", () -> bootstrapServers);
        registry.add("spring.kafka.producer.bootstrap-servers", () -> bootstrapServers);
        registry.add("spring.kafka.consumer.bootstrap-servers", () -> bootstrapServers);
        registry.add("spring.kafka.admin.bootstrap-servers", () -> bootstrapServers);
    }

    @BeforeEach
    void setup() {
        // 手动创建测试主题
        NewTopic topic = TopicBuilder.name(TEST_TOPIC).partitions(1).replicas(1).build();
        kafkaAdmin.createOrModifyTopics(topic);
        // 设置生产者目标主题
        ReflectionTestUtils.setField(kafkaService, "producerTopic", TEST_TOPIC);
    }

    @Test
    void verify_that_expected_event_is_successfully_sent() throws ExecutionException, InterruptedException, TimeoutException {
        TestAvroMessage actual = TestAvroMessage.newBuilder()
                .setId("aDemoId")
                .setName("test")
                .setProcessed(true)
                .setCost(125000.00)
                .setLevel(5)
                .build();

        // 同步发送消息,确保发送成功
        kafkaProducer.sendAvroMessage(actual).get(5, TimeUnit.SECONDS);

        // 等待消息被消费
        boolean messageReceived = latch.await(10, TimeUnit.SECONDS);
        assertThat(messageReceived).isTrue();
        assertThat(actual).isEqualTo(expectedPayload);
    }

    @KafkaListener(topics = TEST_TOPIC,
            groupId = "${test.kafka.group-id}",
            containerFactory = "eventListenerContainerFactory")
    void messageListener(TestAvroMessage event,
                         @Header(value = KafkaHeaders.RECEIVED_KEY, required = false) String messageKey,
                         @Header(KafkaHeaders.RECEIVED_PARTITION) String partitionId,
                         @Header(KafkaHeaders.OFFSET) long offset,
                         @Header(KafkaHeaders.RECEIVED_TOPIC) String topic,
                         @Headers Map<String, Object> headers) {
        latch.countDown();
        expectedPayload = event;
    }
}

内容的提问来源于stack exchange,提问作者Tony Langworthy

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.03 05:21:57