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

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

Kafka Schema Registry容器无法连接Broker(Testcontainers环境)

我需要运行集成测试验证Kafka监听器和Avro序列化功能,依赖Kafka、Schema Registry(间接依赖Zookeeper)。目前用docker-compose能正常搭建环境,但想改用Testcontainers编程式构建,减少人为错误。

Kafka和Zookeeper启动正常,应用能创建主题、监听器订阅,甚至用Kafka控制台生产者发消息都没问题。但Schema Registry容器启动后,能连接Zookeeper,却无法和Broker建立连接,多次重试后超时退出,导致Schema注册/读取失败,测试因序列化报错。

我需要Kafka和Schema Registry都正常工作,不能省略任何一个。也可以继续用docker-compose,但更倾向全编程式环境。

Schema Registry容器日志

2023-02-08 16:56:09 [2023-02-08 15:56:09,556] INFO Session establishment complete on server zookeeper/192.168.144.2:2181, session id = 0x1000085b81e0003, negotiated timeout = 40000 (org.apache.zookeeper.ClientCnxn)
2023-02-08 16:56:09 [2023-02-08 15:56:09,696] INFO Session: 0x1000085b81e0003 closed (org.apache.zookeeper.ZooKeeper)
2023-02-08 16:56:09 [2023-02-08 15:56:09,696] INFO EventThread shut down for session: 0x1000085b81e0003 (org.apache.zookeeper.ClientCnxn)
2023-02-08 16:56:09 [2023-02-08 15:56:09,787] INFO AdminClientConfig values:
/*  Omitted for brevity  */
(org.apache.kafka.clients.admin.AdminClientConfig)
2023-02-08 16:56:10 [2023-02-08 15:56:10,284] INFO Kafka version: 7.3.1-ccs (org.apache.kafka.common.utils.AppInfoParser)
2023-02-08 16:56:10 [2023-02-08 15:56:10,284] INFO Kafka commitId: 8628b0341c3c4676 (org.apache.kafka.common.utils.AppInfoParser)
2023-02-08 16:56:10 [2023-02-08 15:56:10,284] INFO Kafka startTimeMs: 1675871770281 (org.apache.kafka.common.utils.AppInfoParser)
2023-02-08 16:56:10 [2023-02-08 15:56:10,308] INFO [AdminClient clientId=adminclient-1] Node -1 disconnected. (org.apache.kafka.clients.NetworkClient)
2023-02-08 16:56:10 [2023-02-08 15:56:10,313] WARN [AdminClient clientId=adminclient-1] Connection to node -1 (localhost/127.0.0.1:54776) could not be established. Broker may not be available. (org.apache.kafka.clients.NetworkClient)
/*  These lines repeat a few times until the container times out and exits.  */
2023-02-08 16:56:50 [2023-02-08 15:56:50,144] INFO [AdminClient clientId=adminclient-1] Node -1 disconnected. (org.apache.kafka.clients.NetworkClient)
2023-02-08 16:56:50 [2023-02-08 15:56:50,144] WARN [AdminClient clientId=adminclient-1] Connection to node -1 (localhost/127.0.0.1:54776) could not be established. Broker may not be available. (org.apache.kafka.clients.NetworkClient)
2023-02-08 16:56:50 [2023-02-08 15:56:50,298] ERROR Error while getting broker list. (io.confluent.admin.utils.ClusterStatus)
2023-02-08 16:56:50 java.util.concurrent.ExecutionException: org.apache.kafka.common.errors.TimeoutException: Timed out waiting for a node assignment. Call: listNodes
2023-02-08 16:56:50     at java.base/java.util.concurrent.CompletableFuture.reportGet(CompletableFuture.java:395)
2023-02-08 16:56:50     at java.base/java.util.concurrent.CompletableFuture.get(CompletableFuture.java:1999)
2023-02-08 16:56:50     at org.apache.kafka.common.internals.KafkaFutureImpl.get(KafkaFutureImpl.java:165)
2023-02-08 16:56:50     at io.confluent.admin.utils.ClusterStatus.isKafkaReady(ClusterStatus.java:147)
2023-02-08 16:56:50     at io.confluent.admin.utils.cli.KafkaReadyCommand.main(KafkaReadyCommand.java:149)
2023-02-08 16:56:50 Caused by: org.apache.kafka.common.errors.TimeoutException: Timed out waiting for a node assignment. Call: listNodes
2023-02-08 16:56:51 [2023-02-08 15:56:51,103] INFO [AdminClient clientId=adminclient-1] Node -1 disconnected. (org.apache.kafka.clients.NetworkClient)
2023-02-08 16:56:51 [2023-02-08 15:56:51,103] WARN [AdminClient clientId=adminclient-1] Connection to node -1 (localhost/127.0.0.1:54776) could not be established. Broker may not be available. (org.apache.kafka.clients.NetworkClient)
2023-02-08 16:56:51 [2023-02-08 15:56:51,300] INFO Expected 1 brokers but found only 0. Trying to query Kafka for metadata again ... (io.confluent.admin.utils.ClusterStatus)
2023-02-08 16:56:51 [2023-02-08 15:56:51,300] ERROR Expected 1 brokers but found only 0. Brokers found []. (io.confluent.admin.utils.ClusterStatus)
2023-02-08 16:56:51 Using log4j config /etc/schema-registry/log4j.properties

