目录

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 的所有消息都会走确认流程。

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)

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. 消息被拒绝NackRejectrequeue=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