活动

天翼云最新优惠活动,涵盖免费试用,产品折扣等,助您降本增效!
热门活动
  • 618智算钜惠季 爆款云主机2核4G限时秒杀,88元/年起!
  • 免费体验DeepSeek,上天翼云息壤 NEW 新老用户均可免费体验2500万Tokens,限时两周
  • 云上钜惠 HOT 爆款云主机全场特惠,更有万元锦鲤券等你来领!
  • 算力套餐 HOT 让算力触手可及
  • 天翼云脑AOne NEW 连接、保护、办公,All-in-One!
  • 中小企业应用上云专场 产品组合下单即享折上9折起,助力企业快速上云
  • 息壤高校钜惠活动 NEW 天翼云息壤杯高校AI大赛,数款产品享受线上订购超值特惠
  • 天翼云电脑专场 HOT 移动办公新选择,爆款4核8G畅享1年3.5折起,快来抢购!
  • 天翼云奖励推广计划 加入成为云推官,推荐新用户注册下单得现金奖励
免费活动
  • 免费试用中心 HOT 多款云产品免费试用,快来开启云上之旅
  • 天翼云用户体验官 NEW 您的洞察,重塑科技边界

智算服务

打造统一的产品能力,实现算网调度、训练推理、技术架构、资源管理一体化智算服务
智算云(DeepSeek专区)
科研助手
  • 算力商城
  • 应用商城
  • 开发机
  • 并行计算
算力互联调度平台
  • 应用市场
  • 算力市场
  • 算力调度推荐
一站式智算服务平台
  • 模型广场
  • 体验中心
  • 服务接入
智算一体机
  • 智算一体机
大模型
  • DeepSeek-R1-昇腾版(671B)
  • DeepSeek-R1-英伟达版(671B)
  • DeepSeek-V3-昇腾版(671B)
  • DeepSeek-R1-Distill-Llama-70B
  • DeepSeek-R1-Distill-Qwen-32B
  • Qwen2-72B-Instruct
  • StableDiffusion-V2.1
  • TeleChat-12B

应用商城

天翼云精选行业优秀合作伙伴及千余款商品,提供一站式云上应用服务
进入甄选商城进入云市场创新解决方案
办公协同
  • WPS云文档
  • 安全邮箱
  • EMM手机管家
  • 智能商业平台
财务管理
  • 工资条
  • 税务风控云
企业应用
  • 翼信息化运维服务
  • 翼视频云归档解决方案
工业能源
  • 智慧工厂_生产流程管理解决方案
  • 智慧工地
建站工具
  • SSL证书
  • 新域名服务
网络工具
  • 翼云加速
灾备迁移
  • 云管家2.0
  • 翼备份
资源管理
  • 全栈混合云敏捷版(软件)
  • 全栈混合云敏捷版(一体机)
行业应用
  • 翼电子教室
  • 翼智慧显示一体化解决方案

合作伙伴

天翼云携手合作伙伴,共创云上生态,合作共赢
天翼云生态合作中心
  • 天翼云生态合作中心
天翼云渠道合作伙伴
  • 天翼云代理渠道合作伙伴
天翼云服务合作伙伴
  • 天翼云集成商交付能力认证
天翼云应用合作伙伴
  • 天翼云云市场合作伙伴
  • 天翼云甄选商城合作伙伴
天翼云技术合作伙伴
  • 天翼云OpenAPI中心
  • 天翼云EasyCoding平台
天翼云培训认证
  • 天翼云学堂
  • 天翼云市场商学院
天翼云合作计划
  • 云汇计划
天翼云东升计划
  • 适配中心
  • 东升计划
  • 适配互认证

开发者

开发者相关功能入口汇聚
技术社区
  • 专栏文章
  • 互动问答
  • 技术视频
资源与工具
  • OpenAPI中心
开放能力
  • EasyCoding敏捷开发平台
培训与认证
  • 天翼云学堂
  • 天翼云认证
魔乐社区
  • 魔乐社区

支持与服务

为您提供全方位支持与服务,全流程技术保障,助您轻松上云,安全无忧
文档与工具
  • 文档中心
  • 新手上云
  • 自助服务
  • OpenAPI中心
定价
  • 价格计算器
  • 定价策略
基础服务
  • 售前咨询
  • 在线支持
  • 在线支持
  • 工单服务
  • 建议与反馈
  • 用户体验官
  • 服务保障
  • 客户公告
  • 会员中心
