目录

RabbitMQ 全面指南:从基础到实战

目录

1. RabbitMQ 简介

RabbitMQ 是基于 AMQP(Advanced Message Queuing Protocol,高级消息队列协议)规范实现的开源消息中间件,由 Erlang 语言编写,天生具备高并发和高可用能力。

1.1 为什么需要消息队列

场景问题 MQ 解决方式
服务间强耦合 异步解耦,发布/订阅模式
流量洪峰打垮下游 削峰填谷,控制消费速率
同步调用链过长 异步化,减少用户等待时间
多系统同步数据 广播/扇出,一发多收

1.2 常见 MQ 对比

特性 RabbitMQ Kafka RocketMQ NSQ
协议 AMQP 自研 自研 自研
语言 Erlang Scala/Java Java Go
吞吐量 万级 百万级 十万级 十万级
时序消息 不支持 分区内有序 支持 不支持
延迟队列 插件支持 不支持 原生支持 不支持
消息堆积 一般 极强 强 一般
管理界面 好 需插件 好 简单
适用场景 业务解耦、可靠投递 日志、流处理 金融、电商 轻量异步

选型建议:复杂路由、可靠投递优先选 RabbitMQ;超高吞吐、日志流处理选 Kafka;业务消息、顺序消息选 RocketMQ。

1.3 RabbitMQ 核心优势

  • 灵活路由:多种 Exchange 类型,支持复杂的消息路由规则
  • 可靠投递:持久化、确认机制(Publisher Confirm + Consumer ACK)
  • 插件生态:死信队列、延迟队列、Federation 等插件丰富
  • 管理界面:内置 Web UI,运维友好
  • 多语言客户端:官方支持 Go、Java、Python、Ruby 等

2. 核心概念

2.1 整体架构

Producer                    RabbitMQ Broker                    Consumer
   │                    ┌───────────────────────┐                 │
   │  Publish(msg)      │  Virtual Host (vhost) │  Subscribe      │
   │─────────────────►  │                       │ ──────────────► │
   │                    │  Exchange ──► Queue   │                 │
   │                    │    │    Binding  │    │                 │
   │                    └───────────────────────┘                 │
   │                              │                               │
   │◄─────────────────────────────│───────────────────────────────│
   │     Publisher Confirm        │       Consumer ACK            │

2.2 核心术语

术语 说明
Broker RabbitMQ 服务实体
Virtual Host (vhost) 逻辑隔离单元,类似数据库的 database。默认 /
Connection 客户端与 Broker 的 TCP 长连接
Channel(信道) 建立在 Connection 上的逻辑通道,复用 TCP 连接,避免频繁建立 TCP
Exchange(交换器) 接收生产者消息,按路由规则分发到队列
Queue(队列) 消息的存储容器,消费者从队列拉取消息
Binding(绑定) Exchange 与 Queue 的关联规则,由 Binding Key 定义
Routing Key 生产者发消息时携带的路由键,Exchange 根据它路由消息
Message 消息体(payload)+ 消息头(headers/attributes)

2.3 信道(Channel)的意义

TCP 连接的创建/销毁代价高昂,且操作系统对并发 TCP 数有限制。Channel 是建立在同一条 TCP 连接上的虚拟通道:

TCP Connection
├── Channel 1  ← goroutine A
├── Channel 2  ← goroutine B
├── Channel 3  ← goroutine C
└── ...

每个 Channel 有唯一 ID,独立进行 AMQP 操作。Channel 不是线程安全的,每个 goroutine 应使用独立的 Channel。

2.4 消息结构

Message
├── Properties (消息属性)
│   ├── content_type      # "application/json"
│   ├── content_encoding  # "gzip"
│   ├── delivery_mode     # 1=非持久, 2=持久化
│   ├── priority          # 0-9 优先级
│   ├── correlation_id    # RPC 回调关联 ID
│   ├── reply_to          # RPC 回调队列
│   ├── expiration        # 消息 TTL(毫秒字符串)
│   ├── message_id        # 消息唯一 ID
│   ├── timestamp         # 发送时间戳
│   └── headers           # 自定义 K-V 头
└── Body (消息体,[]byte)

3. Exchange 类型详解

Exchange 是路由的核心,共四种类型。

3.1 Direct Exchange(直连)

精确匹配:Routing Key 与 Binding Key 完全一致才投递。

Producer ──[routing_key="error"]──► Exchange(direct)
                                         │
                          ┌──────────────┼──────────────┐
                    [BK="info"]     [BK="error"]    [BK="warn"]
                          │              │              │
                       Queue-A        Queue-B        Queue-C
                     (不收到)        (收到!)        (不收到)

Go 示例:

// 声明 Direct Exchange
err = ch.ExchangeDeclare(
    "logs_direct", // name
    "direct",      // type
    true,          // durable
    false,         // auto-delete
    false,         // internal
    false,         // no-wait
    nil,
)

// 发布消息,指定 routing key
err = ch.Publish(
    "logs_direct", // exchange
    "error",       // routing key
    false,         // mandatory
    false,         // immediate
    amqp.Publishing{
        ContentType:  "text/plain",
        DeliveryMode: amqp.Persistent,
        Body:         []byte("disk full!"),
    },
)

典型场景:日志分级处理、精确任务分发。

3.2 Fanout Exchange(扇出/广播)

无视 Routing Key,将消息广播到所有绑定的队列。

Producer ──[any_key]──► Exchange(fanout)
                              │
                 ┌────────────┼────────────┐
                 │            │            │
              Queue-A      Queue-B      Queue-C
             (都收到!)    (都收到!)    (都收到!)

Go 示例:

err = ch.ExchangeDeclare("notifications", "fanout", true, false, false, false, nil)

// Binding 时 routing key 填空字符串(fanout 忽略)
err = ch.QueueBind(q.Name, "", "notifications", false, nil)

典型场景:用户注册后同时触发欢迎邮件、发券、积分三个服务;图片上传后同时清缓存和发奖励。

3.3 Topic Exchange(主题)

通配符匹配,Routing Key 和 Binding Key 都是以 . 分隔的单词串。

通配符 含义
* 匹配一个单词
# 匹配零个或多个单词
Routing Key: "order.created.vip"

