使用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服务状态:
- 查看Kafka容器日志确认启动正常:
docker-compose logs kafka - 进入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
- 查看Kafka容器日志确认启动正常:
检查端口占用情况:确认宿主机9092端口未被其他程序占用,可通过以下命令排查:
- Linux/macOS:
lsof -i :9092 - Windows:
netstat -ano | findstr :9092
- Linux/macOS:
完善Sarama基础配置:确保Sarama配置包含必要的基础参数,比如版本兼容:
kc.sConfig = sarama.NewConfig() kc.sConfig.Version = sarama.V2_8_0_0 // 匹配你的Kafka版本,可查看wurstmeister/kafka镜像文档确认默认版本
内容的提问来源于stack exchange,提问作者nespondev
相关产品推荐
相关产品推荐

