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

在.NET TestContainers中Kafka与Schema Registry容器协作异常排查

问题:Kafka与Schema Registry容器无法协同工作(.NET TestContainers)

以下是我的容器配置:

private readonly KafkaContainer _kafkaContainer = new KafkaBuilder()
    .WithImage("confluentinc/cp-kafka:7.4.1")
    .WithHostname("broker")
    .WithName("broker")
    .WithNetworkAliases("broker")
    .WithNetwork(NetworkName)
    .WithPortBinding(9092, 9092)
    .WithWaitStrategy(Wait.ForUnixContainer().UntilPortIsAvailable(9092))
    .WithEnvironment(
        new ReadOnlyDictionary<string, string>(new Dictionary<string, string>
        {
            { "KAFKA_ADVERTISED_LISTENERS", "PLAINTEXT://broker:29092,PLAINTEXT_HOST://localhost:9092" }
        }))
    .WithCleanUp(true)
    .Build();

private readonly IContainer _schemaRegistryContainer = new ContainerBuilder()
    .WithImage("confluentinc/cp-schema-registry:7.4.1")
    .WithHostname("schema-registry")
    .WithName("schema-registry")
    .WithCleanUp(true)
    .WithEnvironment(
        new ReadOnlyDictionary<string, string>(new Dictionary<string, string>
        {
            { "SCHEMA_REGISTRY_HOST_NAME", "schema-registry" },
            { "SCHEMA_REGISTRY_KAFKASTORE_BOOTSTRAP_SERVERS", "broker:29092" },
            { "SCHEMA_REGISTRY_LISTENERS", "http://0.0.0.0:8081" }
        }))
    .WithPortBinding(8081, 8081)
    .WithWaitStrategy(Wait.ForUnixContainer().UntilPortIsAvailable(8081))
    .WithNetwork(NetworkName)
    .Build();

Schema Registry容器日志报错:

2023-12-21 13:27:21 [2023-12-21 12:27:21,315] WARN [AdminClient clientId=adminclient-1] Connection to node -1 (broker/172.27.0.2:29092) could not be established. Broker may not be available. (org.apache.kafka.clients.NetworkClient)
2023-12-21 13:27:22 [2023-12-21 12:27:22,364] INFO [AdminClient clientId=adminclient-1] Node -1 disconnected. (org.apache.kafka.clients.NetworkClient)
2023-12-21 13:27:22 [2023-12-21 12:27:22,364] WARN [AdminClient clientId=adminclient-1] Connection to node -1 (broker/172.27.0.2:29092) could not be established. Broker may not be available. (org.apache.kafka.clients.NetworkClient)
2023-12-21 13:27:23 [2023-12-21 12:27:23,604] INFO [AdminClient clientId=adminclient-1] Node -1 disconnected. (org.apache.kafka.clients.NetworkClient)
2023-12-21 13:27:23 [2023-12-21 12:27:23,604] WARN [AdminClient clientId=adminclient-1] Connection to node -1 (broker/172.27.0.2:29092) could not be established. Broker may not be available. (org.apache.kafka.clients.NetworkClient)
2023-12-21 13:27:24 [2023-12-21 12:27:24,741] INFO [AdminClient clientId=adminclient-1] Node -1 disconnected. (org.apache.kafka.clients.NetworkClient)
2023-12-21 13:27:24 [2023-12-21 12:27:24,742] WARN [AdminClient clientId=adminclient-1] Connection to node -1 (broker/172.27.0.2:29092) could not be established. Broker may not be available. (org.apache.kafka.clients.NetworkClient)
2023-12-21 13:27:25 [2023-12-21 12:27:25,772] INFO [AdminClient clientId=adminclient-1] Node -1 disconnected. (org.apache.kafka.clients.NetworkClient)

请问需要修改哪些配置才能让它们成功协作?


解决方案
  • 补充Kafka容器的核心环境变量
    当前Kafka容器缺少关键配置,导致内部未正确监听容器间通信的broker:29092端口:

    1. KAFKA_LISTENERS:指定Kafka监听的地址集合,必须包含容器内部通信的PLAINTEXT://broker:29092和本地访问的PLAINTEXT_HOST://0.0.0.0:9092
    2. KAFKA_LISTENER_SECURITY_PROTOCOL_MAP:映射监听器到安全协议,设置为PLAINTEXT:PLAINTEXT,PLAINTEXT_HOST:PLAINTEXT确保协议匹配
    3. KAFKA_INTER_BROKER_LISTENER_NAME:指定内部broker通信使用的监听器,设置为PLAINTEXT
    4. KAFKA_OFFSETS_TOPIC_REPLICATION_FACTOR:单节点Kafka必须设为1,否则内部主题无法创建,导致服务异常

    修改后的Kafka容器环境变量:

    .WithEnvironment(
        new ReadOnlyDictionary<string, string>(new Dictionary<string, string>
        {
            { "KAFKA_LISTENERS", "PLAINTEXT://broker:29092,PLAINTEXT_HOST://0.0.0.0:9092" },
            { "KAFKA_ADVERTISED_LISTENERS", "PLAINTEXT://broker:29092,PLAINTEXT_HOST://localhost:9092" },
            { "KAFKA_LISTENER_SECURITY_PROTOCOL_MAP", "PLAINTEXT:PLAINTEXT,PLAINTEXT_HOST:PLAINTEXT" },
            { "KAFKA_INTER_BROKER_LISTENER_NAME", "PLAINTEXT" },
            { "KAFKA_OFFSETS_TOPIC_REPLICATION_FACTOR", "1" }
        }))
    
  • 调整Kafka容器的等待策略
    原有的UntilPortIsAvailable(9092)仅检查端口开放,不代表Kafka服务已就绪。改用日志匹配等待,确保Kafka完全启动:

    .WithWaitStrategy(Wait.ForUnixContainer().UntilMessageIsLogged("started (kafka.server.KafkaServer)"))
    
  • 按顺序启动容器
    必须先启动Kafka容器并等待其就绪,再启动Schema Registry容器,避免Registry连接时Kafka还未完全初始化:

    await _kafkaContainer.StartAsync();
    await _schemaRegistryContainer.StartAsync();
    

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.03 22:15:55