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