增值服务
  • 红心服务
  • 首保服务
  • 客户支持计划
  • 专家技术服务
  • 备案管家

了解天翼云

天翼云秉承央企使命,致力于成为数字经济主力军,投身科技强国伟大事业,为用户提供安全、普惠云服务
品牌介绍
  • 关于天翼云
  • 智算云
  • 天翼云4.0
  • 新闻资讯
  • 天翼云APP
基础设施
  • 全球基础设施
  • 信任中心
最佳实践
  • 精选案例
  • 超级探访
  • 云杂志
  • 分析师和白皮书
  • 天翼云·创新直播间
市场活动
  • 2025智能云生态大会
  • 2024智算云生态大会
  • 2023云生态大会
  • 2022云生态大会
  • 天翼云中国行
天翼云
  • 活动
  • 智算服务
  • 产品
  • 解决方案
  • 应用商城
  • 合作伙伴
  • 开发者
  • 支持与服务
  • 了解天翼云
      • 文档
      • 控制中心
      • 备案
      • 管理中心
      文档中心

      分布式消息服务Kafka

      分布式消息服务Kafka

        • 产品动态
        • 服务公告
        • 2024
        • 【优惠】正式开放2年7折,3年5折包年折扣
        • 【优惠】分布式消息服务Kafka增加包年优惠折扣和产品资费进一步下调
        • 【降价】分布式消息服务Kafka产品资费价格下调
        • 【通知】云原生引擎调整为白名单特性
        • 【通知】通用型主机规格调整为白名单特性
        • 产品简介
        • 产品定义
        • 产品优势
        • 功能特性
        • 应用场景
        • 产品规格
        • 开源对比
        • 分布式消息产品选型
        • 安全
        • 认证与访问控制
        • 数据保护技术
        • 审计与日志
        • 服务韧性
        • 监控安全风险
        • 使用限制
        • 名词解释
        • 主子账号和IAM权限管理
        • 与其他服务关系
        • 计费说明
        • 产品资费
        • 新资费
        • 旧资费
        • 计费项
        • 计费模式
        • 续费、到期与欠费
        • 退订
        • 快速入门
        • 入门指引
        • 环境准备
        • 创建实例
        • 创建Topic
        • 编译运行Demo Java工程
        • 配置必须的监控告警
        • 用户指南
        • 权限管理
        • 创建用户并授权使用Kafka
        • 创建Kafka自定义策略
        • 连接Kafka
        • 配置Kafka网络连接
        • 使用VPCEP实现跨VPC访问Kafka
        • 公共接入点接入
        • 安全接入点接入
        • SASL_SSL接入点接入
        • 实例管理
        • 查看实例
        • 设置公网ip
        • 开启IPv6
        • 退订
        • 续订
        • 变更实例规格
        • 计费互转
        • 修改配置参数
        • 重启实例
        • Topic管理
        • 查看Topic
        • 创建Topic
        • 删除Topic
        • 修改Topic
        • 查看分区状态
        • 修改分区平衡
        • 生产消息
        • 删除消息
        • 消费组管理
        • 消费组列表
        • 新建消费组
        • 消息堆积
        • 重置消费位置
        • 删除消费组
        • 用户管理
        • 用户列表
        • 创建应用用户
        • 修改用户信息
        • 删除用户
        • 管理应用用户生产消费权限
        • 消息查询
        • 按位点查询
        • 按时间查询
        • 可观测
        • 监控信息
        • 查看监控数据
        • 支持的监控指标
        • 智能运维
        • 重平衡日志
        • 集群迁移
        • 迁移上云
        • 云实例间迁移
        • 元数据迁移
        • 开发指南
        • 概述
        • 收集连接信息
        • Java
        • Java开发环境搭建
        • Java客户端接入示例
        • Go
        • Python
        • 性能白皮书
        • Kafka性能白皮书
        • 最佳实践
        • 生产者实践
        • 消费者实践
        • 负载均衡
        • 多个订阅
        • 消费位点
        • 消费位点提交
        • 消费位点重置
        • 消息重复和消费幂等
        • 消费失败
        • 消费延迟
        • 消费阻塞以及堆积
        • 提高消费速度
        • 增加 Consumer 实例
        • 增加消费线程
        • 消息过滤
        • 消息广播
        • 订阅关系
        • 通过认证生产与消费加密主题的消息
        • 使用MirrorMaker跨集群数据同步
        • Kafka业务迁移
        • 如何提高消息处理效率
        • Logstash对接Kafka
        • Kafka消费者poll的优化
        • 如何设置消息堆积数超过阈值时,发送告警短信/邮件
        • 消息堆积最佳实践
        • 业务过载最佳实践
        • 业务数据不均衡最佳实践
        • API参考
        • API使用说明
        • 附录
        • 分布式消息服务Kafka资源池
        • SDK参考
        • SDK概述
        • 常见问题
        • 计费与购买类
        • 计费类常见问题
        • 购买类常见问题
        • Kafka支持多可用区?
        • Kafka磁盘选择超高IO还是高IO?
        • 实例问题
        • 实例常见问题
        • 创建的Kafka实例是集群模式么?
        • 连接问题
        • SASL_SSL接入报错
        • 鉴权接入超时问题
        • 连接常见问题
        • 客户端首次接入分布式消息服务Kafka时出现异常的排查方法
        • Kafka实例的连接地址默认有多少个?
        • 如何配置安全组
        • kafka实例连接数有限制吗?
        • 操作类
        • 操作类常见问题
        • 如何配置客户端参数?
        • 如何判断和处理消息堆积?
        • 为什么消费客户端频繁出现Rebalance?
        • 消费端从服务端拉取不到消息或拉取消息缓慢
        • 为什么不推荐使用Sarama Go客户端收发消息?
        • 为什么发送给Topic的消息在分区中分布不均衡?
        • 为什么Group不存在但能消费消息?
        • 消费端挂载NFS是否会影响消费速度?
        • 管理类
        • 相关协议
        • 服务等级协议
        • 服务条款
          无相关产品

          本页目录

          帮助中心分布式消息服务Kafka开发指南Go
          Go
          更新时间 2024-02-05 11:17:13
          • 新浪微博
          • 微信
            扫码分享
          • 复制链接
          最近更新时间: 2024-02-05 11:17:13
          分享文章
          • 新浪微博
          • 微信
            扫码分享
          • 复制链接

          环境准备

          1. 下载Demo包kafka-confluent-go-demo.zip。
          2. 使用开发工具导入Demo。

          配置修改

          1. 如果是ssl连接,需要在控制台下载证书。并且解压压缩包得到ssl.client.truststore.jks,执行以下命令生成caRoot.pem文件。
            keytool -importkeystore -srckeystore ssl.client.truststore.jks -destkeystore caRoot.p12 -deststoretype pkcs12
            openssl pkcs12 -in caRoot.p12 -out caRoot.pem
            
          2. 修改kafka.json文件。(security.protocol仅在ssl连接时需要配置)
            {
              "topic": "XXX",
              "topic2": "XXX",
              "group.id": "XXX",
              "bootstrap.servers" : "XXX:XX",
              "security.protocol" : "SSL"
            }
            

          生产消息

          发送以下命令发送消息。

          go run -mod=vendor producer/producer.go
          

          生产消息示例代码如下:

          package main
          import (
          	"encoding/json"
          	"fmt"
          	"github.com/confluentinc/confluent-kafka-go/kafka"
          	"log"
          	"os"
          	"path/filepath"
          	"strconv"
          	"time"
          )
          const (
          	INT32_MAX = 2147483647 - 1000
          )
          type KafkaConfig struct {
          	Topic      string `json:"topic"`
          	Topic2      string `json:"topic2"`
          	GroupId    string `json:"group.id"`
          	BootstrapServers    string `json:"bootstrap.servers"`
          	SecurityProtocol string `json:"security.protocol"`
          	SslCaLocation string `json:"ssl.ca.location"`
          }
          // config should be a pointer to structure, if not, panic
          func loadJsonConfig() *KafkaConfig {
          	workPath, err := os.Getwd()
          	if err != nil {
          		panic(err)
          	}
          	configPath := filepath.Join(workPath, "conf")
          	fullPath := filepath.Join(configPath, "kafka.json")
          	file, err := os.Open(fullPath);
          	if (err != nil) {
          		msg := fmt.Sprintf("Can not load config at %s. Error: %v", fullPath, err)
          		panic(msg)
          	}
          	defer file.Close()
          	decoder := json.NewDecoder(file)
          	var config = &KafkaConfig{}
          	err = decoder.Decode(config);
          	if (err != nil) {
          		msg := fmt.Sprintf("Decode json fail for config file at %s. Error: %v", fullPath, err)
          		panic(msg)
          	}
          	json.Marshal(config)
          	return  config
          }
          func doInitProducer(cfg *KafkaConfig) *kafka.Producer {
          	fmt.Print("init kafka producer, it may take a few seconds to init the connection\n")
          	//common arguments
          	var kafkaconf = &kafka.ConfigMap{
          		"api.version.request": "true",
          		"message.max.bytes": 1000000,
          		"linger.ms": 500,
          		"sticky.partitioning.linger.ms" : 1000,
          		"retries": INT32_MAX,
          		"retry.backoff.ms": 1000,
          		"acks": "1"}
          	kafkaconf.SetKey("bootstrap.servers", cfg.BootstrapServers)
          	switch cfg.SecurityProtocol {
          	case "PLAINTEXT" :
          		kafkaconf.SetKey("security.protocol", "plaintext");
              case "SSL":
                  kafkaconf.SetKey("security.protocol", "ssl");
                  kafkaconf.SetKey("ssl.ca.location", "/XXX/caRoot.pem")
          	case "SASL_SSL":
          		kafkaconf.SetKey("security.protocol", "sasl_ssl");
          		kafkaconf.SetKey("ssl.ca.location", "/XXX/caRoot.pem");
          		kafkaconf.SetKey("sasl.username", cfg.SaslUsername);
          		kafkaconf.SetKey("sasl.password", cfg.SaslPassword);
          		kafkaconf.SetKey("sasl.mechanism", cfg.SaslMechanism);
                 	kafkaconf.SetKey("enable.ssl.certificate.verification", "false")
          	case "SASL_PLAINTEXT":
          		kafkaconf.SetKey("security.protocol", "sasl_plaintext");
          		kafkaconf.SetKey("sasl.username", cfg.SaslUsername);
          		kafkaconf.SetKey("sasl.password", cfg.SaslPassword);
          		kafkaconf.SetKey("sasl.mechanism", cfg.SaslMechanism)
          	default:
          		panic(kafka.NewError(kafka.ErrUnknownProtocol, "unknown protocol", true))
          	}
          	producer, err := kafka.NewProducer(kafkaconf)
          	if err != nil {
          		panic(err)
          	}
          	fmt.Print("init kafka producer success\n")
          	return producer
          }
          func main() {
          	// Choose the correct protocol
          	cfg := loadJsonConfig();
          	producer := doInitProducer(cfg)
              defer producer.Close()
          	// Delivery report handler for produced messages
          	go func() {
          		for e := range producer.Events() {
          			switch ev := e.(type) {
          			case *kafka.Message:
          				if ev.TopicPartition.Error != nil {
          					log.Printf("Failed to write access log entry:%v", ev.TopicPartition.Error)
          				} else {
          					log.Printf("Send OK topic:%v partition:%v offset:%v content:%s\n", *ev.TopicPartition.Topic,  ev.TopicPartition.Partition, ev.TopicPartition.Offset, ev.Value)
          				}
          			}
          		}
          	}()
              // Produce messages to topic (asynchronously)
          	i := 0
          	for {
          		i = i + 1
          		value := "this is a kafka message from confluent go " + strconv.Itoa(i)
          		var msg *kafka.Message = nil
          		if i % 2 == 0 {
          			msg = &kafka.Message{
          				TopicPartition: kafka.TopicPartition{Topic: &cfg.Topic2, Partition: kafka.PartitionAny},
          				Value:          []byte(value),
          			}
          		} else {
          			msg = &kafka.Message{
          				TopicPartition: kafka.TopicPartition{Topic: &cfg.Topic, Partition: kafka.PartitionAny},
          				Value:          []byte(value),
          			}
          		}
          		producer.Produce(msg, nil)
          		time.Sleep(time.Duration(1) * time.Millisecond)
          	}
          	// Wait for message deliveries before shutting down
          	producer.Flush(15 * 1000)
          }
          
          

          消费消息

          发送以下命令消费消息。

          go run -mod=vendor consumer/consumer.go
          

          消费消息示例代码如下:

          package main
          import (
          	"encoding/json"
          	"fmt"
              "github.com/confluentinc/confluent-kafka-go/kafka"
          	"os"
          	"path/filepath"
          )
          type KafkaConfig struct {
          	Topic      string `json:"topic"`
          	Topic2      string `json:"topic2"`
          	GroupId    string `json:"group.id"`
          	BootstrapServers    string `json:"bootstrap.servers"`
          	SecurityProtocol string `json:"security.protocol"`
          }
          // config should be a pointer to structure, if not, panic
          func loadJsonConfig() *KafkaConfig {
          	workPath, err := os.Getwd()
          	if err != nil {
          		panic(err)
          	}
          	configPath := filepath.Join(workPath, "conf")
          	fullPath := filepath.Join(configPath, "kafka.json")
          	file, err := os.Open(fullPath);
          	if (err != nil) {
          		msg := fmt.Sprintf("Can not load config at %s. Error: %v", fullPath, err)
          		panic(msg)
          	}
          	defer file.Close()
          	decoder := json.NewDecoder(file)
          	var config = &KafkaConfig{}
          	err = decoder.Decode(config);
          	if (err != nil) {
          		msg := fmt.Sprintf("Decode json fail for config file at %s. Error: %v", fullPath, err)
          		panic(msg)
          	}
          	json.Marshal(config)
          	return  config
          }
          func doInitConsumer(cfg *KafkaConfig) *kafka.Consumer {
          	fmt.Print("init kafka consumer, it may take a few seconds to init the connection\n")
          	//common arguments
          	var kafkaconf = &kafka.ConfigMap{
          		"api.version.request": "true",
          		"auto.offset.reset": "latest",
          		"heartbeat.interval.ms": 3000,
          		"session.timeout.ms": 30000,
          		"max.poll.interval.ms": 120000,
          		"fetch.max.bytes": 1024000,
          		"max.partition.fetch.bytes": 256000}
          	kafkaconf.SetKey("bootstrap.servers", cfg.BootstrapServers);
          	kafkaconf.SetKey("group.id", cfg.GroupId)
          	switch cfg.SecurityProtocol {
          	case "PLAINTEXT" :
          		kafkaconf.SetKey("security.protocol", "plaintext");
          	case "SSL":
             		kafkaconf.SetKey("security.protocol", "ssl");
             		kafkaconf.SetKey("ssl.ca.location", "/XXX/caRoot.pem")
          	case "SASL_SSL":
          		kafkaconf.SetKey("security.protocol", "sasl_ssl");
          		kafkaconf.SetKey("ssl.ca.location", "/XXX/caRoot.pem");
          		kafkaconf.SetKey("sasl.username", cfg.SaslUsername);
          		kafkaconf.SetKey("sasl.password", cfg.SaslPassword);
          		kafkaconf.SetKey("sasl.mechanism", cfg.SaslMechanism)
          	case "SASL_PLAINTEXT":
          		kafkaconf.SetKey("security.protocol", "sasl_plaintext");
          		kafkaconf.SetKey("sasl.username", cfg.SaslUsername);
          		kafkaconf.SetKey("sasl.password", cfg.SaslPassword);
          		kafkaconf.SetKey("sasl.mechanism", cfg.SaslMechanism)
          	default:
          		panic(kafka.NewError(kafka.ErrUnknownProtocol, "unknown protocol", true))
          	}
          	consumer, err := kafka.NewConsumer(kafkaconf)
          	if err != nil {
          		panic(err)
          	}
          	fmt.Print("init kafka consumer success\n")
          	return consumer;
          }
          func main() {
          	// Choose the correct protocol
          	cfg := loadJsonConfig();
          	consumer := doInitConsumer(cfg)
          	consumer.SubscribeTopics([]string{cfg.Topic, cfg.Topic2}, nil)
          	for {
          		msg, err := consumer.ReadMessage(-1)
          		if err == nil {
          			fmt.Printf("Message on %s: %s\n", msg.TopicPartition, string(msg.Value))
          		} else {
          			// The client will
          			//automatically try to recover from all errors.
          			fmt.Printf("Consumer error: %v (%v)\n", err, msg)
          		}
          	}
          	consumer.Close()
          }
          
          文档反馈

          建议您登录后反馈,可在建议与反馈里查看问题处理进度

          鼠标选中文档,精准反馈问题

          选中存在疑惑的内容,即可快速反馈问题,我们会跟进处理

          知道了

          上一篇 :  Java客户端接入示例
          下一篇 :  Python
          搜索 关闭
          ©2025 天翼云科技有限公司版权所有 增值电信业务经营许可证A2.B1.B2-20090001
          公司地址:北京市东城区青龙胡同甲1号、3号2幢2层205-32室
          备案 京公网安备11010802043424号 京ICP备 2021034386号
          ©2025天翼云科技有限公司版权所有
          京ICP备 2021034386号
          备案 京公网安备11010802043424号
          增值电信业务经营许可证A2.B1.B2-20090001
          用户协议 隐私政策 法律声明