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

使用Testcontainers启动Kafka与Schema Registry失败求助(Kotlin)

问题分析与解决方案

你遇到的Schema Registry连接Kafka超时问题,核心原因是容器间网络通信配置错误,加上缺少必要的启动等待策略。以下是具体修正方案:

核心问题点

  1. 容器间通信不能使用主机映射的Kafka地址(kafkaContainer.bootstrapServers),必须用容器网络内的内部地址
  2. Schema Registry的监听地址配置错误,导致容器外部无法访问
  3. 未设置Schema Registry的启动等待策略,过早尝试连接未就绪的服务
  4. 动态属性中Schema Registry端口的获取方式错误

修正后的完整配置

@Testcontainers
companion object {

    @Container
    val kafkaContainer = KafkaContainer(DockerImageName.parse("confluentinc/cp-kafka:latest")).apply {
        setWaitStrategy(
            Wait.defaultWaitStrategy()
                .withStartupTimeout(Duration.ofSeconds(75))
        )
        withEmbeddedZookeeper()
            .withEnv("KAFKA_LISTENERS", "PLAINTEXT://0.0.0.0:9093,BROKER://0.0.0.0:9092")
            .withEnv("KAFKA_LISTENER_SECURITY_PROTOCOL_MAP", "BROKER:PLAINTEXT,PLAINTEXT:PLAINTEXT")
            .withEnv("KAFKA_INTER_BROKER_LISTENER_NAME", "BROKER")
            .withEnv("KAFKA_BROKER_ID", "1")
            .withEnv("KAFKA_OFFSETS_TOPIC_REPLICATION_FACTOR", "1")
            .withEnv("KAFKA_OFFSETS_TOPIC_NUM_PARTITIONS", "1")
            .withEnv("KAFKA_TRANSACTION_STATE_LOG_REPLICATION_FACTOR", "1")
            .withEnv("KAFKA_TRANSACTION_STATE_LOG_MIN_ISR", "1")
            .withEnv("KAFKA_LOG_FLUSH_INTERVAL_MESSAGES", Long.MAX_VALUE.toString())
            .withEnv("KAFKA_GROUP_INITIAL_REBALANCE_DELAY_MS", "0")
    }.withReuse(true)

    @Container
    val schemaRegistryContainer = GenericContainer(DockerImageName.parse("confluentinc/cp-schema-registry")).apply {
        withExposedPorts(8081)
        withEnv("SCHEMA_REGISTRY_HOST_NAME", "schema-registry")
        // 允许容器外部访问服务
        withEnv("SCHEMA_REGISTRY_LISTENERS", "http://0.0.0.0:8081")
        // 使用容器网络内的Kafka内部地址(容器默认名称为kafka,对应内部端口9092)
        withEnv("SCHEMA_REGISTRY_KAFKASTORE_BOOTSTRAP_SERVERS", "PLAINTEXT://kafka:9092")
        withNetwork(kafkaContainer.network)
        // 等待Schema Registry服务就绪后再继续
        setWaitStrategy(
            Wait.forHttp("/subjects")
                .forPort(8081)
                .withStartupTimeout(Duration.ofSeconds(60))
        )
    }

    @DynamicPropertySource
    @JvmStatic
    fun kafkaProperties(registry: DynamicPropertyRegistry) {
        registry.add("kafka.bootstrapServers") { kafkaContainer.bootstrapServers }
        // 获取主机映射的实际端口,而非容器内部端口集合
        registry.add("kafka.schemaRegistryUrl") { "http://localhost:${schemaRegistryContainer.getFirstMappedPort()}" }
    }
}

关键修改说明

  1. 移除手动start()调用:Testcontainers会通过@Container注解自动管理容器启动顺序(先启动Kafka,再启动Schema Registry),手动调用会打乱依赖逻辑
  2. 修正Kafka连接地址:容器在同一网络内时,Kafka容器默认名称为kafka,使用PLAINTEXT://kafka:9092作为内部通信地址,而非主机映射的端口
  3. 调整监听地址:将SCHEMA_REGISTRY_LISTENERS改为http://0.0.0.0:8081,确保容器外部(测试代码)能访问服务
  4. 添加等待策略:通过Wait.forHttp("/subjects")等待Schema Registry服务完全就绪,避免启动后立即访问失败
  5. 修正端口获取方式:用getFirstMappedPort()获取主机实际映射的端口,而非直接使用exposedPorts容器内部端口集合

额外注意事项

  • 建议使用同版本的Confluent镜像(如Kafka用cp-kafka:7.5.0,Schema Registry对应cp-schema-registry:7.5.0),避免版本兼容问题
  • 若启用容器复用(withReuse(true)),需确保Testcontainers复用功能已开启:在~/.testcontainers.properties中添加testcontainers.reuse.enable=true

内容的提问来源于stack exchange,提问作者Siddhartha Manka

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.25 13:19:53