Testcontainers基础测试类代码

@Testcontainers
@SpringBootTest
@Slf4j
public class AbstractIT {

  private static final Network network = Network.newNetwork();

  protected static GenericContainer ZOOKEEPER = new GenericContainer<>(
      DockerImageName.parse("confluentinc/cp-zookeeper:7.2.0"))
      .withNetwork(network)
      .withNetworkAliases("zookeeper")
      .withEnv(Map.of(
          "ZOOKEEPER_CLIENT_PORT", "2181",
          "ZOOKEEPER_TICK_TIME", "2000"));
  
  protected static final KafkaContainer KAFKA = new KafkaContainer(
      DockerImageName.parse("confluentinc/cp-kafka"))
      .withExternalZookeeper("zookeeper:2181")
      .dependsOn(ZOOKEEPER)
      .withNetwork(network)
      .withNetworkAliases("broker");

  protected static final GenericContainer SCHEMAREGSISTRY = new GenericContainer<>(
      DockerImageName.parse("confluentinc/cp-schema-registry"))
      .dependsOn(ZOOKEEPER, KAFKA)
      .withEnv(Map.of(
          "SCHEMA_REGISTRY_KAFKASTORE_CONNECTION_URL", "zookeeper:2181",
          "SCHEMA_REGISTRY_HOST_NAME", "schemaregistry",
          "SCHEMA_REGISTRY_LISTENERS", "http://0.0.0.0:8085",
          "SCHEMA_REGISTRY_KAFKASTORE_BOOTSTRAP_SERVERS", "broker:9092"))
      .withNetwork(network)
      .withNetworkAliases("schemaregistry");

  @DynamicPropertySource
  static void registerPgProperties(DynamicPropertyRegistry registry) {
    registry.add("bootstrap.servers", KAFKA::getBootstrapServers);
    registry.add("spring.kafka.bootstrap-servers", KAFKA::getBootstrapServers);
    registry.add("spring.kafka.consumer.auto-offset-reset", () -> "earliest");
    registry.add("spring.data.mongodb.uri", MONGODB::getConnectionString);
    registry.add("spring.data.mongodb.database", () ->"test");
  }

//container startup, shutdown as well as topic creation omitted for brevity

}

目标复刻的docker-compose.yml

version: "3.5"
services:
  zookeeper:
    image: confluentinc/cp-zookeeper:7.2.0
    hostname: zookeeper
    container_name: zookeeper
    ports:
      - "2181:2181"
    environment:
      ZOOKEEPER_CLIENT_PORT: 2181
      ZOOKEEPER_TICK_TIME: 2000

  broker:
    image: confluentinc/cp-kafka:7.2.0
    hostname: broker
    container_name: broker
    restart: always
    depends_on:
      - zookeeper
    ports:
      - "29092:29092"
      - "9092:9092"
    environment:
      KAFKA_BROKER_ID: 1
      KAFKA_ZOOKEEPER_CONNECT: 'zookeeper:2181'
      KAFKA_LISTENER_SECURITY_PROTOCOL_MAP: PLAINTEXT:PLAINTEXT,PLAINTEXT_HOST:PLAINTEXT
      KAFKA_ADVERTISED_LISTENERS: PLAINTEXT://broker:9092,PLAINTEXT_HOST://localhost:29092
      KAFKA_OFFSETS_TOPIC_REPLICATION_FACTOR: 1
      KAFKA_GROUP_INITIAL_REBALANCE_DELAY_MS: 0
      KAFKA_TRANSACTION_STATE_LOG_MIN_ISR: 1
      KAFKA_TRANSACTION_STATE_LOG_REPLICATION_FACTOR: 1
      KAFKA_SCHEMA_REGISTRY_URL: "schemaregistry:8085"


  schemaregistry:
    container_name: schemaregistry
    hostname: schemaregistry
    image: confluentinc/cp-schema-registry:5.1.2
    restart: always
    depends_on:
      - zookeeper
    environment:
      SCHEMA_REGISTRY_KAFKASTORE_CONNECTION_URL: "zookeeper:2181"
      SCHEMA_REGISTRY_HOST_NAME: schemaregistry
      SCHEMA_REGISTRY_LISTENERS: "http://0.0.0.0:8085"
    ports:
      - "8085:8085"
    volumes:
      - "./src/main/avro/:/etc/schema"

