求带SASL用户名密码的Confluent Cloud Kafka Go消费者客户端示例
Confluent Cloud Kafka Go消费者示例(含SASL认证)
以下是完整的Go语言消费者示例,包含Confluent Cloud所需的SASL认证配置:
package main import ( "fmt" "log" "github.com/confluentinc/confluent-kafka-go/v2/kafka" ) func main() { // 核心配置参数 config := &kafka.ConfigMap{ "bootstrap.servers": "你的Confluent Cloud集群Bootstrap地址", // 示例:pkc-abc123.us-east-1.aws.confluent.cloud:9092 "sasl.username": "你的API Key", "sasl.password": "你的API Secret", "security.protocol": "SASL_SSL", "sasl.mechanism": "PLAIN", "group.id": "test-consumer-group", // 自定义消费者组ID "auto.offset.reset": "earliest", // 可选:从最早未消费的消息开始 } // 初始化消费者 consumer, err := kafka.NewConsumer(config) if err != nil { log.Fatalf("创建消费者失败: %v", err) } defer consumer.Close() // 订阅目标主题 err = consumer.SubscribeTopics([]string{"你的目标主题名称"}, nil) if err != nil { log.Fatalf("订阅主题失败: %v", err) } // 持续消费消息 for { msg, err := consumer.ReadMessage(-1) if err != nil { log.Printf("消费出错: %v", err) continue } fmt.Printf("收到消息 | 主题: %s | 分区: %d | 偏移量: %d | 内容: %s\n", *msg.TopicPartition.Topic, msg.TopicPartition.Partition, msg.TopicPartition.Offset, string(msg.Value)) } }
关键配置说明
- bootstrap.servers:必须从Confluent Cloud控制台的「集群设置」中复制完整地址,不要手动修改或省略端口
- sasl.username/sasl.password:对应Confluent Cloud中创建的API Key和Secret,而非你的账户登录密码
- security.protocol:固定设置为
SASL_SSL,Confluent Cloud强制要求SSL加密+SASL认证 - sasl.mechanism:固定设置为
PLAIN,这是Confluent Cloud支持的标准SASL机制
针对「EOF」错误的排查点
你遇到的kafka: client has run out of available brokers to talk to: EOF错误,通常由以下原因导致:
- Bootstrap服务器地址填写错误,检查是否包含完整的域名和端口
- API Key/Secret拼写错误,注意大小写、特殊字符的正确性
- 缺少
security.protocol或sasl.mechanism配置,导致连接方式不兼容 - 网络环境限制:防火墙、代理阻止了9092端口的对外连接
- 消费者组ID包含非法字符,导致集群无法识别
内容的提问来源于stack exchange,提问作者xsqian
相关产品推荐
相关产品推荐

