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

使用Go创建本地Kafka生产者失败,报错无可用Broker

问题:Kafka生产者连接Docker部署的Kafka失败,提示无可用Broker

我在本地用Docker Compose部署了Zookeeper和Kafka服务,配置如下:

services:
  zookeeper:
    image: wurstmeister/zookeeper
    ports:
      - "2181:2181"
    hostname: zookeeper
    tmpfs: "/datalog"
  kafka:
    image: wurstmeister/kafka
    command: [ start-kafka.sh ]
    ports:
      - 9092:9092
    environment:
      KAFKA_CREATE_TOPICS: "Kafka-audit"
      KAFKA_ZOOKEEPER_CONNECT: zookeeper:2181
      KAFKA_ADVERTISED_LISTENERS: PLAINTEXT://kafka:9093,OUTSIDE://localhost:9092
      KAFKA_LISTENER_SECURITY_PROTOCOL_MAP: PLAINTEXT:PLAINTEXT,OUTSIDE:PLAINTEXT
      KAFKA_LISTENERS: PLAINTEXT://0.0.0.0:9093,OUTSIDE://0.0.0.0:9092
      KAFKA_INTER_BROKER_LISTENER_NAME: PLAINTEXT
    volumes:
      - /var/run/docker.sock:/var/run/docker.sock
    depends_on:
      - zookeeper

随后用Go代码(基于Sarama库)尝试创建生产者,代码如下:

func NewAuditKafkaConfig(kafkaConf *config.KafkaConfig) (*AuditKafkaConfig, error) {
    kConfig := &Config{}
    kConfig.SetBrokerList(kafkaConf.BootstrapServers)
    kConfig.SetConnection(kafkaConf.BootstrapServers[0])
    kConfig.SetGroup(kafkaConf.ConsumerGroupID)
    kConfig.SetTopic(kafkaConf.Topic)
    err := setKafkaTLS(kConfig)
    if err != nil {
        log.Error().Msgf("Failed to set kafka tls %s", err.Error())
    }
    producer, err := newProducer(kConfig)
    if err != nil {
        log.Error().Msgf("Failed to create producer instance: %s", err.Error())
    }
    config := &AuditKafkaConfig{
        Producer: producer,
        Config:   kConfig,
    }
    return config, err
}

func newProducer(kConfig *Config) (*Producer, error) {
    producer, err := NewProducer(kConfig)
    return producer, err
}

func setKafkaTLS(kc *Config) error {
    kc.TLS = true
    if config.AppConfig.KafkaConfig.DisableTLS {
        kc.TLS = false
    }
    c := &tls.Config{}
    pemCertBlock, _ := base64.StdEncoding.DecodeString(os.Getenv("KAFKA_CLIENT_CERT"))
    pemKeyBlock, _ := base64.StdEncoding.DecodeString(os.Getenv("KAFKA_CLIENT_KEY"))
    rootCABlock, _ := base64.StdEncoding.DecodeString(os.Getenv("KAFKA_SERVER_CERT"))
    if len(pemCertBlock) > 0 && len(pemKeyBlock) > 0 {
        cert, err := tls.X509KeyPair(pemCertBlock, pemKeyBlock)
        if err != nil {
            return err
        }
        c.Certificates = []tls.Certificate{cert}
    }

    if len(rootCABlock) > 0 {
        certPool := x509.NewCertPool()
        if ok := certPool.AppendCertsFromPEM(rootCABlock); !ok {
            return fmt.Errorf("cert error")
        }
        c.RootCAs = certPool
    }
    kc.SetDefaults(c)
    return nil
}

func NewProducer(c *Config) (*Producer, error) {
    pro := new(Producer)
    var err error
    pro.ErrChannel = make(chan error, 5)
    pro.async, err = sarama.NewAsyncProducer(c.brokerList, c.sConfig)
    if err != nil {
        err = fmt.Errorf("failed to create producer: %s", err)
    }
    return pro, err
}

运行后出现错误:

"Failed to create producer instance: failed to create async producer: kafka: client has run out of available brokers to talk to (Is your cluster reachable?)"


解决建议

  • 修正TLS配置冲突:你的Kafka配置中所有监听器都使用PLAINTEXT(无加密),但Go代码默认开启了TLS。如果本地测试未配置证书,需确保config.AppConfig.KafkaConfig.DisableTLS设置为true,同时在setKafkaTLS函数中明确关闭Sarama的TLS开关:

    func setKafkaTLS(kc *Config) error {
        kc.TLS = !config.AppConfig.KafkaConfig.DisableTLS
        kc.sConfig.Net.TLS.Enable = kc.TLS
        if kc.TLS {
            c := &tls.Config{}
            // 原证书加载逻辑...
            kc.sConfig.Net.TLS.Config = c
        }
        kc.SetDefaults(c)
        return nil
    }
    
  • 确认BootstrapServers配置:将kafkaConf.BootstrapServers设置为["localhost:9092"],对应Docker中Kafka对外暴露的OUTSIDE监听器地址。

  • 验证Kafka服务状态:

    1. 查看Kafka容器日志确认启动正常:docker-compose logs kafka
    2. 进入Kafka容器,用内置工具检查主题和Broker状态:
      docker exec -it <kafka-container-id> kafka-topics.sh --list --bootstrap-server localhost:9092
      docker exec -it <kafka-container-id> kafka-broker-api-versions.sh --bootstrap-server localhost:9092
      
  • 检查端口占用情况:确认宿主机9092端口未被其他程序占用,可通过以下命令排查:

    • Linux/macOS:lsof -i :9092
    • Windows:netstat -ano | findstr :9092
  • 完善Sarama基础配置:确保Sarama配置包含必要的基础参数,比如版本兼容:

    kc.sConfig = sarama.NewConfig()
    kc.sConfig.Version = sarama.V2_8_0_0 // 匹配你的Kafka版本,可查看wurstmeister/kafka镜像文档确认默认版本
    

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.27 22:53:20