Binding Key 匹配情况:
  "order.*.*"       ✓  (两个*各匹配一个单词)
  "order.#"         ✓  (#匹配 created.vip)
  "#"               ✓  (匹配所有)
  "order.created.*" ✓
  "*.created.*"     ✓
  "order.*"         ✗  (只有一个*,无法匹配 created.vip)
  "order.created"   ✗  (精确匹配失败)

Go 示例:

err = ch.ExchangeDeclare("topic_logs", "topic", true, false, false, false, nil)

// 订阅所有 error 级别的日志
err = ch.QueueBind(q.Name, "#.error", "topic_logs", false, nil)

// 订阅 kern 模块的所有日志
err = ch.QueueBind(q.Name, "kern.*", "topic_logs", false, nil)

典型场景:多维度消息过滤,如按地区+业务类型+级别订阅。

3.4 Headers Exchange(头匹配)

根据消息 headers 的 K-V 对匹配,而非 Routing Key。性能较差,几乎不用。

// x-match: "all" 表示所有 header 都必须匹配
// x-match: "any" 表示匹配任意一个即可
err = ch.QueueBind(q.Name, "", "headers_exchange", false, amqp.Table{
    "x-match": "all",
    "format":  "pdf",
    "type":    "report",
})

3.5 默认 Exchange

每个 vhost 有一个名称为空字符串 "" 的默认 Direct Exchange,所有队列自动绑定,Binding Key 等于队列名。

// 直接发送到队列,不需要额外声明 exchange
err = ch.Publish(
    "",         // exchange: 空字符串 = 默认 exchange
    "my-queue", // routing key = 队列名
    false, false,
    amqp.Publishing{Body: []byte("hello")},
)

4. 消息可靠性保障

消息丢失有三个环节:生产者 → Broker → 消费者,每个环节都需要保障。

生产者                    Broker                     消费者
   │                        │                           │
   │  ① Producer Confirm   │  ② Message Persistence   │  ③ Consumer Confirm
   │  Publisher Confirm     │  Exchange + Queue + Msg   │  Consumer ACK
   │  (avoid loss to Broker)│  (survive Broker restart) │  (avoid consume failure)

4.1 生产者确认(Publisher Confirm)

将 Channel 设置为 Confirm 模式,每条消息发送后 Broker 会回发 ACK 或 NACK。

package main

import (
    "fmt"
    "log"

    amqp "github.com/rabbitmq/amqp091-go"
)

func reliablePublish(ch *amqp.Channel, body []byte) error {
    // 开启 confirm 模式
    if err := ch.Confirm(false); err != nil {
        return fmt.Errorf("confirm mode: %w", err)
    }

    // 注册确认监听(异步方式)
    confirms := ch.NotifyPublish(make(chan amqp.Confirmation, 1))

    err := ch.Publish(
        "my-exchange",
        "my-key",
        true,  // mandatory: 找不到队列时返回 Return 事件
        false,
        amqp.Publishing{
            DeliveryMode: amqp.Persistent, // 持久化消息
            ContentType:  "application/json",
            Body:         body,
        },
    )
    if err != nil {
        return fmt.Errorf("publish: %w", err)
    }

    // 等待确认
    if confirmed := <-confirms; !confirmed.Ack {
        return fmt.Errorf("message nacked by broker, tag=%d", confirmed.DeliveryTag)
    }
    return nil
}

注意:Confirm 模式是 Channel 级别的,开启后当前 Channel 的所有消息都会走确认流程。

关于 AMQP 事务:AMQP 还提供了事务机制(txSelect / txCommit / txRollback)来保证投递可靠性,但它是同步阻塞的——每条消息都要等一次完整的提交往返,吞吐量比 Confirm 模式低一到两个数量级。生产环境一律用 Confirm,不要用事务。

4.2 消息持久化

持久化需要三者同时满足,缺一不可:

// 1. 持久化 Exchange
ch.ExchangeDeclare("my-exchange", "direct",
    true,  // durable = true
    false, false, false, nil)

// 2. 持久化 Queue
ch.QueueDeclare("my-queue",
    true,  // durable = true
    false, false, false, nil)

// 3. 持久化 Message
ch.Publish("my-exchange", "my-key", false, false, amqp.Publishing{
    DeliveryMode: amqp.Persistent, // delivery_mode = 2
    Body:         []byte("data"),
})

性能权衡:持久化消息需写入磁盘,吞吐量约降低 10x。对于非关键消息可不持久化。

4.3 消费者确认(Consumer ACK)

RabbitMQ 默认 auto-ack,消息投递后立即删除。生产环境必须手动 ACK。

msgs, err := ch.Consume(
    "my-queue",
    "consumer-1", // consumer tag
    false,        // auto-ack = false,手动确认
    false, false, false, nil,
)

for msg := range msgs {
    if err := processMessage(msg.Body); err != nil {
        // 处理失败:requeue=true 重新入队,requeue=false 进死信队列
        msg.Nack(false, true)
        continue
    }
    // 处理成功:确认消息
    msg.Ack(false) // multiple=false 只确认当前消息
}

ACK 相关 API:

方法 说明
Ack(multiple bool) 确认消息,multiple=true 批量确认 deliveryTag 以下的所有消息
Nack(multiple, requeue bool) 拒绝消息,requeue=true 重新入队,false 丢弃或进死信队列
Reject(requeue bool) 拒绝单条消息(等价于 Nack 的 multiple=false)

消息状态流转:队列中的消息有 Ready 和 Unacked 两种状态。消息投递给消费者后变为 Unacked,此时 Broker 认为该消费者正忙,不会再向它投递新消息,直到:

  • 消费者 ACK/Nack 该消息 → 消息删除或重新入队
  • 消费者断开连接 → 消息自动回到 Ready,重新分发给其他消费者

所以代码里漏写 ACK 不会报错,但该消费者会在处理完第一批消息后静默卡死——这也是排查"消费者不消费了"时首先要看 Unacked 数量的原因。

4.4 QoS(服务质量/预取数量)

防止消费者拉取太多消息堆积在内存中,避免"慢消费者"问题:

// 每次最多预取 10 条,处理完 ACK 后才拉取新消息
err = ch.Qos(
    10,    // prefetch count
    0,     // prefetch size (字节,0=不限)
    false, // global: false=per consumer, true=per channel
)

4.5 mandatory 与 Return 事件

mandatory=true 时,若 Exchange 找不到匹配的队列,消息会通过 Basic.Return 回退给生产者:

returns := ch.NotifyReturn(make(chan amqp.Return, 1))

ch.Publish("my-exchange", "non-existent-key",
    true, false, // mandatory=true
    amqp.Publishing{Body: []byte("test")},
)

select {
case ret := <-returns:
    log.Printf("消息被退回: replyCode=%d, replyText=%s", ret.ReplyCode, ret.ReplyText)
case confirm := <-confirms:
    if confirm.Ack {
        log.Println("消息投递成功")
    }
}

5. 死信队列(Dead Letter Queue)

当消息无法被正常消费时,可以将其路由到死信交换器(DLX),再投入死信队列进行后续处理。

5.1 消息成为死信的三种情况

  1. 消息被拒绝:Nack 或 Reject 且 requeue=false
  2. 消息 TTL 过期:消息在队列中超过设定的存活时间
  3. 队列长度溢出:队列达到 x-max-length 限制,头部消息被丢弃

5.2 配置死信队列

// 第一步:声明死信交换器和死信队列
ch.ExchangeDeclare("dlx.exchange", "direct", true, false, false, false, nil)
ch.QueueDeclare("dlx.queue", true, false, false, false, nil)
ch.QueueBind("dlx.queue", "dlx.key", "dlx.exchange", false, nil)

// 第二步:声明业务队列时绑定 DLX
args := amqp.Table{
    "x-dead-letter-exchange":    "dlx.exchange",  // 死信交换器
    "x-dead-letter-routing-key": "dlx.key",       // 死信路由键(可选)
    "x-message-ttl":             int32(30000),     // 消息 TTL 30s
    "x-max-length":              int32(1000),      // 队列最大长度
}
ch.QueueDeclare("business.queue", true, false, false, false, args)

5.3 死信消息的元数据

死信消息的 headers 中会附加 x-death 字段,记录死信历史:

{
  "x-death": [
    {
      "count": 1,
      "exchange": "business.exchange",
      "queue": "business.queue",
      "reason": "expired",       // rejected / expired / maxlen
      "routing-keys": ["order"],
      "time": "2026-03-30T10:00:00Z"
    }
  ],
  "x-first-death-exchange": "business.exchange",
  "x-first-death-queue": "business.queue",
  "x-first-death-reason": "expired"
}

6. 延迟队列

RabbitMQ 没有原生的延迟队列,有两种实现方式。

6.1 方案一:TTL + 死信队列(无插件)

利用消息 TTL 过期后进入死信队列的机制模拟延迟:

Producer ──► [normal.queue(TTL=30s, DLX=dlx)] ──(过期)──► DLX ──► [delay.queue]
                                                                         │
                                                                     Consumer (30s后消费)
// 发送延迟消息(每条消息可设置不同的TTL)
func sendDelayMessage(ch *amqp.Channel, body []byte, delayMs int) error {
    return ch.Publish(
        "",             // 使用默认 exchange
        "normal.queue", // 路由到带 TTL 的队列
        false, false,
        amqp.Publishing{
            DeliveryMode: amqp.Persistent,
            Expiration:   fmt.Sprintf("%d", delayMs), // 消息级别 TTL
            Body:         body,
        },
    )
}

缺点:队列级 TTL 只能设置统一过期时间;消息级 TTL 时,队列头部未过期的消息会阻塞后面已过期的消息(因为 RabbitMQ 只扫描队头)。

6.2 方案二:rabbitmq-delayed-message-exchange 插件(推荐)

# 安装插件
rabbitmq-plugins enable rabbitmq_delayed_message_exchange
// 声明 x-delayed-message 类型的 Exchange
err = ch.ExchangeDeclare(
    "delayed.exchange",
    "x-delayed-message", // 插件提供的类型
    true, false, false, false,
    amqp.Table{
        "x-delayed-type": "direct", // 实际路由类型
    },
)

// 发送延迟消息,在 header 中指定延迟毫秒数
err = ch.Publish(
    "delayed.exchange",
    "order.timeout",
    false, false,
    amqp.Publishing{
        DeliveryMode: amqp.Persistent,
        Headers: amqp.Table{
            "x-delay": int32(30000), // 延迟 30 秒
        },
        Body: []byte(`{"orderId":"12345"}`),
    },
)

优点:每条消息独立延迟时间,精确可控,无阻塞问题。


7. 消息优先级队列

// 声明支持优先级的队列(优先级范围 0-10)
ch.QueueDeclare("priority.queue", true, false, false, false, amqp.Table{
    "x-max-priority": int32(10),
})

// 发送高优先级消息
ch.Publish("", "priority.queue", false, false, amqp.Publishing{
    Priority: 9, // 数值越大优先级越高
    Body:     []byte("urgent task"),
})

// 发送低优先级消息
ch.Publish("", "priority.queue", false, false, amqp.Publishing{
    Priority: 1,
    Body:     []byte("normal task"),
})

注意:优先级只在有消息堆积时才有意义,消费速度远大于生产速度时优先级无效。


8. 完整 Go 实战案例

8.1 项目结构

rabbitmq-demo/
├── go.mod
├── pkg/
│   └── mq/
│       ├── connection.go    # 连接管理(自动重连)
│       ├── producer.go      # 生产者
│       └── consumer.go      # 消费者
├── examples/
│   ├── work_queue/          # 工作队列
│   ├── pub_sub/             # 发布订阅
│   └── rpc/                 # RPC 模式
└── main.go

8.2 连接管理(自动重连)

// pkg/mq/connection.go
package mq

import (
    "fmt"
    "log"
    "time"

    amqp "github.com/rabbitmq/amqp091-go"
)

type Connection struct {
    conn    *amqp.Connection
    url     string
    closed  chan struct{}
}

func NewConnection(url string) (*Connection, error) {
    c := &Connection{
        url:    url,
        closed: make(chan struct{}),
    }
    if err := c.connect(); err != nil {
        return nil, err
    }
    go c.reconnectLoop()
    return c, nil
}

func (c *Connection) connect() error {
    conn, err := amqp.Dial(c.url)
    if err != nil {
        return fmt.Errorf("dial: %w", err)
    }
    c.conn = conn
    log.Println("RabbitMQ connected")
    return nil
}

func (c *Connection) reconnectLoop() {
    for {
        reason, ok := <-c.conn.NotifyClose(make(chan *amqp.Error))
        if !ok {
            // 主动关闭
            select {
            case <-c.closed:
                return
            default:
            }
        }
        log.Printf("connection closed, reason: %v, reconnecting...", reason)

        // 指数退避重连
        delay := time.Second
        for {
            if err := c.connect(); err == nil {
                break
            }
            log.Printf("reconnect failed, retry in %s", delay)
            time.Sleep(delay)
            if delay < 30*time.Second {
                delay *= 2
            }
        }
    }
}

func (c *Connection) Channel() (*amqp.Channel, error) {
    return c.conn.Channel()
}

func (c *Connection) Close() {
    close(c.closed)
    c.conn.Close()
}

8.3 生产者(可靠投递)

// pkg/mq/producer.go
package mq

import (
    "context"
    "fmt"
    "time"

    amqp "github.com/rabbitmq/amqp091-go"
)

type Producer struct {
    ch       *amqp.Channel
    confirms chan amqp.Confirmation
}

func NewProducer(conn *Connection) (*Producer, error) {
    ch, err := conn.Channel()
    if err != nil {
        return nil, fmt.Errorf("open channel: %w", err)
    }

    // 开启 Publisher Confirm 模式
    if err = ch.Confirm(false); err != nil {
        return nil, fmt.Errorf("confirm mode: %w", err)
    }

    confirms := ch.NotifyPublish(make(chan amqp.Confirmation, 128))
    return &Producer{ch: ch, confirms: confirms}, nil
}

// Publish 发布消息,等待 broker 确认
func (p *Producer) Publish(ctx context.Context, exchange, routingKey string, body []byte) error {
    err := p.ch.Publish(
        exchange,
        routingKey,
        false, false,
        amqp.Publishing{
            ContentType:  "application/json",
            DeliveryMode: amqp.Persistent,
            Timestamp:    time.Now(),
            Body:         body,
        },
    )
    if err != nil {
        return fmt.Errorf("publish: %w", err)
    }

    select {
    case confirm := <-p.confirms:
        if !confirm.Ack {
            return fmt.Errorf("broker nacked message tag=%d", confirm.DeliveryTag)
        }
        return nil
    case <-ctx.Done():
        return ctx.Err()
    case <-time.After(5 * time.Second):
        return fmt.Errorf("confirm timeout")
    }
}

func (p *Producer) Close() {
    p.ch.Close()
}

8.4 消费者(并发处理)

// pkg/mq/consumer.go
package mq

import (
    "context"
    "fmt"
    "log"
    "sync"

    amqp "github.com/rabbitmq/amqp091-go"
)

type HandlerFunc func(ctx context.Context, msg amqp.Delivery) error

type Consumer struct {
    ch          *amqp.Channel
    queue       string
    concurrency int
}

func NewConsumer(conn *Connection, queue string, concurrency int) (*Consumer, error) {
    ch, err := conn.Channel()
    if err != nil {
        return nil, fmt.Errorf("open channel: %w", err)
    }

    // 每次最多预取 concurrency 条消息
    if err = ch.Qos(concurrency, 0, false); err != nil {
        return nil, fmt.Errorf("qos: %w", err)
    }

    return &Consumer{ch: ch, queue: queue, concurrency: concurrency}, nil
}

func (c *Consumer) Start(ctx context.Context, handler HandlerFunc) error {
    msgs, err := c.ch.Consume(
        c.queue,
        "",    // 消费者 tag,空则由 broker 生成
        false, // auto-ack = false
        false, false, false, nil,
    )
    if err != nil {
        return fmt.Errorf("consume: %w", err)
    }

    // 用信号量控制并发度
    sem := make(chan struct{}, c.concurrency)
    var wg sync.WaitGroup

    for {
        select {
        case <-ctx.Done():
            wg.Wait()
            return ctx.Err()
        case msg, ok := <-msgs:
            if !ok {
                wg.Wait()
                return fmt.Errorf("channel closed")
            }

            sem <- struct{}{}
            wg.Add(1)
            go func(m amqp.Delivery) {
                defer func() {
                    <-sem
                    wg.Done()
                }()

                if err := handler(ctx, m); err != nil {
                    log.Printf("handler error: %v, nacking message", err)
                    m.Nack(false, true) // requeue=true 重新入队
                    return
                }
                m.Ack(false)
            }(msg)
        }
    }
}

func (c *Consumer) Close() {
    c.ch.Close()
}

8.5 工作队列案例(任务分发)

场景:多个 Worker 分摊图片压缩任务,确保每个任务只被处理一次。

// examples/work_queue/main.go
package main

import (
    "context"
    "encoding/json"
    "fmt"
    "log"
    "os"
    "os/signal"
    "syscall"
    "time"

    amqp "github.com/rabbitmq/amqp091-go"
    "github.com/your/project/pkg/mq"
)

const (
    queueName = "image.compress"
    amqpURL   = "amqp://guest:guest@localhost:5672/"
)

type ImageTask struct {
    ImageID string `json:"image_id"`
    URL     string `json:"url"`
    Width   int    `json:"width"`
    Height  int    `json:"height"`
}

func main() {
    conn, err := mq.NewConnection(amqpURL)
    if err != nil {
        log.Fatal(err)
    }
    defer conn.Close()

    // 根据命令行参数决定是生产者还是消费者
    if len(os.Args) > 1 && os.Args[1] == "produce" {
        runProducer(conn)
    } else {
        runConsumer(conn)
    }
}

func runProducer(conn *mq.Connection) {
    producer, err := mq.NewProducer(conn)
    if err != nil {
        log.Fatal(err)
    }
    defer producer.Close()

    for i := 1; i <= 20; i++ {
        task := ImageTask{
            ImageID: fmt.Sprintf("img-%04d", i),
            URL:     fmt.Sprintf("https://example.com/images/%d.jpg", i),
            Width:   1920,
            Height:  1080,
        }
        body, _ := json.Marshal(task)
        ctx, cancel := context.WithTimeout(context.Background(), 5*time.Second)
        if err := producer.Publish(ctx, "", queueName, body); err != nil {
            log.Printf("publish failed: %v", err)
        } else {
            log.Printf("published task: %s", task.ImageID)
        }
        cancel()
    }
}

func runConsumer(conn *mq.Connection) {
    consumer, err := mq.NewConsumer(conn, queueName, 5) // 5个并发 worker
    if err != nil {
        log.Fatal(err)
    }
    defer consumer.Close()

    ctx, stop := signal.NotifyContext(context.Background(), os.Interrupt, syscall.SIGTERM)
    defer stop()

    log.Println("worker started, waiting for tasks...")
    consumer.Start(ctx, func(ctx context.Context, msg amqp.Delivery) error {
        var task ImageTask
        if err := json.Unmarshal(msg.Body, &task); err != nil {
            return err
        }

        // 模拟图片压缩处理
        log.Printf("processing image: %s", task.ImageID)
        time.Sleep(500 * time.Millisecond) // 模拟耗时操作
        log.Printf("completed image: %s", task.ImageID)
        return nil
    })
}

8.6 发布订阅案例(事件广播)

场景:用户注册后,同时发送欢迎邮件、发放优惠券、记录日志。

package main

import (
    "context"
    "encoding/json"
    "log"
    "time"

    amqp "github.com/rabbitmq/amqp091-go"
)

const exchangeName = "user.events"

type UserRegisteredEvent struct {
    UserID    int64     `json:"user_id"`
    Email     string    `json:"email"`
    CreatedAt time.Time `json:"created_at"`
}

// 发布用户注册事件
func publishUserRegistered(ch *amqp.Channel, event UserRegisteredEvent) error {
    ch.ExchangeDeclare(exchangeName, "fanout", true, false, false, false, nil)

    body, _ := json.Marshal(event)
    return ch.Publish(
        exchangeName, "",
        false, false,
        amqp.Publishing{
            ContentType:  "application/json",
            DeliveryMode: amqp.Persistent,
            Body:         body,
        },
    )
}

// 邮件服务订阅
func emailSubscriber(ch *amqp.Channel) {
    ch.ExchangeDeclare(exchangeName, "fanout", true, false, false, false, nil)
    q, _ := ch.QueueDeclare("", false, true, true, false, nil) // 临时队列
    ch.QueueBind(q.Name, "", exchangeName, false, nil)

    msgs, _ := ch.Consume(q.Name, "", false, false, false, false, nil)
    for msg := range msgs {
        var event UserRegisteredEvent
        json.Unmarshal(msg.Body, &event)
        log.Printf("[邮件服务] 发送欢迎邮件到: %s", event.Email)
        msg.Ack(false)
    }
}

// 优惠券服务订阅
func couponSubscriber(ch *amqp.Channel) {
    ch.ExchangeDeclare(exchangeName, "fanout", true, false, false, false, nil)
    q, _ := ch.QueueDeclare("coupon.user.register", true, false, false, false, nil) // 持久队列
    ch.QueueBind(q.Name, "", exchangeName, false, nil)

    msgs, _ := ch.Consume(q.Name, "", false, false, false, false, nil)
    for msg := range msgs {
        var event UserRegisteredEvent
        json.Unmarshal(msg.Body, &event)
        log.Printf("[优惠券服务] 为用户 %d 发放新手优惠券", event.UserID)
        msg.Ack(false)
    }
}

func main() {
    // 实际使用中每个服务独立部署,这里仅做演示
    conn, _ := amqp.Dial("amqp://guest:guest@localhost:5672/")
    defer conn.Close()

    ch, _ := conn.Channel()
    defer ch.Close()

    ch2, _ := conn.Channel()
    defer ch2.Close()

    go emailSubscriber(ch)
    go couponSubscriber(ch2)

    // 触发注册事件
    publishUserRegistered(ch, UserRegisteredEvent{
        UserID:    1001,
        Email:     "user@example.com",
        CreatedAt: time.Now(),
    })

    select {}
}

8.7 RPC 模式案例

场景:客户端发送计算请求,等待服务端返回结果(请求-响应模式)。

package main

import (
    "context"
    "encoding/json"
    "fmt"
    "log"
    "math/rand"
    "strconv"
    "time"

    amqp "github.com/rabbitmq/amqp091-go"
)

type RPCRequest struct {
    A, B int `json:"a,b"`
}

type RPCResponse struct {
    Result int    `json:"result"`
    Error  string `json:"error,omitempty"`
}

// RPC 服务端
func rpcServer(ch *amqp.Channel) {
    ch.QueueDeclare("rpc.add", false, false, false, false, nil)
    ch.Qos(1, 0, false)

    msgs, _ := ch.Consume("rpc.add", "", false, false, false, false, nil)
    log.Println("[RPC Server] waiting for requests...")

    for msg := range msgs {
        var req RPCRequest
        json.Unmarshal(msg.Body, &req)

        resp := RPCResponse{Result: req.A + req.B}
        body, _ := json.Marshal(resp)

        // 回复到 reply_to 队列
        ch.Publish(
            "", msg.ReplyTo,
            false, false,
            amqp.Publishing{
                ContentType:   "application/json",
                CorrelationId: msg.CorrelationId, // 回传相同的 correlation_id
                Body:          body,
            },
        )
        msg.Ack(false)
        log.Printf("[RPC Server] %d + %d = %d", req.A, req.B, resp.Result)
    }
}

// RPC 客户端
func rpcCall(ch *amqp.Channel, a, b int) (int, error) {
    // 创建临时回调队列
    q, _ := ch.QueueDeclare("", false, false, true, false, nil)

    replies, _ := ch.Consume(q.Name, "", true, false, false, false, nil)

    corrID := strconv.Itoa(rand.Int())
    req := RPCRequest{A: a, B: b}
    body, _ := json.Marshal(req)

    ch.Publish(
        "", "rpc.add",
        false, false,
        amqp.Publishing{
            ContentType:   "application/json",
            CorrelationId: corrID,
            ReplyTo:       q.Name, // 告诉服务端回复到哪里
            Body:          body,
        },
    )

    ctx, cancel := context.WithTimeout(context.Background(), 5*time.Second)
    defer cancel()

    for {
        select {
        case msg := <-replies:
            if msg.CorrelationId == corrID {
                var resp RPCResponse
                json.Unmarshal(msg.Body, &resp)
                if resp.Error != "" {
                    return 0, fmt.Errorf(resp.Error)
                }
                return resp.Result, nil
            }
        case <-ctx.Done():
            return 0, fmt.Errorf("RPC timeout")
        }
    }
}

func main() {
    conn, _ := amqp.Dial("amqp://guest:guest@localhost:5672/")
    defer conn.Close()

    serverCh, _ := conn.Channel()
    clientCh, _ := conn.Channel()

    go rpcServer(serverCh)
    time.Sleep(100 * time.Millisecond) // 等待服务端就绪

    result, err := rpcCall(clientCh, 30, 12)
    if err != nil {
        log.Fatal(err)
    }
    log.Printf("[RPC Client] 30 + 12 = %d", result) // 输出: 42
}

8.8 订单超时关闭案例(延迟队列)

场景:订单 30 分钟未支付自动关闭。

package main

import (
    "encoding/json"
    "log"
    "time"

    amqp "github.com/rabbitmq/amqp091-go"
)

type Order struct {
    OrderID   string    `json:"order_id"`
    UserID    int64     `json:"user_id"`
    Amount    float64   `json:"amount"`
    CreatedAt time.Time `json:"created_at"`
}

func setupQueues(ch *amqp.Channel) {
    // 死信交换器
    ch.ExchangeDeclare("order.dlx", "direct", true, false, false, false, nil)
    ch.QueueDeclare("order.timeout.queue", true, false, false, false, nil)
    ch.QueueBind("order.timeout.queue", "order.timeout", "order.dlx", false, nil)

    // 业务队列(30分钟TTL,过期进死信)
    ch.QueueDeclare("order.pending.queue", true, false, false, false, amqp.Table{
        "x-dead-letter-exchange":    "order.dlx",
        "x-dead-letter-routing-key": "order.timeout",
        "x-message-ttl":             int32(30 * 60 * 1000), // 30分钟
    })
}

// 创建订单时发送消息
func publishOrder(ch *amqp.Channel, order Order) error {
    body, _ := json.Marshal(order)
    return ch.Publish(
        "", "order.pending.queue",
        false, false,
        amqp.Publishing{
            DeliveryMode: amqp.Persistent,
            ContentType:  "application/json",
            Body:         body,
        },
    )
}

// 处理超时订单
func handleTimeout(ch *amqp.Channel) {
    msgs, _ := ch.Consume("order.timeout.queue", "", false, false, false, false, nil)
    for msg := range msgs {
        var order Order
        json.Unmarshal(msg.Body, &order)

        // 检查订单是否已支付
        if isOrderPaid(order.OrderID) {
            log.Printf("订单 %s 已支付,忽略超时消息", order.OrderID)
        } else {
            closeOrder(order.OrderID)
            log.Printf("订单 %s 超时关闭", order.OrderID)
        }
        msg.Ack(false)
    }
}

func isOrderPaid(orderID string) bool { /* 查询数据库 */ return false }
func closeOrder(orderID string)       { /* 更新订单状态 */ }

func main() {
    conn, _ := amqp.Dial("amqp://guest:guest@localhost:5672/")
    defer conn.Close()
    ch, _ := conn.Channel()
    defer ch.Close()

    setupQueues(ch)

    // 创建订单
    publishOrder(ch, Order{
        OrderID:   "ORDER-2026-001",
        UserID:    1001,
        Amount:    99.99,
        CreatedAt: time.Now(),
    })
    log.Println("订单已创建,30分钟后若未支付将自动关闭")

    // 启动超时处理器
    handleTimeout(ch)
}

9. 集群与高可用

9.1 集群模式

RabbitMQ 集群将多个节点组成一个逻辑 Broker,Exchange、Binding、用户等元数据在所有节点共享,但队列默认只在创建节点上存储。

节点1 (disk)  ←──Erlang OTP──→  节点2 (disk)  ←──Erlang OTP──→  节点3 (ram)
   │                                 │                                 │
  Queue-A                         Queue-B                           (内存节点,性能高)

节点类型:

  • 磁盘节点(disk):元数据持久化到磁盘,集群至少需要一个
  • 内存节点(ram):元数据只在内存,重启丢失,但性能更好

9.2 镜像队列(Classic Mirrored Queue)

RabbitMQ 3.x 的高可用方案,队列数据在多个节点同步副本:

# 设置镜像策略:匹配 "ha." 前缀的队列,所有节点都同步
rabbitmqctl set_policy ha-all "^ha\." \
    '{"ha-mode":"all","ha-sync-mode":"automatic"}'

# 同步到指定数量节点
rabbitmqctl set_policy ha-two "^two\." \
    '{"ha-mode":"exactly","ha-params":2,"ha-sync-mode":"automatic"}'

9.3 仲裁队列(Quorum Queue)—— 推荐

RabbitMQ 3.8+ 推出的新一代高可用队列,基于 Raft 共识算法,比镜像队列更可靠:

// 声明仲裁队列
ch.QueueDeclare("my.quorum.queue", true, false, false, false, amqp.Table{
    "x-queue-type": "quorum",
})
特性 镜像队列 仲裁队列
一致性协议 自研(较弱) Raft(强一致)
脑裂处理 容易丢数据 安全
性能 一般 略低但更安全
推荐程度 即将废弃 推荐使用

9.4 Docker Compose 集群部署

# docker-compose.yml
version: '3.8'
services:
  rabbitmq1:
    image: rabbitmq:3.13-management
    hostname: rabbitmq1
    environment:
      RABBITMQ_ERLANG_COOKIE: "secret-cookie"
      RABBITMQ_DEFAULT_USER: admin
      RABBITMQ_DEFAULT_PASS: admin123
    ports:
      - "5672:5672"
      - "15672:15672"
    volumes:
      - rabbitmq1_data:/var/lib/rabbitmq

  rabbitmq2:
    image: rabbitmq:3.13-management
    hostname: rabbitmq2
    environment:
      RABBITMQ_ERLANG_COOKIE: "secret-cookie"
      RABBITMQ_DEFAULT_USER: admin
      RABBITMQ_DEFAULT_PASS: admin123
    ports:
      - "5673:5672"
      - "15673:15672"
    volumes:
      - rabbitmq2_data:/var/lib/rabbitmq
    depends_on:
      - rabbitmq1

  rabbitmq3:
    image: rabbitmq:3.13-management
    hostname: rabbitmq3
    environment:
      RABBITMQ_ERLANG_COOKIE: "secret-cookie"
      RABBITMQ_DEFAULT_USER: admin
      RABBITMQ_DEFAULT_PASS: admin123
    ports:
      - "5674:5672"
      - "15674:15672"
    volumes:
      - rabbitmq3_data:/var/lib/rabbitmq
    depends_on:
      - rabbitmq1

volumes:
  rabbitmq1_data:
  rabbitmq2_data:
  rabbitmq3_data:
docker-compose up -d

# 节点2加入集群
docker exec rabbitmq2 rabbitmqctl stop_app
docker exec rabbitmq2 rabbitmqctl reset
docker exec rabbitmq2 rabbitmqctl join_cluster rabbit@rabbitmq1
docker exec rabbitmq2 rabbitmqctl start_app

# 节点3加入集群
docker exec rabbitmq3 rabbitmqctl stop_app
docker exec rabbitmq3 rabbitmqctl reset
docker exec rabbitmq3 rabbitmqctl join_cluster rabbit@rabbitmq1
docker exec rabbitmq3 rabbitmqctl start_app

# 查看集群状态
docker exec rabbitmq1 rabbitmqctl cluster_status

9.5 Federation(跨数据中心)

Federation 插件允许不同地区的 RabbitMQ 集群之间同步消息,适合跨机房部署:

rabbitmq-plugins enable rabbitmq_federation
rabbitmq-plugins enable rabbitmq_federation_management

# 配置 upstream(上游集群)
rabbitmqctl set_parameter federation-upstream us-west \
    '{"uri":"amqp://admin:pass@us-west.example.com:5672"}'

# 配置 federation policy
rabbitmqctl set_policy federate-me "^federated\." \
    '{"federation-upstream-set":"all"}' --apply-to exchanges

10. 运维与监控

10.1 常用管理命令

# 查看队列状态
rabbitmqctl list_queues name messages consumers memory

# 查看 Exchange
rabbitmqctl list_exchanges

# 查看 Binding
rabbitmqctl list_bindings

# 查看连接
rabbitmqctl list_connections

# 查看信道
rabbitmqctl list_channels

# 清空队列
rabbitmqctl purge_queue my-queue

# 删除队列
rabbitmqctl delete_queue my-queue

# 查看集群状态
rabbitmqctl cluster_status

# 查看节点内存使用
rabbitmqctl status | grep memory

# 发送测试消息(用于调试)
rabbitmqadmin publish exchange=my-exchange routing_key=test payload="hello"

10.2 Management HTTP API

# 列出所有队列(JSON)
curl -u admin:admin123 http://localhost:15672/api/queues

# 查看某个队列详情
curl -u admin:admin123 http://localhost:15672/api/queues/%2F/my-queue

# 获取队列中的消息(不消费)
curl -u admin:admin123 -X POST http://localhost:15672/api/queues/%2F/my-queue/get \
    -H "Content-Type: application/json" \
    -d '{"count":5,"ackmode":"ack_requeue_true","encoding":"auto"}'

# 发布消息
curl -u admin:admin123 -X POST http://localhost:15672/api/exchanges/%2F/my-exchange/publish \
    -H "Content-Type: application/json" \
    -d '{"routing_key":"test","payload":"hello","payload_encoding":"string","properties":{}}'

10.3 Prometheus + Grafana 监控

# 启用 Prometheus 插件
rabbitmq-plugins enable rabbitmq_prometheus

# 指标暴露地址
# http://localhost:15692/metrics

关键监控指标:

指标 说明 告警阈值
rabbitmq_queue_messages 队列中消息数 > 10000
rabbitmq_queue_messages_unacked 未确认消息数 > 1000
rabbitmq_queue_consumers 消费者数量 = 0 时告警
rabbitmq_connections 连接数 > 5000
rabbitmq_channels 信道数 > 10000
rabbitmq_node_mem_used 内存使用 > 内存报警阈值 80%
rabbitmq_node_disk_free 磁盘剩余 < 磁盘报警阈值

10.4 内存与磁盘报警

# rabbitmq.conf
# 内存高水位:超过40%时停止接收新消息
vm_memory_high_watermark.relative = 0.4

# 磁盘空间低水位:低于50MB时报警
disk_free_limit.absolute = 50MB

# 流量控制:消费者来不及消费时自动降速

11. 高频面试题

Q1: RabbitMQ 如何保证消息不丢失?

三端保障:

① 生产者端 → Publisher Confirm

  • 将 Channel 设为 Confirm 模式
  • 每条消息发送后等待 Broker 的 ACK/NACK 回复
  • 若收到 NACK 或超时,重新发送并记录日志

② Broker 端 → 持久化

  • Exchange 声明时 durable=true
  • Queue 声明时 durable=true
  • 消息发送时 delivery_mode=2(Persistent)
  • 三者缺一不可

③ 消费者端 → 手动 ACK

  • autoAck=false,业务逻辑处理成功后再调用 Ack
  • 处理失败调用 Nack(false, true) 重新入队
  • 结合幂等性设计防止重复消费

Q2: 如何防止消息重复消费?

消息重复消费的根本原因:消费者处理成功但 ACK 未送达(网络抖动、消费者崩溃),Broker 重投消息。

解决方案:幂等性设计

func handleMessage(ctx context.Context, msg amqp.Delivery) error {
    // 提取业务唯一 ID(如订单ID、支付流水号)
    var payload struct {
        OrderID string `json:"order_id"`
    }
    json.Unmarshal(msg.Body, &payload)

    // 方案1:Redis 去重(set NX,过期时间 = 消息保留时间)
    ok, _ := redis.SetNX(ctx, "processed:"+payload.OrderID, 1, 24*time.Hour)
    if !ok {
        // 已处理过,直接 ACK 丢弃
        msg.Ack(false)
        return nil
    }

    // 方案2:数据库唯一索引(insert ignore 或 ON DUPLICATE KEY IGNORE)
    // INSERT INTO order_logs (order_id, ...) VALUES (?) ON DUPLICATE KEY UPDATE id=id

    // 执行业务逻辑
    return processOrder(ctx, payload.OrderID)
}

Q3: 如何保证消息顺序?

RabbitMQ 本身单队列内消息是有序的(FIFO),但以下情况会破坏顺序:

  • 多个 Consumer 并发消费同一队列
  • Consumer 处理失败 Nack 重入队后,顺序被打乱

解决方案:

方案1:单消费者 + 串行处理(简单,吞吐低)
方案2:同一业务 ID 的消息路由到同一队列(如 order_id % N 选队列)
方案3:消费者业务层面排序(消息带序号,乱序时暂存等待)
// 按 order_id 哈希选队列,保证同一订单的消息进同一队列
func routingKey(orderID string) string {
    h := fnv.New32a()
    h.Write([]byte(orderID))
    return fmt.Sprintf("order.queue.%d", h.Sum32()%10) // 10个队列
}

Q4: RabbitMQ 消息积压怎么处理?

紧急处置:

  1. 扩容 Consumer:快速增加消费者数量,提高消费速度
  2. 跳过非关键消息:临时设置 Consumer 快速 ACK(丢弃)非重要消息

根本分析:

消息积压原因:
├── 消费者处理慢(DB 慢查询、外部 API 超时)
│   └── 优化业务逻辑、增加 prefetch count、横向扩容 Consumer
├── 消费者崩溃/下线
│   └── 监控告警、自动重启、健康检查
├── 生产者突发流量
│   └── 生产端限速、Broker 内存/磁盘容量规划
└── 队列配置不当(单队列、无并发)
    └── 拆分队列、增加并发度

临时扩容方案(Go):

// 动态增加消费者
func scaleConsumers(conn *mq.Connection, queue string, n int) []*mq.Consumer {
    consumers := make([]*mq.Consumer, n)
    for i := 0; i < n; i++ {
        c, _ := mq.NewConsumer(conn, queue, 10)
        consumers[i] = c
        go c.Start(context.Background(), processMessage)
    }
    return consumers
}

Q5: 死信队列的使用场景?

场景 说明
消费失败重试N次后入死信 避免无限重试;死信队列人工排查
订单超时关闭 TTL + DLX 实现延迟处理
消息过滤/降级 不符合条件的消息转移到死信,主流程不受影响
审计与问题排查 记录所有处理失败的消息,便于回溯

Q6: direct、fanout、topic 如何选择?

需求 选择
一个消息只给一个消费者,精确匹配 Direct
一个消息广播给所有消费者 Fanout
按多个维度灵活过滤消息 Topic
基于消息属性路由(少用) Headers

Q7: RabbitMQ 集群节点崩溃了怎么办?

情况1:非队列所在节点崩溃

  • 其他节点仍可正常工作,不影响该节点上的队列

情况2:队列所在节点崩溃

  • 普通队列:队列不可用,消息丢失(若未持久化)
  • 镜像队列:自动切换到镜像节点,不中断服务
  • 仲裁队列:Raft 自动选主,只要半数以上节点存活即可

情况3:唯一磁盘节点崩溃

  • 集群可以继续路由消息
  • 但无法创建队列/交换器/绑定等管理操作
  • 恢复后自动同步元数据

Q8: 如何实现消息的延迟投递?

方案对比:

方案 原理 缺点
TTL + DLX 消息过期进死信队列 队头阻塞问题(队列TTL可用,消息TTL不行)
延迟插件 x-delayed-message Exchange 需安装插件;延迟消息存内存,重启可能丢失
业务层轮询 DB 存储延迟任务,定时扫描 实现简单,精度受轮询间隔影响
RocketMQ 原生支持 18 个延迟级别 需换 MQ

生产推荐:延迟精度要求不高(秒级)用 TTL+DLX;精度要求高用延迟插件;消息量极大考虑 RocketMQ。

Q9: 信道(Channel)和连接(Connection)有什么区别?

维度 Connection Channel
本质 TCP 连接 TCP 上的虚拟逻辑通道
创建代价 高(三次握手、TLS 握手、AMQP 协商) 低(AMQP 帧内开启)
数量 建议每应用少量(1-5) 每个 goroutine 独立使用
线程安全 否 否
生命周期 随应用存活,需自动重连 随 Connection 或业务生命周期

最佳实践:

  • 一个应用维护 1~2 个 Connection(发/收可分开)
  • 每个 goroutine 用独立的 Channel
  • Channel 不可跨 goroutine 共享

Q10: AMQP 与其他协议的区别?

协议 特点 使用场景
AMQP 0-9-1 RabbitMQ 默认,功能完整,二进制 企业级消息队列
STOMP 文本协议,简单,WebSocket 友好 浏览器端、简单场景
MQTT 轻量,为物联网设计,低带宽 IoT、移动端
HTTP/HTTPS 通过 Management Plugin 提供 管理 API、简单发消息

12. 最佳实践总结

12.1 生产环境 Checklist

连接管理
  ✅ 实现自动重连逻辑(指数退避)
  ✅ 每个 goroutine 使用独立 Channel,不跨 goroutine 共享
  ✅ 程序退出时优雅关闭 Channel 和 Connection

消息可靠性
  ✅ 生产者开启 Publisher Confirm
  ✅ Exchange、Queue、Message 三层持久化
  ✅ 消费者使用手动 ACK(autoAck=false)
  ✅ 消费失败时 Nack + requeue,避免消息丢失
  ✅ 配置死信队列兜底

消费者设计
  ✅ 设置合理的 prefetch(QoS),避免内存溢出
  ✅ 消费逻辑实现幂等性(Redis/DB 去重)
  ✅ 避免在消费者内做长时间阻塞操作

队列设计
  ✅ 生产环境使用仲裁队列(Quorum Queue)替代镜像队列
  ✅ 为重要队列配置 DLX 死信交换器
  ✅ 设置队列 TTL 和 max-length 防止无限堆积

监控告警
  ✅ 监控队列深度,消息积压及时告警
  ✅ 监控 unacked 消息数量
  ✅ 监控内存和磁盘使用率
  ✅ 消费者存活监控

12.2 常见陷阱

❌ autoAck=true:消息投递后立即删除,处理失败消息丢失
✅ 改为手动 ACK

❌ 在循环中为每条消息创建 Channel/Connection
✅ 复用 Connection,每个 goroutine 持有独立 Channel

❌ 未设置 QoS,消费者拉取大量消息堆积在内存
✅ ch.Qos(10, 0, false) 限制预取数量

❌ 消费失败无限 Nack+requeue,形成消息风暴
✅ 记录重试次数,超过阈值后进死信队列或告警

❌ 依赖消息唯一性保证业务逻辑,未实现幂等
✅ 消费者必须实现幂等性

❌ fanout/topic 交换器期望有消息历史
✅ fanout/topic 无历史消息,中途加入的消费者收不到之前的消息

小结

知识点 核心要点
Exchange 类型 Direct=精确匹配,Fanout=广播,Topic=通配符,Headers=头匹配
消息可靠性 Publisher Confirm + 三层持久化 + 手动 ACK
死信队列 拒绝/TTL/超长 → DLX → 死信队列,用于失败兜底和延迟队列
延迟队列 TTL+DLX(简单)或延迟消息插件(精确)
集群高可用 3.8+ 推荐仲裁队列(Quorum Queue),基于 Raft
消息积压 扩容 Consumer + 优化处理逻辑 + 合理 QoS 设置
幂等消费 Redis SetNX 或 DB 唯一索引去重
信道 vs 连接 Connection=TCP 长连接,Channel=虚拟通道,每 goroutine 独立 Channel