使用Testcontainers启动Kafka与Schema Registry失败求助(Kotlin)
问题分析与解决方案
你遇到的Schema Registry连接Kafka超时问题,核心原因是容器间网络通信配置错误,加上缺少必要的启动等待策略。以下是具体修正方案:
核心问题点
- 容器间通信不能使用主机映射的Kafka地址(
kafkaContainer.bootstrapServers),必须用容器网络内的内部地址 - Schema Registry的监听地址配置错误,导致容器外部无法访问
- 未设置Schema Registry的启动等待策略,过早尝试连接未就绪的服务
- 动态属性中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()}" } } }
关键修改说明
- 移除手动
start()调用:Testcontainers会通过@Container注解自动管理容器启动顺序(先启动Kafka,再启动Schema Registry),手动调用会打乱依赖逻辑 - 修正Kafka连接地址:容器在同一网络内时,Kafka容器默认名称为
kafka,使用PLAINTEXT://kafka:9092作为内部通信地址,而非主机映射的端口 - 调整监听地址:将
SCHEMA_REGISTRY_LISTENERS改为http://0.0.0.0:8081,确保容器外部(测试代码)能访问服务 - 添加等待策略:通过
Wait.forHttp("/subjects")等待Schema Registry服务完全就绪,避免启动后立即访问失败 - 修正端口获取方式:用
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
相关产品推荐
相关产品推荐

