RabbitMQ 深度解析 —— 灵活路由与消息可靠性 / RabbitMQ Deep Dive into Flexible Routing and Reliability
📅 创建时间:2026-07-28 🏷️ 标签:#RabbitMQ #AMQP #Exchange #死信队列 #延迟队列 #QuorumQueues #消息确认 📚 前置知识:[[/03-web/07-middleware/02-mq-kafka]](Kafka 核心机制) [[00-overview]](消息队列全景) 📚 相关知识:[[/03-web/07-middleware/03-mq-others]](MQ 对比选型) [[/03-web/07-middleware/07-mq-advanced]](MQ 高级特性)
📋 本章目标
- 理解 RabbitMQ 的核心架构:Connection、Channel、Exchange、Queue、Binding、Virtual Host
- 掌握四种 Exchange 类型的路由规则:Direct、Topic、Fanout、Headers
- 理解消息确认的完整链路:Publisher Confirm + Consumer Ack/Nack/Reject
- 掌握死信队列的三种触发条件:TTL 超时、队列满、消息被拒绝
- 理解延迟队列的两种实现方式及其优劣
- 了解 RabbitMQ 高可用方案:Mirrored Queues vs Quorum Queues
- 建立 RabbitMQ vs Kafka 的选型决策框架
场景:订单取消的 30 分钟倒计时
┌─────────────────────────────────────────────────────────────┐
│ │
│ 用户下单后,需要在 30 分钟内完成支付。 │
│ 如果超时未支付,系统自动取消订单,释放库存。 │
│ │
│ 直接方案 1:定时任务每分钟扫一次数据库 │
│ → "SELECT * FROM orders WHERE status='pending' │
│ AND created_at < NOW() - INTERVAL 30 MINUTE" │
│ → 订单量 100 万时,这个查询耗时 5 秒,数据库 CPU 飙升 │
│ │
│ 直接方案 2:每个订单起一个 setTimeout │
│ → 100 万个订单 = 100 万个定时器 │
│ → 服务重启 → 所有定时器丢失 │
│ │
│ RabbitMQ 方案: │
│ → 下单时发一条"延迟 30 分钟"的消息 │
│ → 30 分钟后消息到达消费者 │
│ → 消费者检查订单状态,未支付则取消 │
│ → 不需要轮询,不依赖服务内存 │
│ │
└─────────────────────────────────────────────────────────────┘1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
第1部分:RabbitMQ 核心架构
1.1 AMQP 0-9-1 模型
┌─────────────────────────────────────────────────────────────┐
│ RabbitMQ 架构全景 │
├─────────────────────────────────────────────────────────────┤
│ │
│ ┌─────────────────────────────────────────────────────┐ │
│ │ RabbitMQ Broker │ │
│ │ │ │
│ │ ┌──────────────────────────────────────────┐ │ │
│ │ │ Virtual Host (vhost) │ │ │
│ │ │ namespace: "/ecommerce" │ │ │
│ │ │ │ │ │
│ │ │ Producer │ │ │
│ │ │ │ │ │ │
│ │ │ ↓ Connection │ │ │
│ │ │ ┌────────┐ │ │ │
│ │ │ │Channel 1│──→ Exchange ──→ Queue A ──→│ │ │
│ │ │ │Channel 2│──→ Exchange ──→ Queue B ──→│ │ │
│ │ │ └────────┘ (type:topic) Queue C │ │ │
│ │ │ │ │ │ │
│ │ │ ┌─────┼─────┐ │ │ │
│ │ │ ↓ ↓ ↓ │ │ │
│ │ │ Binding Binding Binding │ │ │
│ │ │ ("#.err") ("*.order") ("#") │ │ │
│ │ │ │ │ │
│ │ │ Consumer ←── Queue A (via Channel) │ │ │
│ │ │ Consumer ←── Queue B (via Channel) │ │ │
│ │ └──────────────────────────────────────────┘ │ │
│ │ │ │
│ │ ┌──────────────────────────────────────────┐ │ │
│ │ │ Virtual Host (vhost) │ │ │
│ │ │ namespace: "/monitoring" │ │ │
│ │ │ (完全独立,资源隔离) │ │ │
│ │ └──────────────────────────────────────────┘ │ │
│ └─────────────────────────────────────────────────────┘ │
│ │
└─────────────────────────────────────────────────────────────┘1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
1.2 核心概念详解
┌─────────────────────────────────────────────────────────────┐
│ RabbitMQ 核心概念详解 │
├─────────────────────────────────────────────────────────────┤
│ │
│ Connection(连接) │
│ ┌─────────────────────────────────────────────────────┐ │
│ │ • 客户端和 Broker 之间的 TCP 长连接 │ │
│ │ • 一个 Connection 可以有多个 Channel │ │
│ │ • Connection 是重量级的(TLS 握手、认证) │ │
│ │ • 建议:一个进程一个 Connection │ │
│ └─────────────────────────────────────────────────────┘ │
│ │
│ Channel(信道) │
│ ┌─────────────────────────────────────────────────────┐ │
│ │ • 复用 Connection 的轻量级"虚拟连接" │ │
│ │ • 每个 Channel 有自己的 ID,互相隔离 │ │
│ │ • Channel 是轻量级的(创建/销毁成本低) │ │
│ │ • 建议:一个线程一个 Channel(线程不安全) │ │
│ │ • 为什么要 Channel? │ │
│ │ → 不用为每个操作建立新的 TCP 连接 │ │
│ │ → 一个 TCP 连接上多路复用若干个 Channel │ │
│ └─────────────────────────────────────────────────────┘ │
│ │
│ Exchange(交换机) │
│ ┌─────────────────────────────────────────────────────┐ │
│ │ • Producer 不直接发消息到 Queue,而是发到 Exchange │ │
│ │ • Exchange 根据路由规则将消息分发到 Queue │ │
│ │ • 四种类型:Direct / Topic / Fanout / Headers │ │
│ │ • Exchange 不存储消息,只做路由 │ │
│ └─────────────────────────────────────────────────────┘ │
│ │
│ Queue(队列) │
│ ┌─────────────────────────────────────────────────────┐ │
│ │ • 消息的实际存储单元 │ │
│ │ • Consumer 从 Queue 消费消息 │ │
│ │ • 可配置:持久化、排他性、自动删除、TTL、最大长度 │ │
│ │ • 经典队列 vs 仲裁队列(Quorum Queue) │ │
│ └─────────────────────────────────────────────────────┘ │
│ │
│ Binding(绑定) │
│ ┌─────────────────────────────────────────────────────┐ │
│ │ • Exchange 和 Queue 之间的"连线" │ │
│ │ • 每条 Binding 有一个 Routing Key │ │
│ │ • Exchange 根据消息的 Routing Key 匹配 Binding │ │
│ │ • 一个 Queue 可以绑定到多个 Exchange │ │
│ │ • 一个 Exchange 可以绑定到多个 Queue │ │
│ └─────────────────────────────────────────────────────┘ │
│ │
│ Virtual Host(虚拟主机) │
│ ┌─────────────────────────────────────────────────────┐ │
│ │ • RabbitMQ 内部的逻辑隔离单位 │ │
│ │ • 每个 vhost 有独立的 Exchange、Queue、权限 │ │
│ │ • 典型用法: │ │
│ │ /ecommerce → 电商业务 │ │
│ │ /monitoring → 监控数据 │ │
│ │ / → 默认 vhost(不推荐生产用) │ │
│ └─────────────────────────────────────────────────────┘ │
│ │
└─────────────────────────────────────────────────────────────┘1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
第2部分:四种 Exchange 类型
2.1 Direct Exchange
┌─────────────────────────────────────────────────────────────┐
│ Direct Exchange —— 精确匹配 │
├─────────────────────────────────────────────────────────────┤
│ │
│ Routing Key 完全匹配才路由 │
│ │
│ ┌─────────────────────────────────────────────────────┐ │
│ │ │ │
│ │ Direct Exchange │ │
│ │ "order.ex" │ │
│ │ │ │
│ │ msg(routing_key="create") │ │
│ │ │ │ │
│ │ ├──→ Binding("create") → Queue A ✅ │ │
│ │ ├──→ Binding("delete") → Queue B ❌ │ │
│ │ └──→ Binding("update") → Queue C ❌ │ │
│ │ │ │
│ │ msg(routing_key="delete") │ │
│ │ │ │ │
│ │ ├──→ Binding("create") → Queue A ❌ │ │
│ │ ├──→ Binding("delete") → Queue B ✅ │ │
│ │ └──→ Binding("update") → Queue C ❌ │ │
│ │ │ │
│ └─────────────────────────────────────────────────────┘ │
│ │
│ 适用场景: │
│ • 单播:一个命令发给确定的处理者 │
│ • RPC:请求-响应模式 │
│ • 任务分配:不同的任务类型走不同的队列 │
│ │
└─────────────────────────────────────────────────────────────┘1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
2.2 Topic Exchange
┌─────────────────────────────────────────────────────────────┐
│ Topic Exchange —— 模式匹配 │
├─────────────────────────────────────────────────────────────┤
│ │
│ Routing Key 按 "." 分段,支持通配符: │
│ * (star) → 匹配恰好一段 │
│ # (hash) → 匹配零段或多段 │
│ │
│ ┌─────────────────────────────────────────────────────┐ │
│ │ │ │
│ │ Topic Exchange │ │
│ │ "events.ex" │ │
│ │ │ │
│ │ 消息: "order.us.electronics" │ │
│ │ │ │ │
│ │ ├── "order.*.electronics" → ✅ │ │
│ │ ├── "order.#" → ✅ │ │
│ │ ├── "*.us.*" → ✅ │ │
│ │ ├── "payment.#" → ❌ │ │
│ │ └── "order.eu.*" → ❌ │ │
│ │ │ │
│ │ 消息: "payment.cn.alipay" │ │
│ │ │ │ │
│ │ ├── "payment.*.alipay" → ✅ │ │
│ │ ├── "payment.#" → ✅ │ │
│ │ └── "order.#" → ❌ │ │
│ │ │ │
│ └─────────────────────────────────────────────────────┘ │
│ │
│ 适用场景: │
│ • 多级分类路由(如日志级别 + 服务名 + 环境) │
│ • 按地域/业务线分流消息 │
│ • 灵活的多对多路由 │
│ │
└─────────────────────────────────────────────────────────────┘1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
2.3 Fanout Exchange
┌─────────────────────────────────────────────────────────────┐
│ Fanout Exchange —— 广播 │
├─────────────────────────────────────────────────────────────┤
│ │
│ 忽略 Routing Key,将消息广播到所有绑定的 Queue │
│ │
│ ┌─────────────────────────────────────────────────────┐ │
│ │ │ │
│ │ Fanout Exchange │ │
│ │ "broadcast.ex" │ │
│ │ │ │
│ │ msg ──┬──→ Queue A (订单服务) │ │
│ │ ├──→ Queue B (风控服务) │ │
│ │ ├──→ Queue C (推送服务) │ │
│ │ └──→ Queue D (分析服务) │ │
│ │ │ │
│ │ 每个 Queue 都收到完全相同的消息副本 │ │
│ │ │ │
│ └─────────────────────────────────────────────────────┘ │
│ │
│ 适用场景: │
│ • 事件广播(如"用户注册成功"通知多个系统) │
│ • 配置变更通知 │
│ • 缓存失效广播 │
│ │
└─────────────────────────────────────────────────────────────┘1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
2.4 Headers Exchange
┌─────────────────────────────────────────────────────────────┐
│ Headers Exchange —— 消息头匹配 │
├─────────────────────────────────────────────────────────────┤
│ │
│ 不看 Routing Key,而是根据消息的 Headers 属性匹配 │
│ │
│ ┌─────────────────────────────────────────────────────┐ │
│ │ │ │
│ │ 消息 Headers: │ │
│ │ { │ │
│ │ "format": "pdf", │ │
│ │ "type": "report", │ │
│ │ "x-match": "all" ← 所有条件都要满足 │ │
│ │ } │ │
│ │ │ │
│ │ Binding 1: format=pdf, type=report → ✅ │ │
│ │ Binding 2: format=image → ❌ │ │
│ │ Binding 3: type=report, x-match=any → ✅ │ │
│ │ │ │
│ │ x-match 选项: │ │
│ │ "all" → 所有 Header 条件都满足才路由 │ │
│ │ "any" → 任意一个 Header 条件满足就路由 │ │
│ └─────────────────────────────────────────────────────┘ │
│ │
│ 适用场景: │
│ • 复杂的多条件路由(替代多层 Topic 匹配) │
│ • 注:Headers Exchange 性能低于 Topic Exchange │
│ • 大多数场景 Topic Exchange 就够了 │
│ │
└─────────────────────────────────────────────────────────────┘1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
2.5 四种 Exchange 对比总结
| 类型 | 路由规则 | Routing Key | 性能 | 典型场景 |
|---|---|---|---|---|
| Direct | 精确匹配 | 必须 | 最高 | 单播、任务分配 |
| Topic | 模式匹配 (*/#) | 必须 | 高 | 多级分类路由 |
| Fanout | 广播 | 忽略 | 高 | 事件广播 |
| Headers | Header 匹配 | 忽略 | 较低 | 多条件复杂路由 |
typescript
// Node.js (amqplib) —— 声明 Exchange/Queue/Binding 的典型代码
import * as amqp from "amqplib";
async function setupTopology(): Promise<void> {
const conn = await amqp.connect("amqp://localhost");
const channel = await conn.createChannel();
// 1. 声明 Exchange
await channel.assertExchange("order-events", "topic", {
durable: true, // Exchange 持久化(重启保留)
autoDelete: false, // 没有绑定的 Queue 时不自动删除
});
// 2. 声明 Queue
await channel.assertQueue("order-created-queue", {
durable: true, // Queue 持久化
exclusive: false, // 不排他(多个 Consumer 可连接)
autoDelete: false, // Consumer 断开时不删除
});
await channel.assertQueue("order-all-queue", {
durable: true,
});
// 3. 绑定 Queue 到 Exchange
// Binding: order-created-queue 只收 order.created 消息
await channel.bindQueue(
"order-created-queue",
"order-events",
"order.created" // Routing Key 精确匹配
);
// Binding: order-all-queue 收所有 order.* 消息
await channel.bindQueue(
"order-all-queue",
"order-events",
"order.#" // 通配符匹配所有 order.xxx
);
// 4. 发布消息
channel.publish(
"order-events",
"order.created", // Routing Key
Buffer.from(JSON.stringify({ orderId: "123", status: "created" })),
{
persistent: true, // 消息持久化
headers: {
"source": "order-service",
"version": "1.0",
},
}
);
}1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
第3部分:消息确认机制
3.1 Publisher Confirm(生产者确认)
┌─────────────────────────────────────────────────────────────┐
│ Publisher Confirm 流程 │
├─────────────────────────────────────────────────────────────┤
│ │
│ 默认情况下,Producer 发完消息就忘了(fire-and-forget) │
│ Publisher Confirm 让 Producer 知道消息是否到达 Broker │
│ │
│ ┌─────────────────────────────────────────────────────┐ │
│ │ │ │
│ │ Producer RabbitMQ Broker │ │
│ │ │ │ │ │
│ │ │── publish(msg1) ───────→│ │ │
│ │ │── publish(msg2) ───────→│ │ │
│ │ │── publish(msg3) ───────→│ │ │
│ │ │ │ (持久化到磁盘) │ │
│ │ │←── confirm(msg1) ──────│ ack │ │
│ │ │←── confirm(msg2) ──────│ ack │ │
│ │ │←── confirm(msg3) ──────│ ack │ │
│ │ │ │
│ │ 三种使用方式: │ │
│ │ ① 逐条确认(串行,吞吐低) │ │
│ │ ② 批量确认(等 N 条一起确认,延迟高) │ │
│ │ ③ 异步确认(推荐:发消息不等待,异步回调处理) │ │
│ └─────────────────────────────────────────────────────┘ │
│ │
└─────────────────────────────────────────────────────────────┘1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
typescript
// Node.js —— Publisher Confirm 异步模式
async function publishWithConfirm(
channel: amqp.ConfirmChannel,
exchange: string,
routingKey: string,
message: object
): Promise<void> {
return new Promise<void>((resolve, reject) => {
channel.publish(
exchange,
routingKey,
Buffer.from(JSON.stringify(message)),
{ persistent: true },
(err, ok) => {
// 这个回调就是 Publisher Confirm
if (err) {
console.error("❌ Publisher confirm failed:", err);
reject(err);
} else {
console.log("✅ Message confirmed by broker");
resolve();
}
}
);
});
}
// 批量异步确认模式(更高吞吐)
async function publishInBatch(
channel: amqp.ConfirmChannel,
messages: Array<{ exchange: string; routingKey: string; body: object }>
): Promise<void> {
const pendingConfirms = messages.map((msg) => {
return new Promise<void>((resolve, reject) => {
channel.publish(
msg.exchange,
msg.routingKey,
Buffer.from(JSON.stringify(msg.body)),
{ persistent: true },
(err) => (err ? reject(err) : resolve())
);
});
});
await Promise.all(pendingConfirms);
console.log(`✅ All ${messages.length} messages confirmed`);
}1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
3.2 Consumer Ack / Nack / Reject
┌─────────────────────────────────────────────────────────────┐
│ Consumer 确认机制 │
├─────────────────────────────────────────────────────────────┤
│ │
│ RabbitMQ 消费者有三种应答方式: │
│ │
│ ┌─────────────────────────────────────────────────────┐ │
│ │ │ │
│ │ ① basic.ack (肯定确认) │ │
│ │ 处理成功 → 告诉 Broker 可以删掉这条消息了 │ │
│ │ Broker 收到 ack → 从 Queue 中移除消息 │ │
│ │ │ │
│ │ ② basic.nack (否定确认) │ │
│ │ 处理失败 → 告诉 Broker 这条消息处理不了 │ │
│ │ requeue=true → 重新入队,给其他 Consumer │ │
│ │ requeue=false → 不重新入队(配合 DLX 进入死信)│ │
│ │ │ │
│ │ ③ basic.reject (拒绝) │ │
│ │ 类似 nack,但只能拒绝单条消息 │ │
│ │ nack 可以批量拒绝(multiple=true) │ │
│ └─────────────────────────────────────────────────────┘ │
│ │
│ 自动确认 vs 手动确认: │
│ ┌─────────────────────────────────────────────────────┐ │
│ │ │ │
│ │ 自动确认 (noAck: true): │ │
│ │ Broker 发出消息 → 立即从 Queue 删除 │ │
│ │ 风险:Consumer 崩溃 → 消息丢失 │ │
│ │ │ │
│ │ 手动确认 (noAck: false): │ │
│ │ Consumer 显式调用 ack/nack │ │
│ │ 推荐:生产环境始终使用手动确认 │ │
│ └─────────────────────────────────────────────────────┘ │
│ │
└─────────────────────────────────────────────────────────────┘1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
typescript
// Node.js —— 手动确认 + 重试逻辑
async function consumeWithManualAck(
channel: amqp.Channel,
queueName: string
): Promise<void> {
const MAX_RETRIES = 3;
channel.consume(queueName, async (msg) => {
if (!msg) return;
const content = JSON.parse(msg.content.toString());
let retries = (msg.properties.headers?.["x-retry-count"] ?? 0) as number;
try {
await processMessage(content);
// ✅ 成功 → 确认
channel.ack(msg);
console.log(`✅ Acked: ${content.id}`);
} catch (error) {
console.error(`❌ Failed: ${content.id}, retry ${retries + 1}/${MAX_RETRIES}`);
if (retries < MAX_RETRIES) {
// 重试:重新发布消息,携带重试计数
channel.publish(
msg.fields.exchange, // 发回原来的 Exchange
msg.fields.routingKey, // 原来的 Routing Key
msg.content,
{
persistent: true,
headers: {
...msg.properties.headers,
"x-retry-count": retries + 1,
},
}
);
// 确认原消息(因为已经重新发布了)
channel.ack(msg);
} else {
// 达到最大重试次数 → nack 且不 requeue → 进入死信队列
channel.nack(msg, false, false);
console.error(`💀 Sent to DLQ: ${content.id}`);
}
}
}, { noAck: false }); // 手动确认模式
}1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
第4部分:死信队列(DLX / Dead Letter Exchange)
4.1 什么是死信?
┌─────────────────────────────────────────────────────────────┐
│ 消息变成死信的三种情况 │
├─────────────────────────────────────────────────────────────┤
│ │
│ 条件1:消息被拒绝(reject/nack 且 requeue=false) │
│ ┌─────────────────────────────────────────────────────┐ │
│ │ │ │
│ │ Consumer 收到消息 → 处理失败 → nack(requeue=false) │ │
│ │ → 消息变成死信 → 转发到 DLX │ │
│ │ │ │
│ │ 适用:业务逻辑无法处理的消息(如数据格式错误) │ │
│ └─────────────────────────────────────────────────────┘ │
│ │
│ 条件2:消息 TTL 超时(Time To Live 到期) │
│ ┌─────────────────────────────────────────────────────┐ │
│ │ │ │
│ │ 消息在 Queue 中超过 TTL 时间未被消费 │ │
│ │ → 消息自动过期 → 转发到 DLX │ │
│ │ │ │
│ │ 适用:超时未处理的订单取消等延迟场景 │ │
│ └─────────────────────────────────────────────────────┘ │
│ │
│ 条件3:队列已满(达到最大长度) │
│ ┌─────────────────────────────────────────────────────┐ │
│ │ │ │
│ │ Queue 消息数量达到 max-length 上限 │ │
│ │ → 新消息到达时,队列头部的旧消息被"挤出" │ │
│ │ → 被挤出的消息转发到 DLX │ │
│ │ │ │
│ │ 适用:保护 Queue 不爆满,同时不丢失溢出的消息 │ │
│ └─────────────────────────────────────────────────────┘ │
│ │
└─────────────────────────────────────────────────────────────┘1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
4.2 死信队列配置
┌─────────────────────────────────────────────────────────────┐
│ 死信队列架构图 │
├─────────────────────────────────────────────────────────────┤
│ │
│ ┌─────────────────────────────────────────────────────┐ │
│ │ │ │
│ │ Producer → Exchange ──→ 主队列 (business-queue) │ │
│ │ │ │ │ │
│ │ │ │ 死信产生时 │ │
│ │ │ ↓ │ │
│ │ │ Dead Letter Exchange (DLX) │ │
│ │ │ │ │ │
│ │ │ ↓ │ │
│ │ └─────→ 死信队列 (dead-letter-queue) │ │
│ │ │ │ │
│ │ ↓ │ │
│ │ 监控/告警 Consumer │ │
│ │ (人工处理/自动重试) │ │
│ │ │ │
│ └─────────────────────────────────────────────────────┘ │
│ │
└─────────────────────────────────────────────────────────────┘1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
typescript
// Node.js —— 配置死信队列
async function setupDLQ(channel: amqp.Channel): Promise<void> {
// 1. 声明死信交换机
await channel.assertExchange("dlx-exchange", "direct", {
durable: true,
});
// 2. 声明死信队列
await channel.assertQueue("dead-letter-queue", {
durable: true,
});
// 3. 绑定死信队列到死信交换机
await channel.bindQueue("dead-letter-queue", "dlx-exchange", "dead");
// 4. 声明业务队列,配置 DLX
await channel.assertQueue("order-queue", {
durable: true,
arguments: {
// 指定死信交换机
"x-dead-letter-exchange": "dlx-exchange",
// 死信消息的 Routing Key
"x-dead-letter-routing-key": "dead",
// 消息 TTL(毫秒)—— 30 分钟后未消费则变成死信
"x-message-ttl": 30 * 60 * 1000,
// 队列最大长度
"x-max-length": 100000,
},
});
}
// 消费死信队列 —— 人工排查/告警
async function consumeDLQ(channel: amqp.Channel): Promise<void> {
channel.consume("dead-letter-queue", (msg) => {
if (!msg) return;
const deadInfo = {
content: JSON.parse(msg.content.toString()),
// 死信原因相关的 Header
reason: msg.properties.headers?.["x-death"]?.[0]?.reason, // "expired" | "rejected" | "maxlen"
originalQueue: msg.properties.headers?.["x-death"]?.[0]?.queue,
originalExchange: msg.properties.headers?.["x-death"]?.[0]?.exchange,
originalRoutingKey: msg.properties.headers?.["x-first-death-exchange"],
deathTime: msg.properties.headers?.["x-death"]?.[0]?.time,
};
console.error("💀 Dead Letter:", JSON.stringify(deadInfo, null, 2));
// 发出告警
alertingService.send({
level: "warning",
message: `Dead letter in order-queue: ${deadInfo.reason}`,
detail: deadInfo,
});
// 确认死信消息(否则会一直堆积在死信队列)
channel.ack(msg);
});
}1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
第5部分:延迟队列
5.1 两种实现方式
┌─────────────────────────────────────────────────────────────┐
│ 延迟队列的两种实现 │
├─────────────────────────────────────────────────────────────┤
│ │
│ 方案A:死信 + TTL(传统方式,不需要插件) │
│ ┌─────────────────────────────────────────────────────┐ │
│ │ │ │
│ │ Producer │ │
│ │ │ │ │
│ │ ↓ │ │
│ │ order.exchange ──→ delay-queue (TTL=30min) │ │
│ │ (没有 Consumer 消费它) │ │
│ │ │ │ │
│ │ 30 分钟后 TTL 到期 │ │
│ │ │ │ │
│ │ ↓ │ │
│ │ DLX Exchange │ │
│ │ │ │ │
│ │ ↓ │ │
│ │ target-queue │ │
│ │ │ │ │
│ │ ↓ │ │
│ │ 真正的 Consumer │ │
│ │ │ │
│ │ 问题:如果需要不同的延迟时间(如 5min, 10min, │ │
│ │ 30min),就得创建多个不同 TTL 的"延迟队列"。 │ │
│ └─────────────────────────────────────────────────────┘ │
│ │
│ 方案B:rabbitmq_delayed_message_exchange 插件 │
│ ┌─────────────────────────────────────────────────────┐ │
│ │ │ │
│ │ Producer │ │
│ │ │ │ │
│ │ ↓ │ │
│ │ delayed-exchange (type: "x-delayed-message") │ │
│ │ │ │ │
│ │ │ 消息 header: x-delay = 300000 (ms) │ │
│ │ │ │ │
│ │ │ 等待 300 秒后... │ │
│ │ │ │ │
│ │ ↓ │ │
│ │ target-queue → Consumer │ │
│ │ │ │
│ │ 优势:一个 Exchange 支持任意延迟时间 │ │
│ │ 劣势:需要安装插件,不是 RabbitMQ 标准功能 │ │
│ │ 插件使用 Mnesia 存储,重启可能丢失延迟消息 │ │
│ └─────────────────────────────────────────────────────┘ │
│ │
└─────────────────────────────────────────────────────────────┘1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
5.2 延迟插件配置与使用
typescript
// Node.js —— 使用 delayed message exchange 插件
async function setupDelayedExchange(channel: amqp.Channel): Promise<void> {
// 声明延迟交换机(类型是 x-delayed-message)
await channel.assertExchange("delayed-exchange", "x-delayed-message", {
durable: true,
arguments: {
// 延迟交换机必须指定内部使用的"真实"类型
"x-delayed-type": "direct",
},
});
// 声明目标队列
await channel.assertQueue("order-close-queue", { durable: true });
// 绑定
await channel.bindQueue("order-close-queue", "delayed-exchange", "order.close");
}
async function publishDelayedMessage(
channel: amqp.Channel,
orderId: string,
delayMs: number
): Promise<void> {
channel.publish(
"delayed-exchange",
"order.close",
Buffer.from(JSON.stringify({
orderId,
action: "closeOrder",
timestamp: Date.now(),
})),
{
persistent: true,
headers: {
// 关键:x-delay 指定延迟时间(毫秒)
"x-delay": delayMs,
},
}
);
console.log(
`📨 Published delayed message for order ${orderId}, will trigger in ${delayMs}ms`
);
}
// 下单时发送 30 分钟后的关单检查消息
await publishDelayedMessage(channel, orderId, 30 * 60 * 1000);1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
第6部分:高可用
6.1 Mirrored Queues vs Quorum Queues
┌─────────────────────────────────────────────────────────────┐
│ Mirrored Queues vs Quorum Queues │
├─────────────────────────────────────────────────────────────┤
│ │
│ 经典镜像队列 (Mirrored Queues) —— 传统 HA 方案 │
│ ┌─────────────────────────────────────────────────────┐ │
│ │ │ │
│ │ RabbitMQ 集群 │ │
│ │ ┌──────────┐ ┌──────────┐ ┌──────────┐ │ │
│ │ │ Node 1 │ │ Node 2 │ │ Node 3 │ │ │
│ │ │ Queue │ │ Queue │ │ Queue │ │ │
│ │ │ Master │←→│ Mirror │←→│ Mirror │ │ │
│ │ └──────────┘ └──────────┘ └──────────┘ │ │
│ │ │ │
│ │ 工作方式: │ │
│ │ • 一个 Master + 多个 Mirror(Slave) │ │
│ │ • 所有读写经过 Master │ │
│ │ • Master 宕机 → 最老的 Mirror 提升为 Master │ │
│ │ │ │
│ │ 优点: │ │
│ │ ✅ 读/写低延迟 │ │
│ │ ✅ 支持所有队列特性(TTL、DLX 等) │ │
│ │ │ │
│ │ 缺点: │ │
│ │ ❌ Master 同步 Mirror 是异步的 → 可能丢消息 │ │
│ │ ❌ 网络分区时有"脑裂"风险 │ │
│ │ ❌ 已被标记为 deprecated(RabbitMQ 3.9+) │ │
│ └─────────────────────────────────────────────────────┘ │
│ │
│ 仲裁队列 (Quorum Queues) —— 基于 Raft 的现代方案 │
│ ┌─────────────────────────────────────────────────────┐ │
│ │ │ │
│ │ RabbitMQ 集群 │ │
│ │ ┌──────────┐ ┌──────────┐ ┌──────────┐ │ │
│ │ │ Node 1 │ │ Node 2 │ │ Node 3 │ │ │
│ │ │ Leader │←→│Follower│←→│Follower│ │ │
│ │ │ (Raft) │ │ (Raft) │ │ (Raft) │ │ │
│ │ └──────────┘ └──────────┘ └──────────┘ │ │
│ │ │ │
│ │ 工作方式: │ │
│ │ • 基于 Raft 共识协议 │ │
│ │ • 消息写入 Leader → 复制到多数 Follower 才确认 │ │
│ │ • Leader 宕机 → Raft 自动选举新 Leader │ │
│ │ │ │
│ │ 优点: │ │
│ │ ✅ 强一致性(不丢消息) │ │
│ │ ✅ 自动故障转移(Raft 协议) │ │
│ │ ✅ 没有脑裂问题 │ │
│ │ ✅ RabbitMQ 官方推荐 │ │
│ │ │ │
│ │ 缺点: │ │
│ │ ❌ 写延迟较高(需要多数派确认) │ │
│ │ ❌ 不支持 TTL per-message、优先级等特性 │ │
│ │ ❌ 消息量大时逐条 fsync,性能低于经典队列 │ │
│ └─────────────────────────────────────────────────────┘ │
│ │
└─────────────────────────────────────────────────────────────┘1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
6.2 Quorum Queue 配置
typescript
// Node.js —— 创建仲裁队列
async function setupQuorumQueue(channel: amqp.Channel): Promise<void> {
await channel.assertQueue("critical-orders", {
durable: true,
arguments: {
// 关键参数:指定队列类型为仲裁队列
"x-queue-type": "quorum",
// 初始副本数(奇数,推荐 3 或 5)
"x-quorum-initial-group-size": 3,
// 死信交换机(仲裁队列同样支持 DLX)
"x-dead-letter-exchange": "dlx-exchange",
"x-dead-letter-routing-key": "critical-dead",
},
});
console.log("✅ Created quorum queue with 3 replicas");
}
// 对比:经典镜像队列(已废弃,仅作参考)
// await channel.assertQueue("legacy-queue", {
// durable: true,
// arguments: {
// "x-ha-policy": "all", // 在所有节点上镜像
// },
// });1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
第7部分:RabbitMQ vs Kafka 选型决策
7.1 选型对比图
┌─────────────────────────────────────────────────────────────┐
│ RabbitMQ vs Kafka 选型决策框架 │
├─────────────────────────────────────────────────────────────┤
│ │
│ 你的需求是... │
│ │
│ ┌─────────────────────────────────────────────────────┐ │
│ │ │ │
│ │ 需要灵活的路由规则? │ │
│ │ → YES: RabbitMQ(Direct/Topic/Fanout/Headers) │ │
│ │ → NO: Kafka(Key → Partition 就够了) │ │
│ │ │ │
│ │ 需要消息优先级? │ │
│ │ → YES: RabbitMQ(内置 per-message priority) │ │
│ │ → NO: Kafka(不支持,需要多 Topic 模拟) │ │
│ │ │ │
│ │ 需要延迟消息? │ │
│ │ → YES: RabbitMQ(插件或 DLX+TTL) │ │
│ │ → NO: 都可以 │ │
│ │ │ │
│ │ 需要消息回溯/重放? │ │
│ │ → YES: Kafka(消息持久化,任意 Offset 开始消费) │ │
│ │ → NO: RabbitMQ(消息消费后删除) │ │
│ │ │ │
│ │ 日处理消息量 > 百万条/秒? │ │
│ │ → YES: Kafka(设计目标就是高吞吐) │ │
│ │ → NO: RabbitMQ(万级 QPS 够用) │ │
│ │ │ │
│ │ 需要事务消息(分布式事务)? │ │
│ │ → YES: RocketMQ(内置事务消息) │ │
│ │ → 或: Kafka(事务 API,但复杂) │ │
│ │ │ │
│ │ 运维团队小(< 3 人)? │ │
│ │ → YES: RabbitMQ(部署运维简单) │ │
│ │ → NO: Kafka 也可(但运维投入大) │ │
│ │ │ │
│ └─────────────────────────────────────────────────────┘ │
│ │
└─────────────────────────────────────────────────────────────┘1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
7.2 RabbitMQ 的独特优势场景
┌─────────────────────────────────────────────────────────────┐
│ RabbitMQ 比 Kafka 更适合的场景 │
├─────────────────────────────────────────────────────────────┤
│ │
│ 场景1:复杂路由需求 │
│ ┌─────────────────────────────────────────────────────┐ │
│ │ 电商系统:日志消息需要根据级别+服务名+环境路由 │ │
│ │ Topic Exchange: "error.order-service.prod" │ │
│ │ → 运维队列: "error.#" │ │
│ │ → 开发队列: "*.order-service.*" │ │
│ │ → Kafka 模拟:需要多个 Topic + 重复发送 │ │
│ └─────────────────────────────────────────────────────┘ │
│ │
│ 场景2:RPC / 请求-响应模式 │
│ ┌─────────────────────────────────────────────────────┐ │
│ │ Client 发请求 → 等待响应 → 拿到结果 │ │
│ │ RabbitMQ 的 Direct Reply-to 机制天然支持 RPC │ │
│ │ Kafka 设计上不适合 RPC(Pull 模型 + 分区) │ │
│ └─────────────────────────────────────────────────────┘ │
│ │
│ 场景3:消息优先级队列 │
│ ┌─────────────────────────────────────────────────────┐ │
│ │ VIP 用户的消息优先处理 │ │
│ │ RabbitMQ: 设置消息 priority=10 │ │
│ │ Kafka: 需要单独的 VIP Topic + 专用 Consumer │ │
│ └─────────────────────────────────────────────────────┘ │
│ │
│ 场景4:中小团队、快速起步 │
│ ┌─────────────────────────────────────────────────────┐ │
│ │ RabbitMQ: docker-compose 一行启动,开箱即用 │ │
│ │ Kafka: 需要 ZK/KRaft + 多个 Broker + 配置调优 │ │
│ └─────────────────────────────────────────────────────┘ │
│ │
└─────────────────────────────────────────────────────────────┘1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
核心总结
总结1:RabbitMQ 核心架构
Producer → Exchange → [Binding Rules] → Queue → Consumer
↑ ↑
路由规则 存储+分发
Connection(TCP 长连接)
└── Channel(轻量级虚拟连接,多路复用)1
2
3
4
5
6
2
3
4
5
6
总结2:四种 Exchange
| 类型 | 路由方式 | 使用场景 |
|---|---|---|
| Direct | Routing Key 精确匹配 | 单播、任务分配 |
| Topic | Routing Key 模式匹配 (*/#) | 多级分类 |
| Fanout | 广播到所有绑定 Queue | 事件广播 |
| Headers | Header 属性匹配 (x-match) | 复杂条件路由 |
总结3:消息可靠性链路
Publisher Confirm(发端保证 Broker 收到)
→ Queue 持久化 + 消息持久化(Broker 保证不丢)
→ Consumer Ack/Nack(收端保证处理完成)
→ 失败消息 → Dead Letter Queue(兜底)1
2
3
4
2
3
4
总结4:延迟队列
方案A:死信 + TTL → 不需要插件,但不同延迟需要多个队列
方案B:delayed-message-exchange 插件 → 一个 Exchange 支持任意延迟
选择:生产环境推荐插件方案;简单场景用方案 A 也能跑1
2
3
2
3
总结5:高可用
经典镜像队列 → 已废弃,异步复制可能丢消息
仲裁队列 (Quorum Queue) → 基于 Raft,强一致性,官方推荐1
2
2
总结6:选型一句话
需要灵活路由/延迟消息/RPC/优先级 → RabbitMQ
需要高吞吐/消息回溯/流处理 → Kafka
两个都要 → Kafka 做主链路 + RabbitMQ 做灵活路由辅助1
2
3
2
3
章节测试
测试1:Channel 的作用是什么?为什么需要它?
测试2:一条 Routing Key 为 order.us.electronics 的消息,能匹配以下哪些 Binding?
A. order.*.* B. order.# C. order.us.* D. payment.# E. *.us.electronics
测试3:Publisher Confirm 的三种使用方式分别是什么?
测试4:Consumer 什么时候应该用 nack(requeue=false)?
测试5:消息变成死信的三种条件是什么?
测试6:Quorum Queue 和 Mirrored Queue 的核心区别是什么?
测试7:以下场景适合用 RabbitMQ 还是 Kafka?
A. 电商订单系统:下单后 30 分钟未支付自动取消 B. 千万级用户行为日志采集和分析 C. VIP 用户的客服工单需要优先处理 D. 微服务间的 RPC 调用
参考答案
测试1答案
答案:
- Channel 是复用 TCP Connection 的轻量级"虚拟连接"
- 为什么需要:避免为每个操作(发布/消费)创建新的 TCP 连接
- 一个 Connection 上可以有多个 Channel,各自独立通信
- 建议:一个线程一个 Channel(Channel 不是线程安全的)
测试2答案
答案:
- A
order.*.*✅(*匹配恰好一段:us、electronics 共两段) - B
order.#✅(#匹配零段或多段) - C
order.us.*✅(匹配 us 后恰好一段:electronics) - D
payment.#❌(第一段不匹配) - E
*.us.electronics✅(*匹配恰好一段:order,然后 us.electronics 精确匹配)
测试3答案
答案:
- 逐条确认:发一条等一条确认,吞吐最低,实现最简单
- 批量确认:发 N 条后等一批确认,延迟较高,吞吐中等
- 异步确认:发消息不等待,异步回调处理确认结果,推荐方案
测试4答案
答案: 当消息达到最大重试次数仍然无法处理时,使用 nack(requeue=false):
requeue=false→ 消息不重新入队,结合 DLX 进入死信队列- 防止坏消息无限重试、阻塞队列
- 由死信队列的专属 Consumer 或人工介入处理
测试5答案
答案:
- 消息被拒绝:Consumer 调用
basic.reject或basic.nack且requeue=false - 消息 TTL 超时:消息在队列中超过设置的 TTL 时间仍未被消费
- 队列已满:队列达到最大长度,头部的旧消息被挤出
测试6答案
答案:
| 维度 | Mirrored Queue | Quorum Queue |
|---|---|---|
| 一致性协议 | 异步复制(无共识) | Raft(强一致性) |
| 数据可靠性 | 可能丢消息 | 不丢消息(多数派确认) |
| 脑裂风险 | 有 | 无(Raft 保证) |
| 性能 | 读写延迟低 | 写延迟较高 |
| 状态 | 已废弃(deprecated) | 官方推荐 |
测试7答案
答案:
- A:RabbitMQ(延迟消息是 RabbitMQ 的强项)
- B:Kafka(高吞吐日志管道,内置流处理)
- C:RabbitMQ(消息优先级队列)
- D:RabbitMQ(Direct Reply-to 天然支持 RPC)
相关笔记
- [[00-overview]] - 消息队列全景
- [[/03-web/07-middleware/02-mq-kafka]] - Kafka 核心机制
- [[/03-web/07-middleware/03-mq-others]] - 消息队列对比选型
- [[/03-web/07-middleware/07-mq-advanced]] - MQ 高级特性
- [[03-kafka-streams-connect]] - Kafka Streams 与 Connect
下一步学习
- [ ] 阅读 [[/03-web/07-middleware/07-mq-advanced]] - 死信队列、延迟消息、消息积压
- [ ] 实践:用 Docker Compose 搭建 RabbitMQ + 延迟队列 + 死信队列的完整 demo
学习状态:🟡 开始学习