使用Ory Dockertest部署Confluent Kafka/ZooKeeper遇域名解析问题求助
问题描述
无法通过Ory Dockertest启动Confluentinc的ZooKeeper和Kafka实例,需求如下:
- 测试实例使用独立端口,不干扰生产环境Kafka
- 支持非localhost主机访问
当前代码:
package main import ( "fmt" "github.com/ory/dockertest/v3" "github.com/ory/dockertest/v3/docker" "time" ) const namePrefix = "test_container" var kafkaPort string //used in tests func main() { containers, err := StartDockerContainers() if err != nil { panic(err) } containers.Stop() } type KafkaDockerTest struct { pool *dockertest.Pool network *docker.Network zookeeper *dockertest.Resource kafka *dockertest.Resource } func StartDockerContainers() (*KafkaDockerTest, error) { pool, err := dockertest.NewPool("") if err != nil { return nil, fmt.Errorf("could not connect to docker: %w", err) } err = pool.Client.Ping() if err != nil { return nil, fmt.Errorf("could not connect to docker: %w", err) } timestampAssignation := time.Now().UnixNano() network, err := pool.Client.CreateNetwork(docker.CreateNetworkOptions{Name: fmt.Sprintf("zookeeper_kafka_network_%d", timestampAssignation)}) if err != nil { return nil, fmt.Errorf("could not create a network to zookeeper and kafka: %w", err) } zookeeperOptions := &dockertest.RunOptions{ Name: fmt.Sprintf("%s-zookeeper-%d", namePrefix, timestampAssignation), Repository: "confluentinc/cp-zookeeper", NetworkID: network.ID, Tag: "latest", Hostname: "zookeeper", Env: []string{ "ZOOKEEPER_CLIENT_PORT=2181", }, } zookeeperResource, err := pool.RunWithOptions(zookeeperOptions, func(config *docker.HostConfig) { // set AutoRemove to true so that stopped container goes away by itself config.AutoRemove = false //for debug config.RestartPolicy = docker.RestartPolicy{ Name: "no", } }) if err != nil { return nil, fmt.Errorf("could not start zookeeper: %s", err) } //todo ping - pool.Retry(func() error {} kafkaOptions := &dockertest.RunOptions{ Name: fmt.Sprintf("%s-kafka-%d", namePrefix, timestampAssignation), Repository: "confluentinc/cp-kafka", Tag: "latest", Hostname: "kafka", NetworkID: network.ID, Env: []string{ "KAFKA_ADVERTISED_LISTENERS=INSIDE://kafka:9092,OUTSIDE://localhost:9093", "KAFKA_LISTENERS=INSIDE://0.0.0.0:9092,OUTSIDE://0.0.0.0:9093", "KAFKA_BROKER_ID=1", "KAFKA_OFFSETS_TOPIC_REPLICATION_FACTOR=1", "KAFKA_ZOOKEEPER_CONNECT=zookeeper:2181", "KAFKA_INTER_BROKER_LISTENER_NAME=INSIDE", "KAFKA_LISTENER_SECURITY_PROTOCOL_MAP=INSIDE:PLAINTEXT,OUTSIDE:PLAINTEXT", }, PortBindings: map[docker.Port][]docker.PortBinding{ "9093/tcp": {{HostIP: "localhost", HostPort: "9093/tcp"}}, }, ExposedPorts: []string{"9093/tcp"}, } kafkaResource, err := pool.RunWithOptions(kafkaOptions, func(config *docker.HostConfig) { // set AutoRemove to true so that stopped container goes away by itself config.AutoRemove = false config.RestartPolicy = docker.RestartPolicy{ Name: "no", } }) if err != nil { return nil, fmt.Errorf("could not start kafka: %s", err) } //todo ping - pool.Retry(func() error {} kafkaPort = kafkaResource.GetPort("9093/tcp") return &KafkaDockerTest{ pool: pool, network: network, zookeeper: zookeeperResource, kafka: kafkaResource, }, nil } func (kdt *KafkaDockerTest) Stop() error { if kdt == nil || kdt.zookeeper == nil { return fmt.Errorf("could not stop zookeeper container") } if err := kdt.pool.Purge(kdt.zookeeper); err != nil { return fmt.Errorf("could not purge resource %v :%s", kdt.zookeeper, err) } if kdt == nil || kdt.kafka == nil { return fmt.Errorf("could not stop kafka container") } if err := kdt.pool.Purge(kdt.kafka); err != nil { fmt.Errorf("could not purge resource %v :%s", kdt.kafka, err) } if err := kdt.pool.Client.RemoveNetwork(kdt.network.ID); err != nil { return fmt.Errorf("could not remove network %s :%s", kdt.network.ID, err) } return nil }
运行后报错:
dial tcp: lookup kafka: no such host
解决方案
错误根源
- 未等待ZooKeeper完全就绪就启动Kafka,导致Kafka无法连接ZK
- Kafka外部监听绑定
localhost,限制了非本地访问 - 固定端口9093可能与生产环境冲突
- 主机无法解析容器内的
kafka主机名,缺少必要的配置
具体修改步骤
1. 等待ZooKeeper就绪
启动Kafka前,添加重试逻辑确保ZK服务可用(需导入net包):
// 等待ZooKeeper就绪 if err := pool.Retry(func() error { conn, err := net.Dial("tcp", fmt.Sprintf("localhost:%s", zookeeperResource.GetPort("2181/tcp"))) if err != nil { return err } conn.Close() return nil }); err != nil { return nil, fmt.Errorf("zookeeper failed to start: %w", err) }
2. 修正Kafka配置
- 动态分配主机端口避免冲突
- 允许非localhost访问
- 添加主机名解析配置:
kafkaOptions := &dockertest.RunOptions{ Name: fmt.Sprintf("%s-kafka-%d", namePrefix, timestampAssignation), Repository: "confluentinc/cp-kafka", Tag: "latest", Hostname: "kafka", NetworkID: network.ID, Env: []string{ "KAFKA_ADVERTISED_LISTENERS=INSIDE://kafka:9092,OUTSIDE://0.0.0.0:9093", "KAFKA_LISTENERS=INSIDE://0.0.0.0:9092,OUTSIDE://0.0.0.0:9093", "KAFKA_BROKER_ID=1", "KAFKA_OFFSETS_TOPIC_REPLICATION_FACTOR=1", "KAFKA_ZOOKEEPER_CONNECT=zookeeper:2181", "KAFKA_INTER_BROKER_LISTENER_NAME=INSIDE", "KAFKA_LISTENER_SECURITY_PROTOCOL_MAP=INSIDE:PLAINTEXT,OUTSIDE:PLAINTEXT", // 解决主机名解析问题 "KAFKA_ADVERTISED_HOST_NAME=kafka", }, // 动态分配端口,避免与生产环境冲突 PortBindings: map[docker.Port][]docker.PortBinding{ "9093/tcp": {{HostIP: "0.0.0.0", HostPort: "0"}}, }, ExposedPorts: []string{"9093/tcp"}, }
3. 等待Kafka就绪
启动Kafka后添加重试逻辑:
// 等待Kafka就绪 if err := pool.Retry(func() error { conn, err := net.Dial("tcp", fmt.Sprintf("localhost:%s", kafkaResource.GetPort("9093/tcp"))) if err != nil { return err } conn.Close() return nil }); err != nil { return nil, fmt.Errorf("kafka failed to start: %w", err) }
4. 修复Stop方法的错误
原方法中清理Kafka时未返回错误,修正如下:
if err := kdt.pool.Purge(kdt.kafka); err != nil { return fmt.Errorf("could not purge resource %v :%s", kdt.kafka, err) }
完整修改后代码
package main import ( "fmt" "net" "github.com/ory/dockertest/v3" "github.com/ory/dockertest/v3/docker" "time" ) const namePrefix = "test_container" var kafkaPort string //used in tests func main() { containers, err := StartDockerContainers() if err != nil { panic(err) } defer containers.Stop() fmt.Printf("Kafka is running on port: %s\n", kafkaPort) // 此处可添加测试逻辑 } type KafkaDockerTest struct { pool *dockertest.Pool network *docker.Network zookeeper *dockertest.Resource kafka *dockertest.Resource } func StartDockerContainers() (*KafkaDockerTest, error) { pool, err := dockertest.NewPool("") if err != nil { return nil, fmt.Errorf("could not connect to docker: %w", err) } err = pool.Client.Ping() if err != nil { return nil, fmt.Errorf("could not connect to docker: %w", err) } timestampAssignation := time.Now().UnixNano() network, err := pool.Client.CreateNetwork(docker.CreateNetworkOptions{Name: fmt.Sprintf("zookeeper_kafka_network_%d", timestampAssignation)}) if err != nil { return nil, fmt.Errorf("could not create a network to zookeeper and kafka: %w", err) } zookeeperOptions := &dockertest.RunOptions{ Name: fmt.Sprintf("%s-zookeeper-%d", namePrefix, timestampAssignation), Repository: "confluentinc/cp-zookeeper", NetworkID: network.ID, Tag: "latest", Hostname: "zookeeper", Env: []string{ "ZOOKEEPER_CLIENT_PORT=2181", }, } zookeeperResource, err := pool.RunWithOptions(zookeeperOptions, func(config *docker.HostConfig) { config.AutoRemove = false //for debug config.RestartPolicy = docker.RestartPolicy{ Name: "no", } }) if err != nil { return nil, fmt.Errorf("could not start zookeeper: %s", err) } // 等待ZooKeeper就绪 if err := pool.Retry(func() error { conn, err := net.Dial("tcp", fmt.Sprintf("localhost:%s", zookeeperResource.GetPort("2181/tcp"))) if err != nil { return err } conn.Close() return nil }); err != nil { return nil, fmt.Errorf("zookeeper failed to start: %w", err) } kafkaOptions := &dockertest.RunOptions{ Name: fmt.Sprintf("%s-kafka-%d", namePrefix, timestampAssignation), Repository: "confluentinc/cp-kafka", Tag: "latest", Hostname: "kafka", NetworkID: network.ID, Env: []string{ "KAFKA_ADVERTISED_LISTENERS=INSIDE://kafka:9092,OUTSIDE://0.0.0.0:9093", "KAFKA_LISTENERS=INSIDE://0.0.0.0:9092,OUTSIDE://0.0.0.0:9093", "KAFKA_BROKER_ID=1", "KAFKA_OFFSETS_TOPIC_REPLICATION_FACTOR=1", "KAFKA_ZOOKEEPER_CONNECT=zookeeper:2181", "KAFKA_INTER_BROKER_LISTENER_NAME=INSIDE", "KAFKA_LISTENER_SECURITY_PROTOCOL_MAP=INSIDE:PLAINTEXT,OUTSIDE:PLAINTEXT", "KAFKA_ADVERTISED_HOST_NAME=kafka", }, // 动态分配端口,避免冲突 PortBindings: map[docker.Port][]docker.PortBinding{ "9093/tcp": {{HostIP: "0.0.0.0", HostPort: "0"}}, }, ExposedPorts: []string{"9093/tcp"}, } kafkaResource, err := pool.RunWithOptions(kafkaOptions, func(config *docker.HostConfig) { config.AutoRemove = false config.RestartPolicy = docker.RestartPolicy{ Name: "no", } }) if err != nil { return nil, fmt.Errorf("could not start kafka: %s", err) } // 等待Kafka就绪 if err := pool.Retry(func() error { conn, err := net.Dial("tcp", fmt.Sprintf("localhost:%s", kafkaResource.GetPort("9093/tcp"))) if err != nil { return err } conn.Close() return nil }); err != nil { return nil, fmt.Errorf("kafka failed to start: %w", err) } kafkaPort = kafkaResource.GetPort("9093/tcp") return &KafkaDockerTest{ pool: pool, network: network, zookeeper: zookeeperResource, kafka: kafkaResource, }, nil } func (kdt *KafkaDockerTest) Stop() error { if kdt == nil || kdt.zookeeper == nil { return fmt.Errorf("could not stop zookeeper container") } if err := kdt.pool.Purge(kdt.zookeeper); err != nil { return fmt.Errorf("could not purge resource %v :%s", kdt.zookeeper, err) } if kdt == nil || kdt.kafka == nil { return fmt.Errorf("could not stop kafka container") } if err := kdt.pool.Purge(kdt.kafka); err != nil { return fmt.Errorf("could not purge resource %v :%s", kdt.kafka, err) } if err := kdt.pool.Client.RemoveNetwork(kdt.network.ID); err != nil { return fmt.Errorf("could not remove network %s :%s", kdt.network.ID, err) } return nil }
关键说明
- 动态端口:通过
HostPort: "0"让Docker自动分配空闲端口,避免与生产环境冲突 - 服务就绪检测:使用
pool.Retry确保ZK和Kafka完全启动后再执行后续逻辑 - 非本地访问支持:将
HostIP设为0.0.0.0,让Kafka监听所有网卡,外部主机可通过主机IP+分配的端口访问
内容的提问来源于stack exchange,提问作者mikr
相关产品推荐
相关产品推荐

