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

Testcontainers中SchemaRegistry容器无法连接Kafka容器求助

Testcontainers Schema Registry连接Kafka容器失败的可行解决方案

我尝试启动Schema Registry测试容器并连接Kafka,参考相关方法后仍未成功,以下是两次尝试的细节:

首次尝试代码

public SchemaRegistryContainer withKafka( final KafkaContainer kafkaContainer ) {
    withNetwork( kafkaContainer.getNetwork() );
    withEnv( "SCHEMA_REGISTRY_HOST_NAME", "kafka-confluent-cp-schema-registry" );
    withEnv( "SCHEMA_REGISTRY_LISTENERS", "http://0.0.0.0:" + SCHEMA_REGISTRY_PORT );
    withEnv( "SCHEMA_REGISTRY_KAFKASTORE_BOOTSTRAP_SERVERS", kafkaContainer.getBootstrapServers() );
    return self();
}

报错信息

java.util.concurrent.ExecutionException: org.apache.kafka.common.errors.TimeoutException: Timed out waiting for a node assignment. Call: listNodes
        at java.base/java.util.concurrent.CompletableFuture.reportGet(CompletableFuture.java:395)
        at java.base/java.util.concurrent.CompletableFuture.get(CompletableFuture.java:1999)
        at org.apache.kafka.common.internals.KafkaFutureImpl.get(KafkaFutureImpl.java:165)
        at io.confluent.admin.utils.ClusterStatus.isKafkaReady(ClusterStatus.java:147)
        at io.confluent.admin.utils.cli.KafkaReadyCommand.main(KafkaReadyCommand.java:149)
Caused by: org.apache.kafka.common.errors.TimeoutException: Timed out waiting for a node assignment. Call: listNodes

第二次尝试代码

public SchemaRegistryContainer withKafka( final KafkaContainer kafkaContainer ) {
    withNetwork( kafkaContainer.getNetwork() );
    withEnv( "SCHEMA_REGISTRY_HOST_NAME", "schema-registry" );
    withEnv( "SCHEMA_REGISTRY_LISTENERS", "http://0.0.0.0:" + SCHEMA_REGISTRY_PORT );
    withEnv( "SCHEMA_REGISTRY_KAFKASTORE_BOOTSTRAP_SERVERS", "PLAINTEXT://" + kafkaContainer.getNetworkAliases().get( 0 ) + ":9092" );
    return self();
}

报错信息

2023-04-27 07:46:35,687] WARN Couldn't resolve server PLAINTEXT://tc-HJrDbND5:9092 from bootstrap.servers as DNS resolution failed for tc-HJrDbND5 (org.apache.kafka.clients.ClientUtils)
[2023-04-27 07:46:35,687] ERROR Error while running kafka-ready. (io.confluent.admin.utils.cli.KafkaReadyCommand)
org.apache.kafka.common.KafkaException: Failed to create new KafkaAdminClient
        at org.apache.kafka.clients.admin.KafkaAdminClient.createInternal(KafkaAdminClient.java:553)
        at org.apache.kafka.clients.admin.Admin.create(Admin.java:144)
        at org.apache.kafka.clients.admin.AdminClient.create(AdminClient.java:49)
        at io.confluent.admin.utils.ClusterStatus.isKafkaReady(ClusterStatus.java:136)
        at io.confluent.admin.utils.cli.KafkaReadyCommand.main(KafkaReadyCommand.java:149)
Caused by: org.apache.kafka.common.config.ConfigException: No resolvable bootstrap urls given in bootstrap.servers
        at org.apache.kafka.clients.ClientUtils.parseAndValidateAddresses(ClientUtils.java:89)
        at org.apache.kafka.clients.ClientUtils.parseAndValidateAddresses(ClientUtils.java:48)
        at org.apache.kafka.clients.admin.KafkaAdminClient.createInternal(KafkaAdminClient.java:505)
        ... 4 more

可行的测试方案

1. 为Kafka容器设置固定网络别名

给Kafka容器添加固定的网络别名,确保Schema Registry能通过别名稳定解析到Kafka服务:

// 创建Kafka容器时指定固定网络别名
KafkaContainer kafkaContainer = new KafkaContainer(DockerImageName.parse("confluentinc/cp-kafka:7.3.0"))
        .withNetworkAliases("kafka");

// 配置Schema Registry时使用该固定别名
public SchemaRegistryContainer withKafka(final KafkaContainer kafkaContainer) {
    withNetwork(kafkaContainer.getNetwork());
    withEnv("SCHEMA_REGISTRY_HOST_NAME", "schema-registry");
    withEnv("SCHEMA_REGISTRY_LISTENERS", "http://0.0.0.0:" + SCHEMA_REGISTRY_PORT);
    withEnv("SCHEMA_REGISTRY_KAFKASTORE_BOOTSTRAP_SERVERS", "PLAINTEXT://kafka:9092");
    return self();
}

2. 显式配置Kafka的监听与公告地址

确保Kafka容器同时配置内部和外部的监听地址,避免Schema Registry无法获取正确的Kafka节点地址:

KafkaContainer kafkaContainer = new KafkaContainer(DockerImageName.parse("confluentinc/cp-kafka:7.3.0"))
        .withEnv("KAFKA_ADVERTISED_LISTENERS", "PLAINTEXT://kafka:9092,PLAINTEXT_HOST://localhost:9093")
        .withEnv("KAFKA_LISTENERS", "PLAINTEXT://0.0.0.0:9092,PLAINTEXT_HOST://0.0.0.0:9093")
        .withNetworkAliases("kafka");

3. 使用Docker Compose管理容器依赖

通过Docker Compose定义Kafka、Zookeeper和Schema Registry的完整依赖关系,利用Testcontainers加载Compose文件:
创建docker-compose.yml文件:

version: '3.8'
services:
  zookeeper:
    image: confluentinc/cp-zookeeper:7.3.0
    environment:
      ZOOKEEPER_CLIENT_PORT: 2181

  kafka:
    image: confluentinc/cp-kafka:7.3.0
    depends_on:
      - zookeeper
    environment:
      KAFKA_ZOOKEEPER_CONNECT: zookeeper:2181
      KAFKA_ADVERTISED_LISTENERS: PLAINTEXT://kafka:9092,PLAINTEXT_HOST://localhost:9093
      KAFKA_LISTENERS: PLAINTEXT://0.0.0.0:9092,PLAINTEXT_HOST://0.0.0.0:9093

  schema-registry:
    image: confluentinc/cp-schema-registry:7.3.0
    depends_on:
      - kafka
    environment:
      SCHEMA_REGISTRY_HOST_NAME: schema-registry
      SCHEMA_REGISTRY_LISTENERS: http://0.0.0.0:8081
      SCHEMA_REGISTRY_KAFKASTORE_BOOTSTRAP_SERVERS: PLAINTEXT://kafka:9092

测试中加载并启动:

DockerComposeContainer composeContainer = new DockerComposeContainer(new File("docker-compose.yml"))
        .withExposedService("schema-registry", 8081)
        .withExposedService("kafka", 9093);
composeContainer.start();

// 获取Schema Registry的外部访问地址
String schemaRegistryUrl = String.format("http://%s:%d",
        composeContainer.getServiceHost("schema-registry", 8081),
        composeContainer.getServicePort("schema-registry", 8081));

4. 添加Kafka就绪等待逻辑

在启动Schema Registry前,先确保Kafka容器完全就绪可连接,避免因Kafka未启动完成导致连接超时:

kafkaContainer.start();

// 使用Kafka AdminClient等待集群就绪
try (AdminClient adminClient = AdminClient.create(Map.of(
        AdminClientConfig.BOOTSTRAP_SERVERS_CONFIG, kafkaContainer.getBootstrapServers()))) {
    while (true) {
        try {
            // 尝试获取集群节点信息,超时10秒
            adminClient.describeCluster().nodes().get(10, TimeUnit.SECONDS);
            break;
        } catch (Exception e) {
            // 未就绪则等待1秒后重试
            Thread.sleep(1000);
        }
    }
}

// 启动Schema Registry容器
schemaRegistryContainer.start();

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.23 14:54:56