问题原因与修复方案

核心问题分析

从日志可见,Schema Registry尝试连接localhost:54776而非broker:9092,说明Kafka Broker告知Schema Registry的地址是容器内的localhost随机端口,而非容器网络内的broker:9092。这是因为Testcontainers默认的KafkaContainer未配置符合容器网络的advertised.listeners,导致Broker在容器内注册的地址无法被其他容器访问。

另外,docker-compose中Kafka明确配置了KAFKA_ADVERTISED_LISTENERS,但Testcontainers代码里的KafkaContainer缺少该配置,这是关键差异;同时Schema Registry与Kafka版本差距过大(5.1.2 vs 7.2.0),存在兼容性问题。

具体修复步骤

1. 配置KafkaContainer的Advertised Listeners

修改KafkaContainer定义,添加容器网络内的监听配置,确保Broker在容器集群内暴露正确地址:

protected static final KafkaContainer KAFKA = new KafkaContainer(
    DockerImageName.parse("confluentinc/cp-kafka:7.2.0")) // 和docker-compose保持版本一致
    .withExternalZookeeper("zookeeper:2181")
    .dependsOn(ZOOKEEPER)
    .withNetwork(network)
    .withNetworkAliases("broker")
    .withEnv("KAFKA_LISTENER_SECURITY_PROTOCOL_MAP", "PLAINTEXT:PLAINTEXT,PLAINTEXT_HOST:PLAINTEXT")
    .withEnv("KAFKA_ADVERTISED_LISTENERS", "PLAINTEXT://broker:9092,PLAINTEXT_HOST://localhost:${random.port}")
    .withEnv("KAFKA_OFFSETS_TOPIC_REPLICATION_FACTOR", "1");

2. 统一Schema Registry版本并修正配置

Schema Registry版本必须与Kafka一致,避免兼容性问题,同时明确指定Broker连接地址:

protected static final GenericContainer SCHEMAREGSISTRY = new GenericContainer<>(
    DockerImageName.parse("confluentinc/cp-schema-registry:7.2.0")) // 和Kafka版本对齐
    .dependsOn(ZOOKEEPER, KAFKA)
    .withEnv(Map.of(
        "SCHEMA_REGISTRY_KAFKASTORE_CONNECTION_URL", "zookeeper:2181",
        "SCHEMA_REGISTRY_HOST_NAME", "schemaregistry",
        "SCHEMA_REGISTRY_LISTENERS", "http://0.0.0.0:8085",
        "SCHEMA_REGISTRY_KAFKASTORE_BOOTSTRAP_SERVERS", "PLAINTEXT://broker:9092"))
    .withNetwork(network)
    .withNetworkAliases("schemaregistry")
    .withHealthCheck(HealthChecks.httpGet("/subjects").forPort(8085).withStartupTimeout(Duration.ofMinutes(2))); // 添加健康检查,等待Kafka就绪

3. 添加Schema Registry地址到Spring配置

在@DynamicPropertySource方法中添加Schema Registry的访问地址,让应用能正常调用:

registry.add("spring.kafka.properties.schema.registry.url", () -> 
    "http://" + SCHEMAREGSISTRY.getHost() + ":" + SCHEMAREGSISTRY.getMappedPort(8085));

验证修复

启动测试后,查看Schema Registry日志,若不再出现Connection to node -1 (localhost/127.0.0.1:xxxx)的错误,且能看到成功连接Broker的日志,说明修复生效,Avro序列化/反序列化测试可正常执行。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.01 02:35:26