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

求带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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.30 14:46:01