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

使用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
解决方案

错误根源

  1. 未等待ZooKeeper完全就绪就启动Kafka,导致Kafka无法连接ZK
  2. Kafka外部监听绑定localhost,限制了非本地访问
  3. 固定端口9093可能与生产环境冲突
  4. 主机无法解析容器内的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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.11 14:34:51