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

如何使用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的核心原因有两个:

  1. Kafka的监听配置错误,容器间通信应该用内部监听端口而非外部端口
  2. 容器启动顺序依赖处理不当,@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,提问作者박수민

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.12 23:39:58