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
相关产品推荐
相关产品推荐

