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 消息成为死信的三种情况
- 消息被拒绝:
Nack或Reject且requeue=false - 消息 TTL 过期:消息在队列中超过设定的存活时间
- 队列长度溢出:队列达到
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 消息积压怎么处理?
紧急处置:
- 扩容 Consumer:快速增加消费者数量,提高消费速度
- 跳过非关键消息:临时设置 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 |
xingliuhua