Spring Boot中Kafka Testcontainers动态覆盖Bootstrap Servers配置问题
解决Kafka Testcontainers集成测试中Bootstrap Servers配置不生效的问题
问题根源分析
- Kafka容器的Advertised Listeners配置错误:默认的
confluentinc/cp-kafka:latest镜像会把KAFKA_ADVERTISED_LISTENERS设为容器内部的PLAINTEXT://localhost:29092,但Testcontainers映射的是随机外部端口,客户端拿到的元数据是内部端口,导致无法连接到实际启动的Kafka服务。 - 手动修改KafkaProperties时机过晚:测试方法中修改
kafkaProperties.setBootstrapServers时,Spring已经完成了Kafka生产者/消费者客户端的初始化,此时修改配置不会生效。 - 测试主题未自动创建:如果未开启主题自动创建,发送消息时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
相关产品推荐
相关产品推荐

