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

Testcontainers中Kafka Schema Registry无法连接Kafka Broker问题

Schema Registry无法连接Testcontainers启动的Kafka Broker(Apple M2环境)

在Apple M2 Max环境下,用Testcontainers启动Confluent Kafka和Schema Registry容器后,Schema Registry报连接错误:

Connection to node -1 (localhost/127.0.0.1:54380) could not be established. Broker may not be available.

核心问题

容器之间不能用localhost通信,Schema Registry容器需要访问Kafka容器的内部网络地址,而非主机的localhost。同时需正确配置Kafka的监听地址,兼顾容器内部与外部访问需求。

具体解决步骤

1. 让Kafka和Schema Registry共享同一Testcontainers网络

创建自定义网络,确保两个容器处于同一网络域,可通过容器名互相访问:

Network network = Network.newNetwork();

2. 正确配置Kafka容器环境变量

关键是设置KAFKA_ADVERTISED_LISTENERS,区分容器内部与主机的访问地址:

KafkaContainer kafka = new KafkaContainer(DockerImageName.parse("confluentinc/cp-kafka:7.8.0"))
        .withNetwork(network)
        .withNetworkAliases("kafka") // 给Kafka设置容器别名,方便Schema Registry寻址
        .withEnv("KAFKA_LISTENERS", "PLAINTEXT://0.0.0.0:9092,PLAINTEXT_HOST://0.0.0.0:9093")
        .withEnv("KAFKA_ADVERTISED_LISTENERS", "PLAINTEXT://kafka:9092,PLAINTEXT_HOST://localhost:" + kafka.getMappedPort(9093))
        .withEnv("KAFKA_LISTENER_SECURITY_PROTOCOL_MAP", "PLAINTEXT:PLAINTEXT,PLAINTEXT_HOST:PLAINTEXT")
        .withEnv("KAFKA_OFFSETS_TOPIC_REPLICATION_FACTOR", "1")
        .withEnv("KAFKA_TRANSACTION_STATE_LOG_REPLICATION_FACTOR", "1")
        .withEnv("KAFKA_TRANSACTION_STATE_LOG_MIN_ISR", "1");

3. 配置Schema Registry指向Kafka内部地址

设置SCHEMA_REGISTRY_KAFKASTORE_BOOTSTRAP_SERVERS为Kafka容器别名+内部端口:

SchemaRegistryContainer schemaRegistry = new SchemaRegistryContainer(DockerImageName.parse("confluentinc/cp-schema-registry:7.8.0"))
        .withNetwork(network)
        .withEnv("SCHEMA_REGISTRY_HOST_NAME", "schema-registry")
        .withEnv("SCHEMA_REGISTRY_KAFKASTORE_BOOTSTRAP_SERVERS", "PLAINTEXT://kafka:9092")
        .withEnv("SCHEMA_REGISTRY_LISTENERS", "http://0.0.0.0:8081");

4. 测试代码中访问服务的正确地址

  • Kafka生产者/消费者用主机localhost + Kafka容器映射的PLAINTEXT_HOST端口:直接调用kafka.getBootstrapServers()(Testcontainers会自动返回正确的主机地址)
  • Schema Registry用主机localhost + 映射的8081端口:调用schemaRegistry.getSchemaRegistryUrl()

5. Apple M2环境额外注意

确保拉取的Confluent镜像支持arm64架构,Confluent 7.0+版本已原生支持arm64。若之前拉过amd64镜像,建议清理后重新拉取:

docker rmi confluentinc/cp-kafka:7.8.0 confluentinc/cp-schema-registry:7.8.0
docker pull confluentinc/cp-kafka:7.8.0
docker pull confluentinc/cp-schema-registry:7.8.0

完整测试类示例(JUnit 5)

import org.junit.jupiter.api.Test;
import org.testcontainers.containers.Network;
import org.testcontainers.containers.kafka.KafkaContainer;
import org.testcontainers.containers.schema.SchemaRegistryContainer;
import org.testcontainers.junit.jupiter.Container;
import org.testcontainers.junit.jupiter.Testcontainers;
import org.testcontainers.utility.DockerImageName;

@Testcontainers
public class KafkaAvroIntegrationTest {

    private static final Network NETWORK = Network.newNetwork();
    private static final DockerImageName KAFKA_IMAGE = DockerImageName.parse("confluentinc/cp-kafka:7.8.0");
    private static final DockerImageName SCHEMA_REGISTRY_IMAGE = DockerImageName.parse("confluentinc/cp-schema-registry:7.8.0");

    @Container
    private final KafkaContainer kafka = new KafkaContainer(KAFKA_IMAGE)
            .withNetwork(NETWORK)
            .withNetworkAliases("kafka")
            .withEnv("KAFKA_LISTENERS", "PLAINTEXT://0.0.0.0:9092,PLAINTEXT_HOST://0.0.0.0:9093")
            .withEnv("KAFKA_ADVERTISED_LISTENERS", "PLAINTEXT://kafka:9092,PLAINTEXT_HOST://localhost:" + "%d")
            .withEnv("KAFKA_LISTENER_SECURITY_PROTOCOL_MAP", "PLAINTEXT:PLAINTEXT,PLAINTEXT_HOST:PLAINTEXT")
            .withEnv("KAFKA_OFFSETS_TOPIC_REPLICATION_FACTOR", "1")
            .withEnv("KAFKA_TRANSACTION_STATE_LOG_REPLICATION_FACTOR", "1")
            .withEnv("KAFKA_TRANSACTION_STATE_LOG_MIN_ISR", "1");

    @Container
    private final SchemaRegistryContainer schemaRegistry = new SchemaRegistryContainer(SCHEMA_REGISTRY_IMAGE)
            .withNetwork(NETWORK)
            .withEnv("SCHEMA_REGISTRY_KAFKASTORE_BOOTSTRAP_SERVERS", "PLAINTEXT://kafka:9092")
            .withEnv("SCHEMA_REGISTRY_LISTENERS", "http://0.0.0.0:8081");

    @Test
    void testSchemaRegistryConnection() {
        // 验证Schema Registry可用性
        String schemaRegistryUrl = schemaRegistry.getSchemaRegistryUrl();
        // 可添加Schema Registry REST API测试逻辑,比如获取配置

        // 获取Kafka地址用于AVRO生产者/消费者测试
        String bootstrapServers = kafka.getBootstrapServers();
        // 编写AVRO消息生产、消费逻辑
    }
}

排查技巧

  • 查看Kafka容器日志,确认是否正常启动并监听9092端口:docker logs <kafka-container-id>
  • 查看Schema Registry容器日志,确认bootstrap.servers配置为kafka:9092而非localhost:docker logs <schema-registry-container-id>
  • 进入Schema Registry容器内部,尝试连接Kafka:docker exec -it <schema-registry-container-id> bash,执行kafka-topics.sh --list --bootstrap-server kafka:9092,检查是否能返回topic列表

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.15 05:14:52