如何使用Testcontainers配置Confluent Kafka与Schema Registry
Testcontainers配置Confluent Kafka+Schema Registry的问题解决
我是Testcontainers新手,想要配置带Schema Registry的Confluent Kafka,当前配置代码如下:
@TestConfiguration(proxyBeanMethods = false) class KafkaContainerConfig { companion object { val network = Network.newNetwork() } @Bean fun kafkaContainer(): ConfluentKafkaContainer { val kafkaContainer = ConfluentKafkaContainer("confluentinc/cp-kafka:7.5.0") .apply { withNetwork(network) withNetworkAliases("kafka") withEnv("KAFKA_CFG_NODE_ID", "1") withEnv("KAFKA_CFG_PROCESS_ROLES", "controller,broker") withEnv("KAFKA_CFG_LISTENERS", "INTERNAL://:9091,CONTROLLER://:9093,EXTERNAL://:9092") withEnv("KAFKA_CFG_ADVERTISED_LISTENERS", "INTERNAL://kafka:9091,EXTERNAL://kafka:9092") withEnv("KAFKA_CFG_LISTENER_SECURITY_PROTOCOL_MAP", "INTERNAL:PLAINTEXT,CONTROLLER:PLAINTEXT,EXTERNAL:PLAINTEXT") withEnv("KAFKA_CFG_CONTROLLER_LISTENER_NAMES", "CONTROLLER") withEnv("KAFKA_CFG_TRANSACTION_STATE_LOG_REPLICATION_FACTOR", "1") withEnv("KAFKA_CFG_TRANSACTION_STATE_LOG_ISR", "1") withEnv("KAFKA_CFG_OFFSETS_TOPIC_REPLICATION_FACTOR", "1") withEnv("KAFKA_CFG_INTER_BROKER_LISTENER_NAME", "CONTROLLER") withEnv("KAFKA_CFG_CONTROLLER_QUORUM_VOTERS", "1@kafka:9093") withExposedPorts(9092, 9093, 9091) } return kafkaContainer } @Bean @DependsOn("kafkaContainer") fun schemaRegistryContainer(kafkaContainer: ConfluentKafkaContainer): GenericContainer<Nothing> { val schemaRegistryContainer = GenericContainer<Nothing>("confluentinc/cp-schema-registry:7.5.0") .apply { withExposedPorts(8085) withNetwork(network) withNetworkAliases("schema-registry") withEnv("SCHEMA_REGISTRY_HOST_NAME", "schema-registry") withEnv("SCHEMA_REGISTRY_CUB_KAFKA_MIN_BROKERS", "1") withEnv("SCHEMA_REGISTRY_KAFKASTORE_BOOTSTRAP_SERVERS", "kafka:9092") withEnv("SCHEMA_REGISTRY_LISTENERS", "http://0.0.0.0:8085") } return schemaRegistryContainer } @Bean fun setKafkaProperties(kafkaContainer: ConfluentKafkaContainer, schemaRegistryContainer: GenericContainer<Nothing>): DynamicPropertyRegistrar { return DynamicPropertyRegistrar { registry -> registry.add("spring.kafka.properties.bootstrap-servers") { kafkaContainer.bootstrapServers } registry.add("kafka.schema-registry-url") { "localhost:${schemaRegistryContainer.exposedPorts[0]}" } registry.add("spring.kafka.properties.schema.registry.url") { "localhost:${schemaRegistryContainer.exposedPorts[0]}" } } } }
运行测试时,Schema Registry容器抛出如下错误:
2025-06-04 17:33:03 [2025-06-04 08:33:03,440] WARN Couldn't resolve server kafka:9092 from bootstrap.servers as DNS resolution failed for kafka (org.apache.kafka.clients.ClientUtils) 2025-06-04 17:33:03 [2025-06-04 08:33:03,441] ERROR Error while running kafka-ready. (io.confluent.admin.utils.cli.KafkaReadyCommand) 2025-06-04 17:33:03 org.apache.kafka.common.KafkaException: Failed to create new KafkaAdminClient 2025-06-04 17:33:03 at org.apache.kafka.clients.admin.KafkaAdminClient.createInternal(KafkaAdminClient.java:551) 2025-06-04 17:33:03 at org.apache.kafka.clients.admin.Admin.create(Admin.java:144) 2025-06-04 17:33:03 at org.apache.kafka.clients.admin.AdminClient.create(AdminClient.java:49) 2025-06-04 17:33:03 at io.confluent.admin.utils.ClusterStatus.isKafkaReady(ClusterStatus.java:136) 2025-06-04 17:33:03 at io.confluent.admin.utils.cli.KafkaReadyCommand.main(KafkaReadyCommand.java:149) 2025-06-04 17:33:03 Caused by: org.apache.kafka.common.config.ConfigException: No resolvable bootstrap urls given in bootstrap.servers 2025-06-04 17:33:03 at org.apache.kafka.clients.ClientUtils.parseAndValidateAddresses(ClientUtils.java:89) 2025-06-04 17:33:03 at org.apache.kafka.clients.ClientUtils.parseAndValidateAddresses(ClientUtils.java:48) 2025-06-04 17:33:03 at org.apache.kafka.clients.admin.KafkaAdminClient.createInternal(KafkaAdminClient.java:508) 2025-06-04 17:33:03 ... 4 more 2025-06-04 17:33:03 Using log4j config /etc/schema-registry/log4j.properties
问题分析
你的判断正确,Schema Registry无法解析kafka:9092的核心原因有两个:
- Kafka的监听配置错误,容器间通信应该用内部监听端口而非外部端口
- 容器启动顺序依赖处理不当,
@DependsOn无法确保Kafka完全就绪后再启动Schema Registry
正确配置方案
以下是修正后的完整配置代码:
@TestConfiguration(proxyBeanMethods = false) class KafkaContainerConfig { companion object { val network = Network.newNetwork() } @Bean fun kafkaContainer(): ConfluentKafkaContainer { return ConfluentKafkaContainer("confluentinc/cp-kafka:7.5.0") .apply { withNetwork(network) withNetworkAliases("kafka") // 分离内部/外部监听:容器间用INTERNAL,测试代码用EXTERNAL withEnv("KAFKA_CFG_LISTENERS", "INTERNAL://:9091,CONTROLLER://:9093,EXTERNAL://0.0.0.0:9092") withEnv("KAFKA_CFG_ADVERTISED_LISTENERS", "INTERNAL://kafka:9091,EXTERNAL://localhost:${getMappedPort(9092)}") withEnv("KAFKA_CFG_LISTENER_SECURITY_PROTOCOL_MAP", "INTERNAL:PLAINTEXT,CONTROLLER:PLAINTEXT,EXTERNAL:PLAINTEXT") withEnv("KAFKA_CFG_CONTROLLER_LISTENER_NAMES", "CONTROLLER") withEnv("KAFKA_CFG_INTER_BROKER_LISTENER_NAME", "INTERNAL") // 修正为内部监听,确保broker间通信正常 withEnv("KAFKA_CFG_NODE_ID", "1") withEnv("KAFKA_CFG_PROCESS_ROLES", "controller,broker") withEnv("KAFKA_CFG_CONTROLLER_QUORUM_VOTERS", "1@kafka:9093") withEnv("KAFKA_CFG_TRANSACTION_STATE_LOG_REPLICATION_FACTOR", "1") withEnv("KAFKA_CFG_TRANSACTION_STATE_LOG_ISR", "1") withEnv("KAFKA_CFG_OFFSETS_TOPIC_REPLICATION_FACTOR", "1") } } @Bean fun schemaRegistryContainer(kafkaContainer: ConfluentKafkaContainer): GenericContainer<Nothing> { return GenericContainer<Nothing>("confluentinc/cp-schema-registry:7.5.0") .apply { withNetwork(network) withNetworkAliases("schema-registry") withExposedPorts(8085) withEnv("SCHEMA_REGISTRY_HOST_NAME", "schema-registry") withEnv("SCHEMA_REGISTRY_LISTENERS", "http://0.0.0.0:8085") // 连接Kafka内部监听端口,容器网络内可直接解析kafka别名 withEnv("SCHEMA_REGISTRY_KAFKASTORE_BOOTSTRAP_SERVERS", "kafka:9091") withEnv("SCHEMA_REGISTRY_CUB_KAFKA_MIN_BROKERS", "1") // 依赖Kafka容器,Testcontainers会等待Kafka健康检查通过后再启动 dependsOn(kafkaContainer) } } @Bean fun kafkaProperties(kafkaContainer: ConfluentKafkaContainer, schemaRegistryContainer: GenericContainer<Nothing>): DynamicPropertyRegistrar { return DynamicPropertyRegistrar { registry -> // Testcontainers自动返回外部可访问的bootstrap地址 registry.add("spring.kafka.bootstrap-servers") { kafkaContainer.bootstrapServers } // 获取实际映射的端口,避免随机端口冲突 val schemaRegistryPort = schemaRegistryContainer.getMappedPort(8085) registry.add("spring.kafka.properties.schema.registry.url") { "http://localhost:$schemaRegistryPort" } registry.add("kafka.schema-registry-url") { "http://localhost:$schemaRegistryPort" } } } }
关键修正点说明
- Kafka监听配置:容器间通信使用
INTERNAL://kafka:9091,测试代码通过bootstrapServers获取Testcontainers自动映射的外部地址 - Schema Registry连接地址:改为Kafka的内部端口
9091,容器网络内可直接解析kafka别名 - 容器依赖:使用
dependsOn(kafkaContainer)替代@DependsOn,确保Schema Registry等待Kafka完全就绪后启动 - 端口处理:用
getMappedPort获取实际映射端口,避免固定端口冲突问题
内容的提问来源于stack exchange,提问作者박수민
相关产品推荐
相关产品推荐

