活动

天翼云最新优惠活动,涵盖免费试用,产品折扣等,助您降本增效!
热门活动
  • 618智算钜惠季 爆款云主机2核4G限时秒杀,88元/年起!
  • 免费体验DeepSeek,上天翼云息壤 NEW 新老用户均可免费体验2500万Tokens,限时两周
  • 云上钜惠 HOT 爆款云主机全场特惠,更有万元锦鲤券等你来领!
  • 算力套餐 HOT 让算力触手可及
  • 天翼云脑AOne NEW 连接、保护、办公,All-in-One!
  • 中小企业服务商合作专区 国家云助力中小企业腾飞,高额上云补贴重磅上线
  • 出海产品促销专区 NEW 爆款云主机低至2折,高性价比,不限新老速来抢购!
  • 天翼云电脑专场 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云生态大会
  • 天翼云中国行
天翼云
  • 活动
  • 智算服务
  • 产品
  • 解决方案
  • 应用商城
  • 合作伙伴
  • 开发者
  • 支持与服务
  • 了解天翼云
      • 文档
      • 控制中心
      • 备案
      • 管理中心
      文档中心

      分布式消息服务RabbitMQ

      分布式消息服务RabbitMQ

        • 产品动态
        • 服务公告
        • 2024
        • 【优惠】正式开放2年7折,3年5折包年折扣
        • 【优惠】分布式消息服务RabbitMQ增加包年优惠折扣和产品资费进一步下调
        • 【降价】分布式消息服务RabbitMQ产品资费价格下调
        • 【通知】云原生引擎调整为白名单特性
        • 【通知】通用型主机规格调整为白名单特性
        • 产品简介
        • 产品定义
        • 产品优势
        • 功能特性
        • 应用场景
        • 分布式消息产品选型
        • 产品规格
        • 安全方案
        • 使用限制
        • 名词解释
        • 与其他服务关系
        • 计费说明
        • 产品资费
        • 新资费
        • 旧资费
        • 计费项
        • 计费模式
        • 续费、到期与欠费
        • 退订
        • 变更配置
        • 快速入门
        • 入门指引
        • 环境准备
        • 购买实例
        • 创建资源
        • 编译工程生产消费
        • 用户指南
        • 创建实例
        • 实例管理
        • 查看实例
        • 实例概览
        • 连接实例
        • 修改实例
        • 实例退订
        • 实例扩容
        • 按需转包周期
        • 虚拟主机管理
        • 创建虚拟主机
        • 查看虚拟主机
        • 删除虚拟主机
        • 用户管理
        • 连接管理
        • 信道管理
        • 操作策略管理
        • 虚拟主机限制管理
        • 交换器管理
        • 队列管理
        • 监控指标
        • 高级特性
        • 惰性队列
        • 消息持久化
        • 死信和TTL
        • Rabbitmq消息确认机制
        • 预取值
        • 心跳检测
        • 单一活跃消费者
        • 仲裁队列
        • 开发指南
        • 概述
        • 收集连接信息
        • Java
        • Python
        • .net
        • PHP
        • Go
        • 性能白皮书
        • RabbitMQ性能白皮书
        • 最佳实践
        • RabbitMQ元数据迁移
        • 如何实现RabbitMQ的高性能
        • RabbitMQ接入
        • 获取SDK
        • 接入方式
        • 代码示例
        • 消息幂等
        • 网络异常自动恢复
        • 节点重启后消费者如何重连
        • 使用AMQProxy解决PHP等客户端Connection复用问题
        • API参考
        • API使用说明
        • SDK参考
        • SDK概述
        • 常见问题
        • 计费类
        • 购买类
        • 操作类
        • 管理类
        • 相关协议
        • 服务等级协议
        • 服务条款
          无相关产品

          本页目录

          帮助中心分布式消息服务RabbitMQ开发指南Go
          Go
          更新时间 2025-07-03 14:33:08
          • 新浪微博
          • 微信
            扫码分享
          • 复制链接
          最近更新时间: 2025-07-03 14:33:08
          分享文章
          • 新浪微博
          • 微信
            扫码分享
          • 复制链接

          编译工程生产消费

          引入依赖

          go.mod

          module test_go
          
          require (
          github.com/rabbitmq/amqp091-go v1.10.0
          golang.org/x/net v0.26.0
          )
          
          go 1.20

          生产消息

          package main
          
          import (
          "flag"
          "fmt"
          amqp "github.com/rabbitmq/amqp091-go"
          "log"
          )
          
          var (
          uri          = flag.String("uri", "amqp://USERNAME:PASSWORD@33.0.1.35:5672", "AMQP URI")
          exchangeName = flag.String("exchange", "test-exchange", "Durable AMQP exchange name")
          exchangeType = flag.String("exchange-type", "direct", "Exchange type - direct|fanout|topic|x-custom")
          routingKey   = flag.String("key", "test-key", "AMQP routing key")
          body         = flag.String("body", "foobar", "Body of message")
          reliable     = flag.Bool("reliable", true, "Wait for the publisher confirmation before exiting")
          )
          
          func init() {
          flag.Parse()
          }
          
          func main() {
          if err := publish(*uri, *exchangeName, *exchangeType, *routingKey, *body, *reliable); err != nil {
          log.Fatalf("%s", err)
          }
          log.Printf("published %dB OK", len(*body))
          }
          
          func publish(amqpURI, exchange, exchangeType, routingKey, body string, reliable bool) error {
          
          // This function dials, connects, declares, publishes, and tears down,
          // all in one go. In a real service, you probably want to maintain a
          // long-lived connection as state, and publish against that.
          
          log.Printf("dialing %q", amqpURI)
          connection, err := amqp.Dial(amqpURI)
          if err != nil {
          return fmt.Errorf("Dial: %s", err)
          }
          defer connection.Close()
          
          log.Printf("got Connection, getting Channel")
          channel, err := connection.Channel()
          if err != nil {
          return fmt.Errorf("Channel: %s", err)
          }
          
          log.Printf("got Channel, declaring %q Exchange (%q)", exchangeType, exchange)
          if err := channel.ExchangeDeclare(
          exchange,     // name
          exchangeType, // type
          true,         // durable
          false,        // auto-deleted
          false,        // internal
          false,        // noWait
          nil,          // arguments
          ); err != nil {
          return fmt.Errorf("Exchange Declare: %s", err)
          }
          
          // Reliable publisher confirms require confirm.select support from the
          // connection.
          if reliable {
          log.Printf("enabling publishing confirms.")
          if err := channel.Confirm(false); err != nil {
          return fmt.Errorf("Channel could not be put into confirm mode: %s", err)
          }
          
          confirms := channel.NotifyPublish(make(chan amqp.Confirmation, 1))
          
          defer confirmOne(confirms)
          }
          
          log.Printf("declared Exchange, publishing %dB body (%q)", len(body), body)
          if err = channel.Publish(
          exchange,   // publish to an exchange
          routingKey, // routing to 0 or more queues
          false,      // mandatory
          false,      // immediate
          amqp.Publishing{
          Headers:         amqp.Table{},
          ContentType:     "text/plain",
          ContentEncoding: "",
          Body:            []byte(body),
          DeliveryMode:    amqp.Transient, // 1=non-persistent, 2=persistent
          Priority:        0,              // 0-9
          // a bunch of application/implementation-specific fields
          },
          ); err != nil {
          return fmt.Errorf("Exchange Publish: %s", err)
          }
          
          return nil
          }
          
          // One would typically keep a channel of publishings, a sequence number, and a
          // set of unacknowledged sequence numbers and loop until the publishing channel
          // is closed.
          func confirmOne(confirms <-chan amqp.Confirmation) {
          log.Printf("waiting for confirmation of one publishing")
          
          if confirmed := <-confirms; confirmed.Ack {
          log.Printf("confirmed delivery with delivery tag: %d", confirmed.DeliveryTag)
          } else {
          log.Printf("failed delivery of delivery tag: %d", confirmed.DeliveryTag)
          }
          }

          消费消息

          package main
          
          import (
          "flag"
          "fmt"
          amqp "github.com/rabbitmq/amqp091-go"
          "log"
          "time"
          )
          
          var (
          uri          = flag.String("uri", "amqp://USERNAME:PASSWORD@10.10.33.196:5672", "AMQP URI")
          exchange     = flag.String("exchange", "test-exchange", "Durable, non-auto-deleted AMQP exchange name")
          exchangeType = flag.String("exchange-type", "direct", "Exchange type - direct|fanout|topic|x-custom")
          queue        = flag.String("queue", "test-queue", "Ephemeral AMQP queue name")
          bindingKey   = flag.String("key", "test-key", "AMQP binding key")
          consumerTag  = flag.String("consumer-tag", "simple-consumer", "AMQP consumer tag (should not be blank)")
          lifetime     = flag.Duration("lifetime", 5*time.Second, "lifetime of process before shutdown (0s=infinite)")
          )
          
          func init() {
          flag.Parse()
          }
          
          func main() {
          c, err := NewConsumer(*uri, *exchange, *exchangeType, *queue, *bindingKey, *consumerTag)
          if err != nil {
          log.Fatalf("%s", err)
          }
          
          if *lifetime > 0 {
          log.Printf("running for %s", *lifetime)
          time.Sleep(*lifetime)
          } else {
          log.Printf("running forever")
          select {}
          }
          
          log.Printf("shutting down")
          
          if err := c.Shutdown(); err != nil {
          log.Fatalf("error during shutdown: %s", err)
          }
          }
          
          type Consumer struct {
          conn    *amqp.Connection
          channel *amqp.Channel
          tag     string
          done    chan error
          }
          
          func NewConsumer(amqpURI, exchange, exchangeType, queueName, key, ctag string) (*Consumer, error) {
          c := &Consumer{
          conn:    nil,
          channel: nil,
          tag:     ctag,
          done:    make(chan error),
          }
          
          var err error
          
          log.Printf("dialing %q", amqpURI)
          c.conn, err = amqp.Dial(amqpURI)
          if err != nil {
          return nil, fmt.Errorf("Dial: %s", err)
          }
          
          go func() {
          fmt.Printf("closing: %s", <-c.conn.NotifyClose(make(chan *amqp.Error)))
          }()
          
          log.Printf("got Connection, getting Channel")
          c.channel, err = c.conn.Channel()
          if err != nil {
          return nil, fmt.Errorf("Channel: %s", err)
          }
          
          log.Printf("got Channel, declaring Exchange (%q)", exchange)
          if err = c.channel.ExchangeDeclare(
          exchange,     // name of the exchange
          exchangeType, // type
          true,         // durable
          false,        // delete when complete
          false,        // internal
          false,        // noWait
          nil,          // arguments
          ); err != nil {
          return nil, fmt.Errorf("Exchange Declare: %s", err)
          }
          
          log.Printf("declared Exchange, declaring Queue %q", queueName)
          queue, err := c.channel.QueueDeclare(
          queueName, // name of the queue
          true,      // durable
          false,     // delete when unused
          false,     // exclusive
          false,     // noWait
          nil,       // arguments
          )
          if err != nil {
          return nil, fmt.Errorf("Queue Declare: %s", err)
          }
          
          log.Printf("declared Queue (%q %d messages, %d consumers), binding to Exchange (key %q)",
          queue.Name, queue.Messages, queue.Consumers, key)
          
          if err = c.channel.QueueBind(
          queue.Name, // name of the queue
          key,        // bindingKey
          exchange,   // sourceExchange
          false,      // noWait
          nil,        // arguments
          ); err != nil {
          return nil, fmt.Errorf("Queue Bind: %s", err)
          }
          
          log.Printf("Queue bound to Exchange, starting Consume (consumer tag %q)", c.tag)
          deliveries, err := c.channel.Consume(
          queue.Name, // name
          c.tag,      // consumerTag,
          false,      // noAck
          false,      // exclusive
          false,      // noLocal
          false,      // noWait
          nil,        // arguments
          )
          if err != nil {
          return nil, fmt.Errorf("Queue Consume: %s", err)
          }
          
          go handle(deliveries, c.done)
          
          return c, nil
          }
          
          func (c *Consumer) Shutdown() error {
          // will close() the deliveries channel
          if err := c.channel.Cancel(c.tag, true); err != nil {
          return fmt.Errorf("Consumer cancel failed: %s", err)
          }
          
          if err := c.conn.Close(); err != nil {
          return fmt.Errorf("AMQP connection close error: %s", err)
          }
          
          defer log.Printf("AMQP shutdown OK")
          
          // wait for handle() to exit
          return <-c.done
          }
          
          func handle(deliveries <-chan amqp.Delivery, done chan error) {
          for d := range deliveries {
          log.Printf(
          "got %dB delivery: [%v] %q",
          len(d.Body),
          d.DeliveryTag,
          d.Body,
          )
          d.Ack(false)
          }
          log.Printf("handle: deliveries channel closed")
          done <- nil
          }

          ssl生产消息

          package main
          
          import (
          "crypto/tls"
          "crypto/x509"
          "flag"
          "fmt"
          amqp "github.com/rabbitmq/amqp091-go"
          "io/ioutil"
          "log"
          )
          
          var (
          uri          = flag.String("uri", "amqps://USERNAME:PASSWORD@10.10.33.196:5671", "AMQP URI")
          exchangeName = flag.String("exchange", "go-exchange", "Durable AMQP exchange name")
          exchangeType = flag.String("exchange-type", "direct", "Exchange type - direct|fanout|topic|x-custom")
          routingKey   = flag.String("key", "test-key", "AMQP routing key")
          body         = flag.String("body", "foobar", "Body of message")
          reliable     = flag.Bool("reliable", true, "Wait for the publisher confirmation before exiting")
          )
          
          func init() {
          flag.Parse()
          }
          
          func main() {
          if err := publish(*uri, *exchangeName, *exchangeType, *routingKey, *body, *reliable); err != nil {
          log.Fatalf("%s", err)
          }
          log.Printf("published %dB OK", len(*body))
          }
          
          func publish(amqpsURI, exchange, exchangeType, routingKey, body string, reliable bool) error {
          caCert, err := ioutil.ReadFile("D:\\tmp\\hzmq-test-0520_rabbitmq_ssl_client\\ca_certificate.pem")
          if err != nil {
          return err
          }
          
          cert, err := tls.LoadX509KeyPair("D:\\tmp\\hzmq-test-0520_rabbitmq_ssl_client\\client_rabbitmq_certificate.pem", "D:\\tmp\\hzmq-test-0520_rabbitmq_ssl_client\\client_rabbitmq_key.pem")
          if err != nil {
          return err
          }
          
          rootCAs := x509.NewCertPool()
          rootCAs.AppendCertsFromPEM(caCert)
          
          tlsConf := &tls.Config{
          RootCAs:      rootCAs,
          Certificates: []tls.Certificate{cert},
          //ServerName:   "localhost", // Optional
          InsecureSkipVerify: true,
          }
          
          connection, err := amqp.DialTLS(amqpsURI, tlsConf)
          if err != nil {
          return fmt.Errorf("Dial: %s", err)
          }
          defer connection.Close()
          
          log.Printf("got Connection, getting Channel")
          channel, err := connection.Channel()
          if err != nil {
          return fmt.Errorf("Channel: %s", err)
          }
          
          log.Printf("got Channel, declaring %q Exchange (%q)", exchangeType, exchange)
          if err := channel.ExchangeDeclare(
          exchange,     // name
          exchangeType, // type
          true,         // durable
          false,        // auto-deleted
          false,        // internal
          false,        // noWait
          nil,          // arguments
          ); err != nil {
          return fmt.Errorf("Exchange Declare: %s", err)
          }
          
          // Reliable publisher confirms require confirm.select support from the
          // connection.
          if reliable {
          log.Printf("enabling publishing confirms.")
          if err := channel.Confirm(false); err != nil {
          return fmt.Errorf("Channel could not be put into confirm mode: %s", err)
          }
          
          confirms := channel.NotifyPublish(make(chan amqp.Confirmation, 1))
          
          defer confirmOne(confirms)
          }
          
          log.Printf("declared Exchange, publishing %dB body (%q)", len(body), body)
          if err = channel.Publish(
          exchange,   // publish to an exchange
          routingKey, // routing to 0 or more queues
          false,      // mandatory
          false,      // immediate
          amqp.Publishing{
          Headers:         amqp.Table{},
          ContentType:     "text/plain",
          ContentEncoding: "",
          Body:            []byte(body),
          DeliveryMode:    amqp.Transient, // 1=non-persistent, 2=persistent
          Priority:        0,              // 0-9
          },
          ); err != nil {
          return fmt.Errorf("Exchange Publish: %s", err)
          }
          
          return nil
          }
          
          func confirmOne(confirms <-chan amqp.Confirmation) {
          log.Printf("waiting for confirmation of one publishing")
          
          if confirmed := <-confirms; confirmed.Ack {
          log.Printf("confirmed delivery with delivery tag: %d", confirmed.DeliveryTag)
          } else {
          log.Printf("failed delivery of delivery tag: %d", confirmed.DeliveryTag)
          }
          }

          ssl消费消息

          package main
          
          import (
          "crypto/tls"
          "crypto/x509"
          "flag"
          "fmt"
          amqp "github.com/rabbitmq/amqp091-go"
          "io/ioutil"
          "log"
          "time"
          )
          
          var (
          uri          = flag.String("uri", "amqps://USERNAME:PASSWORD@10.10.33.196:5671", "AMQP URI")
          exchange     = flag.String("exchange", "test-exchange", "Durable, non-auto-deleted AMQP exchange name")
          exchangeType = flag.String("exchange-type", "direct", "Exchange type - direct|fanout|topic|x-custom")
          queue        = flag.String("queue", "test-queue", "Ephemeral AMQP queue name")
          bindingKey   = flag.String("key", "test-key", "AMQP binding key")
          consumerTag  = flag.String("consumer-tag", "simple-consumer", "AMQP consumer tag (should not be blank)")
          lifetime     = flag.Duration("lifetime", 5*time.Second, "lifetime of process before shutdown (0s=infinite)")
          )
          
          func init() {
          flag.Parse()
          }
          
          func main() {
          c, err := NewConsumer(*uri, *exchange, *exchangeType, *queue, *bindingKey, *consumerTag)
          if err != nil {
          log.Fatalf("%s", err)
          }
          
          if *lifetime > 0 {
          log.Printf("running for %s", *lifetime)
          time.Sleep(*lifetime)
          } else {
          log.Printf("running forever")
          select {}
          }
          
          log.Printf("shutting down")
          
          if err := c.Shutdown(); err != nil {
          log.Fatalf("error during shutdown: %s", err)
          }
          }
          
          type Consumer struct {
          conn    *amqp.Connection
          channel *amqp.Channel
          tag     string
          done    chan error
          }
          
          func NewConsumer(amqpsURI, exchange, exchangeType, queueName, key, ctag string) (*Consumer, error) {
          c := &Consumer{
          conn:    nil,
          channel: nil,
          tag:     ctag,
          done:    make(chan error),
          }
          
          var err error
          
          log.Printf("dialing %q", amqpsURI)
          
          caCert, err := ioutil.ReadFile("D:\\codeTest\\amqp\\_examples\\ssl\\ca_certificate.pem")
          if err != nil {
          return nil, err
          }
          
          cert, err := tls.LoadX509KeyPair("D:\\codeTest\\amqp\\_examples\\ssl\\client_openstack_certificate.pem", "D:\\codeTest\\amqp\\_examples\\ssl\\client_openstack_key.pem")
          if err != nil {
          return nil, err
          }
          
          rootCAs := x509.NewCertPool()
          rootCAs.AppendCertsFromPEM(caCert)
          
          tlsConf := &tls.Config{
          RootCAs:      rootCAs,
          Certificates: []tls.Certificate{cert},
          //ServerName:   "localhost", // Optional
          InsecureSkipVerify: true,
          }
          
          c.conn, err = amqp.DialTLS(amqpsURI, tlsConf)
          if err != nil {
          return nil, fmt.Errorf("Dial: %s", err)
          }
          
          go func() {
          fmt.Printf("closing: %s", <-c.conn.NotifyClose(make(chan *amqp.Error)))
          }()
          
          log.Printf("got Connection, getting Channel")
          c.channel, err = c.conn.Channel()
          if err != nil {
          return nil, fmt.Errorf("Channel: %s", err)
          }
          
          log.Printf("got Channel, declaring Exchange (%q)", exchange)
          if err = c.channel.ExchangeDeclare(
          exchange,     // name of the exchange
          exchangeType, // type
          true,         // durable
          false,        // delete when complete
          false,        // internal
          false,        // noWait
          nil,          // arguments
          ); err != nil {
          return nil, fmt.Errorf("Exchange Declare: %s", err)
          }
          
          log.Printf("declared Exchange, declaring Queue %q", queueName)
          queue, err := c.channel.QueueDeclare(
          queueName, // name of the queue
          true,      // durable
          false,     // delete when unused
          false,     // exclusive
          false,     // noWait
          nil,       // arguments
          )
          if err != nil {
          return nil, fmt.Errorf("Queue Declare: %s", err)
          }
          
          log.Printf("declared Queue (%q %d messages, %d consumers), binding to Exchange (key %q)",
          queue.Name, queue.Messages, queue.Consumers, key)
          
          if err = c.channel.QueueBind(
          queue.Name, // name of the queue
          key,        // bindingKey
          exchange,   // sourceExchange
          false,      // noWait
          nil,        // arguments
          ); err != nil {
          return nil, fmt.Errorf("Queue Bind: %s", err)
          }
          
          log.Printf("Queue bound to Exchange, starting Consume (consumer tag %q)", c.tag)
          deliveries, err := c.channel.Consume(
          queue.Name, // name
          c.tag,      // consumerTag,
          false,      // noAck
          false,      // exclusive
          false,      // noLocal
          false,      // noWait
          nil,        // arguments
          )
          if err != nil {
          return nil, fmt.Errorf("Queue Consume: %s", err)
          }
          
          go handle(deliveries, c.done)
          
          return c, nil
          }
          
          func (c *Consumer) Shutdown() error {
          // will close() the deliveries channel
          if err := c.channel.Cancel(c.tag, true); err != nil {
          return fmt.Errorf("Consumer cancel failed: %s", err)
          }
          
          if err := c.conn.Close(); err != nil {
          return fmt.Errorf("AMQP connection close error: %s", err)
          }
          
          defer log.Printf("AMQP shutdown OK")
          
          // wait for handle() to exit
          return <-c.done
          }
          
          func handle(deliveries <-chan amqp.Delivery, done chan error) {
          for d := range deliveries {
          log.Printf(
          "got %dB delivery: [%v] %q",
          len(d.Body),
          d.DeliveryTag,
          d.Body,
          )
          d.Ack(false)
          }
          log.Printf("handle: deliveries channel closed")
          done <- nil
          }
          文档反馈

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

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

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

          知道了

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