消息队列全景——为什么你的系统需要一个中间人 / Message Queue Overview and Why Your System Needs a Middleman
📅 创建时间:2026-07-28 🏷️ 标签:#消息队列 #MQ #Kafka #RabbitMQ #异步解耦 #削峰填谷 #消息可靠性 📚 前置知识:[[/03-web/07-middleware/00-overview]](中间件全景) [[/03-web/07-middleware/01-redis-deep]](Redis 缓存) 📚 相关知识:[[/03-web/07-middleware/02-mq-kafka]](Kafka 核心机制) [[/03-web/07-middleware/03-mq-others]](MQ 对比选型) [[/03-web/07-middleware/07-mq-advanced]](MQ 高级特性)
📋 本章目标
- 理解消息队列在分布式系统中的三大核心价值:异步解耦、削峰填谷、最终一致性
- 掌握 MQ 核心概念:Producer、Broker、Consumer、Topic、Queue、Partition、Consumer Group
- 区分三种消息模型:队列模型、发布-订阅模型、流模型
- 理解消息可靠性的四层保障:生产者确认、Broker 持久化、消费者确认、死信队列
- 掌握三种消费语义的区别:at-most-once、at-least-once、exactly-once
- 理解消息顺序性问题的本质,以及分区有序的解决方案
- 建立 Kafka vs RabbitMQ vs RocketMQ vs Redis Streams 的全景认知
场景:为什么订单总是"丢了"
┌─────────────────────────────────────────────────────────────┐
│ │
│ 双十一零点,秒杀活动开始。 │
│ │
│ 10 万用户同时点击"下单"。 │
│ 后端服务线程池只有 200 个线程。 │
│ │
│ 第 0.1 秒:200 个请求进入处理,其余 99800 排队。 │
│ 第 1 秒:已有 5000 个请求超时,用户看到"系统繁忙"。 │
│ 第 5 秒:数据库连接池耗尽,所有请求报 500 错误。 │
│ │
│ 结果: │
│ 10 万个下单请求,成功 2000 个,丢失 3000 个, │
│ 重复 1000 个,剩下 94000 个直接报错。 │
│ │
│ 老板:"为什么不能先把请求收下来,慢慢处理?" │
│ │
└─────────────────────────────────────────────────────────────┘1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
这就是消息队列要解决的核心问题:请求不能丢,也不能让下游崩掉。
第1部分:为什么需要消息队列?
1.1 三大核心价值
┌─────────────────────────────────────────────────────────────┐
│ 消息队列的三大核心价值 │
├─────────────────────────────────────────────────────────────┤
│ │
│ ┌─────────────────────────────────────────────────────┐ │
│ │ 价值1:异步解耦(Async Decoupling) │ │
│ │ │ │
│ │ 同步调用:A → B(A 必须等 B 返回) │ │
│ │ 异步解耦:A → MQ → B(A 发完就完事) │ │
│ │ │ │
│ │ 好处:A 和 B 不再直接依赖 │ │
│ │ • B 挂了,A 不受影响 │ │
│ │ • 加新消费者 C,A 不需要改代码 │ │
│ │ • B 升级维护,消息暂存 MQ,不丢 │ │
│ └─────────────────────────────────────────────────────┘ │
│ │
│ ┌─────────────────────────────────────────────────────┐ │
│ │ 价值2:削峰填谷(Peak Shaving / Valley Filling) │ │
│ │ │ │
│ │ 请求量 │ │
│ │ ▲ │ │
│ │ │ ╱╲ │ │
│ │ │ ╱ ╲ ← 峰值 10 万 QPS,DB 直接打挂 │ │
│ │ │╱ ╲___ │ │
│ │ └──────────→ 时间 │ │
│ │ │ │
│ │ 加 MQ 后: │ │
│ │ 峰值请求 → MQ 缓冲 → Consumer 匀速消费 │ │
│ │ 高峰期积压的消息,低谷期慢慢消化 │ │
│ └─────────────────────────────────────────────────────┘ │
│ │
│ ┌─────────────────────────────────────────────────────┐ │
│ │ 价值3:最终一致性(Eventual Consistency) │ │
│ │ │ │
│ │ 下单 → 扣库存 → 发短信 → 加积分 → 更新 BI │ │
│ │ │ │
│ │ 这些操作不需要在同一事务里强一致完成 │ │
│ │ 只要最终都完成了,业务就是正确的 │ │
│ │ 消息队列保证:"一定会处理,但不一定立刻处理" │ │
│ └─────────────────────────────────────────────────────┘ │
│ │
└─────────────────────────────────────────────────────────────┘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
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
1.2 同步 vs 异步:一个下单场景的对比
┌─────────────────────────────────────────────────────────────┐
│ 同步调用 vs 异步消息 │
├─────────────────────────────────────────────────────────────┤
│ │
│ 同步调用(无 MQ): │
│ │
│ 用户下单 ──→ 订单服务 ──→ 库存服务(等 50ms) │
│ ──→ 支付服务(等 100ms) │
│ ──→ 短信服务(等 80ms) │
│ ──→ 积分服务(等 60ms) │
│ │
│ 总响应时间 = 50 + 100 + 80 + 60 = 290ms │
│ 任何一步失败 → 整个下单流程中断 │
│ │
│ ─────────────────────────────────────────────────────── │
│ │
│ 异步消息(有 MQ): │
│ │
│ 用户下单 ──→ 订单服务 ──→ MQ ──→ 库存服务(异步) │
│ │ ──→ 支付服务(异步) │
│ │ ──→ 短信服务(异步) │
│ │ ──→ 积分服务(异步) │
│ │ │
│ └── 立即返回 "下单成功"(50ms) │
│ │
│ 后续各服务独立消费,互不影响 │
│ │
└─────────────────────────────────────────────────────────────┘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
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
1.3 不用 MQ 能行吗?Redis List 行不行?
┌─────────────────────────────────────────────────────────────┐
│ Redis List 做队列 vs 专业 MQ │
├─────────────────────────────────────────────────────────────┤
│ │
│ Redis List(LPUSH / BRPOP)可以做简单队列: │
│ │
│ Producer: LPUSH orders:queue '{"orderId": 123}' │
│ Consumer: BRPOP orders:queue 0 │
│ │
│ 但有以下致命缺陷: │
│ │
│ ❌ 没有消费者确认机制 —— 消费失败消息就丢了 │
│ ❌ 没有消息持久化保证 —— Redis 重启/RDB 间隔内数据丢失 │
│ ❌ 没有重试/死信机制 —— 坏消息阻塞整个队列 │
│ ❌ 没有 Consumer Group —— 无法做并发消费和 Rebalance │
│ ❌ 不支持消息回溯 —— 消费了就没了 │
│ ❌ 数据量受内存限制 —— 10GB 消息就 OOM │
│ │
│ 结论:简单场景可以用 Redis List(如任务队列), │
│ 但核心业务链路必须用专业 MQ。 │
│ │
└─────────────────────────────────────────────────────────────┘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
第2部分:MQ 核心概念
2.1 核心角色
┌─────────────────────────────────────────────────────────────┐
│ MQ 核心角色关系图 │
├─────────────────────────────────────────────────────────────┤
│ │
│ ┌──────────┐ ┌──────────────────┐ ┌──────────┐│
│ │ │ │ │ │ ││
│ │ Producer │──────→│ Broker │──────→│ Consumer ││
│ │ (生产者) │ │ (消息代理) │ │ (消费者) ││
│ │ │ │ │ │ ││
│ └──────────┘ │ ┌────────────┐ │ └──────────┘│
│ │ │ Topic A │ │ │
│ ┌──────────┐ │ │ ┌────┬────┐│ │ ┌──────────┐│
│ │ │ │ │ │ P0 │ P1 ││ │ │ ││
│ │ Producer │──────→│ │ └────┴────┘│ │──────→│ Consumer ││
│ │ │ │ └────────────┘ │ │ ││
│ └──────────┘ │ │ └──────────┘│
│ └──────────────────┘ │
│ │
└─────────────────────────────────────────────────────────────┘1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
2.2 核心概念详解
| 概念 | 类比 | 说明 |
|---|---|---|
| Producer | 寄件人 | 生产并发送消息的应用 |
| Broker | 邮局 | 接收、存储、转发消息的中间节点 |
| Consumer | 收件人 | 接收并处理消息的应用 |
| Topic | 邮件主题分类 | 消息的逻辑分类,生产者发到 Topic,消费者订阅 Topic |
| Partition | 邮局的分拣台 | Topic 的物理分片,一个 Topic 可以有多个 Partition |
| Consumer Group | 收件小组 | 一组消费者协同消费,同组内每个 Partition 只被一个 Consumer 消费 |
| Offset | 邮件编号 | 消息在 Partition 中的唯一序号,消费者用 Offset 记录消费位置 |
| Queue | 私人信箱(RabbitMQ 概念) | 消息的存储单元,Consumer 直接从 Queue 消费 |
┌─────────────────────────────────────────────────────────────┐
│ Topic / Partition / Consumer Group 的关系 │
├─────────────────────────────────────────────────────────────┤
│ │
│ Topic: "orders" │
│ ┌─────────────────────────────────────────────────────┐ │
│ │ Partition 0: [msg0] [msg1] [msg2] [msg3] ... │ │
│ │ Partition 1: [msg4] [msg5] [msg6] [msg7] ... │ │
│ │ Partition 2: [msg8] [msg9] [msgA] [msgB] ... │ │
│ └─────────────────────────────────────────────────────┘ │
│ │
│ Consumer Group "order-processors": │
│ ┌─────────────────────────────────────────────────────┐ │
│ │ Consumer A → 消费 Partition 0 │ │
│ │ Consumer B → 消费 Partition 1 │ │
│ │ Consumer C → 消费 Partition 2 │ │
│ └─────────────────────────────────────────────────────┘ │
│ │
│ 关键规则: │
│ • Partition 数 ≥ Consumer 数(同组),否则有 Consumer 空闲 │
│ • 同一 Partition 内消息有序 │
│ • 不同 Partition 间消息无序 │
│ • 不同 Consumer Group 独立消费同一 Topic │
│ │
└─────────────────────────────────────────────────────────────┘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
TypeScript 示例 —— Producer 和 Consumer 的基本形态:
typescript
// Producer 端:发送用户注册事件
interface UserRegisteredEvent {
userId: string;
email: string;
registeredAt: string;
source: "web" | "app" | "api";
}
async function publishUserRegistered(
producer: KafkaProducer,
event: UserRegisteredEvent
): Promise<void> {
await producer.send({
topic: "user-events",
messages: [{
// key 决定分区路由 —— 同一 userId 的事件进同一分区
key: event.userId,
value: JSON.stringify(event),
headers: {
"event-type": "user.registered",
"source": event.source,
},
}],
});
}
// Consumer 端:发送欢迎邮件
async function sendWelcomeEmail(consumer: KafkaConsumer): Promise<void> {
await consumer.subscribe({ topic: "user-events" });
await consumer.run({
eachMessage: async ({ topic, partition, message }) => {
const event: UserRegisteredEvent = JSON.parse(
message.value!.toString()
);
if (message.headers?.["event-type"]?.toString() === "user.registered") {
await emailService.sendWelcomeEmail(event.email);
}
// 消息默认自动提交 offset
// 手动提交可以获得更好的可靠性控制
},
});
}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
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
第3部分:消息模型对比
3.1 三种消息模型
┌─────────────────────────────────────────────────────────────┐
│ 队列模型 vs 发布-订阅模型 vs 流模型 │
├─────────────────────────────────────────────────────────────┤
│ │
│ ① 队列模型(Queue Model)—— RabbitMQ 经典模式 │
│ ┌─────────────────────────────────────────────────────┐ │
│ │ │ │
│ │ Producer ──→ Queue ──→ Consumer A │ │
│ │ └──→ Consumer B(竞争消费) │ │
│ │ │ │
│ │ 特点: │ │
│ │ • 一条消息只能被一个 Consumer 消费(竞争模式) │ │
│ │ • 消费后消息删除(或确认后删除) │ │
│ │ • 适合:任务分发、RPC 命令 │ │
│ └─────────────────────────────────────────────────────┘ │
│ │
│ ② 发布-订阅模型(Pub/Sub Model)—— 所有 MQ 都支持 │
│ ┌─────────────────────────────────────────────────────┐ │
│ │ │ │
│ │ Producer ──→ Topic ──→ Consumer Group A(订单组) │ │
│ │ └──→ Consumer Group B(风控组) │ │
│ │ └──→ Consumer Group C(分析组) │ │
│ │ │ │
│ │ 特点: │ │
│ │ • 一条消息被多个 Consumer Group 各自消费一次 │ │
│ │ • 每个 Consumer Group 管理自己的消费位置 │ │
│ │ • 适合:事件广播、多系统同步 │ │
│ └─────────────────────────────────────────────────────┘ │
│ │
│ ③ 流模型(Stream Model)—— Kafka 核心 │
│ ┌─────────────────────────────────────────────────────┐ │
│ │ │ │
│ │ Producer ──→ 不可变日志(Immutable Log) │ │
│ │ │ │ │
│ │ │ msg0 msg1 msg2 msg3 msg4 ... │ │
│ │ │ 0 1 2 3 4 │ │
│ │ │ ↑ ↑ │ │
│ │ │ Consumer A (offset=1) │ │ │
│ │ │ Consumer B (offset=4) │ │
│ │ │ │
│ │ 特点: │ │
│ │ • 消息不删除,持久化在磁盘(可回溯重放) │ │
│ │ • Consumer 自行管理 Offset(消费到哪了) │ │
│ │ • 适合:事件溯源、大数据管道、日志聚合 │ │
│ └─────────────────────────────────────────────────────┘ │
│ │
└─────────────────────────────────────────────────────────────┘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 各模型典型代表
| 模型 | 代表 MQ | 消息保留 | 消费模式 | 典型场景 |
|---|---|---|---|---|
| 队列模型 | RabbitMQ | 消费后删除 | 竞争消费 | 任务分发、RPC |
| 发布-订阅 | 所有 MQ | 取决于实现 | 广播 | 事件通知、数据同步 |
| 流模型 | Kafka | 持久化(可配置保留时间) | 拉模式 | 日志聚合、事件溯源 |
第4部分:消息可靠性
4.1 消息可能在哪里丢失?
┌─────────────────────────────────────────────────────────────┐
│ 消息丢失的三个环节 │
├─────────────────────────────────────────────────────────────┤
│ │
│ ┌──────────┐ ┌──────────┐ ┌──────────┐ │
│ │ │ ① │ │ ② │ │ │
│ │ Producer │───→───│ Broker │───→───│ Consumer │ │
│ │ │ │ │ │ │ │
│ └──────────┘ └──────────┘ └──────────┘ │
│ │
│ ① 发送环节:Producer → Broker │
│ 网络超时,Producer 认为发送失败但实际 Broker 收到了 │
│ → 可能重复发送,也可能真正丢失 │
│ │
│ ② 存储环节:Broker 内部 │
│ Leader 宕机,Follower 还没同步,消息丢失 │
│ → Kafka: acks=all + min.insync.replicas ≥ 2 │
│ → RabbitMQ: 持久化队列 + 持久化消息 + Publisher Confirm │
│ │
│ ③ 消费环节:Broker → Consumer │
│ Consumer 拉到消息后还没处理完就崩溃 │
│ → 需要消费者确认机制 │
│ │
└─────────────────────────────────────────────────────────────┘1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
4.2 四层可靠性保障
┌─────────────────────────────────────────────────────────────┐
│ 消息可靠性的四层保障 │
├─────────────────────────────────────────────────────────────┤
│ │
│ 第1层:生产者确认(Producer Acknowledgement) │
│ ┌─────────────────────────────────────────────────────┐ │
│ │ │ │
│ │ Kafka acks 配置: │ │
│ │ acks=0: 发出去就不管(可能丢) │ │
│ │ acks=1: Leader 确认即成功(Leader 宕机可能丢) │ │
│ │ acks=all: 所有 ISR 副本确认(最安全) │ │
│ │ │ │
│ │ RabbitMQ Publisher Confirm: │ │
│ │ 生产者发送后等待 Broker 的 confirm 回执 │ │
│ │ 未收到 confirm → 重发 │ │
│ └─────────────────────────────────────────────────────┘ │
│ │
│ 第2层:Broker 持久化(Persistence) │
│ ┌─────────────────────────────────────────────────────┐ │
│ │ │ │
│ │ Kafka:消息写入 Page Cache + 刷盘策略 │ │
│ │ • 默认依赖 OS 刷盘(可能丢末尾数据) │ │
│ │ • 关键业务配置 flush.messages=1 每一条都刷盘 │ │
│ │ │ │
│ │ RabbitMQ: │ │
│ │ • durable=true 队列(重启后队列还在) │ │
│ │ • persistent=true 消息(重启后消息还在) │ │
│ └─────────────────────────────────────────────────────┘ │
│ │
│ 第3层:消费者确认(Consumer Acknowledgement) │
│ ┌─────────────────────────────────────────────────────┐ │
│ │ │ │
│ │ 自动提交(Auto Commit): │ │
│ │ Consumer 拉到消息 → 自动提交 offset → 处理 │ │
│ │ 风险:处理失败,消息已"消费",实际丢失 │ │
│ │ │ │
│ │ 手动提交(Manual Commit): │ │
│ │ Consumer 拉到消息 → 处理 → 成功则提交 offset │ │
│ │ 风险:处理成功但提交失败 → 重复消费(需要幂等) │ │
│ └─────────────────────────────────────────────────────┘ │
│ │
│ 第4层:死信队列(Dead Letter Queue) │
│ ┌─────────────────────────────────────────────────────┐ │
│ │ │ │
│ │ 消息重试 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
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
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
TypeScript 示例 —— 手动提交实现可靠消费:
typescript
async function reliableConsumer(consumer: KafkaConsumer): Promise<void> {
await consumer.subscribe({ topic: "orders" });
await consumer.run({
// 关闭自动提交
autoCommit: false,
eachMessage: async ({ topic, partition, message }) => {
const MAX_RETRIES = 3;
for (let attempt = 1; attempt <= MAX_RETRIES; attempt++) {
try {
const order = JSON.parse(message.value!.toString());
await orderService.process(order);
// 处理成功 → 手动提交 offset
await consumer.commitOffsets([
{ topic, partition, offset: (Number(message.offset) + 1).toString() },
]);
console.log(`✅ Processed offset ${message.offset}`);
return; // 成功,跳出循环
} catch (error) {
console.error(
`❌ Attempt ${attempt}/${MAX_RETRIES} failed for offset ${message.offset}`,
error
);
if (attempt === MAX_RETRIES) {
// 达到最大重试次数 → 发送到死信 Topic
await sendToDLQ({
originalTopic: topic,
offset: message.offset,
payload: message.value!.toString(),
error: String(error),
});
// 即使失败了也提交 offset,避免阻塞
await consumer.commitOffsets([
{ topic, partition, offset: (Number(message.offset) + 1).toString() },
]);
} else {
// 等待后重试
await sleep(attempt * 1000);
}
}
}
},
});
}
async function sendToDLQ(record: DLQRecord): Promise<void> {
// 写入独立的死信 Topic,供人工排查
await kafka.producer().send({
topic: "orders-dlq",
messages: [{
key: record.offset,
value: JSON.stringify(record),
}],
});
}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
60
61
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
60
61
第5部分:消费语义
5.1 三种语义的定义
┌─────────────────────────────────────────────────────────────┐
│ At-Most-Once / At-Least-Once / Exactly-Once │
├─────────────────────────────────────────────────────────────┤
│ │
│ ① At-Most-Once(最多一次) │
│ ┌─────────────────────────────────────────────────────┐ │
│ │ │ │
│ │ Producer 发消息 → 不重试 │ │
│ │ Consumer 自动提交 → 处理失败消息丢失 │ │
│ │ │ │
│ │ 结果:消息可能丢失,但绝不重复 │ │
│ │ 适用:监控指标、日志采集(丢几条无所谓) │ │
│ │ 可靠性:⭐⭐ │ │
│ └─────────────────────────────────────────────────────┘ │
│ │
│ ② At-Least-Once(至少一次) │
│ ┌─────────────────────────────────────────────────────┐ │
│ │ │ │
│ │ Producer 发消息 → 超时重试 → 可能重复 │ │
│ │ Consumer 处理完再提交 → 提交失败则重复消费 │ │
│ │ │ │
│ │ 结果:消息绝不丢失,但可能重复 │ │
│ │ 适用:大多数业务场景(配合幂等) │ │
│ │ 可靠性:⭐⭐⭐⭐ │ │
│ └─────────────────────────────────────────────────────┘ │
│ │
│ ③ Exactly-Once(精确一次) │
│ ┌─────────────────────────────────────────────────────┐ │
│ │ │ │
│ │ "消息既不会丢,也不会重复" │ │
│ │ → 实际上是一个分布式系统的圣杯,实现极其复杂 │ │
│ │ │ │
│ │ Kafka 实现方式: │ │
│ │ • 幂等生产者(Idempotent Producer): │ │
│ │ Producer 给每条消息分配 PID + SeqNum │ │
│ │ Broker 去重 │ │
│ │ • 事务(Transactions): │ │
│ │ 原子写入多个 Partition + 消费者端事务提交 │ │
│ │ │ │
│ │ 现实建议: │ │
│ │ 大多数场景 at-least-once + 消费者幂等 = 实际上的 │ │
│ │ exactly-once,远比 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
40
41
42
43
44
45
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
5.2 幂等消费 —— 让 At-Least-Once 变成 Effective Exactly-Once
┌─────────────────────────────────────────────────────────────┐
│ 幂等消费的实现策略 │
├─────────────────────────────────────────────────────────────┤
│ │
│ 策略1:数据库唯一约束(最常用) │
│ ┌─────────────────────────────────────────────────────┐ │
│ │ │ │
│ │ 消息里带一个业务唯一 ID(如 orderId) │ │
│ │ INSERT INTO orders (order_id, ...) VALUES (...) │ │
│ │ → 重复消息 INSERT 时报唯一约束冲突 → 忽略 │ │
│ │ │ │
│ │ 优点:简单、可靠 │ │
│ │ 缺点:仅适用于写操作 │ │
│ └─────────────────────────────────────────────────────┘ │
│ │
│ 策略2:Redis 去重(高性能场景) │
│ ┌─────────────────────────────────────────────────────┐ │
│ │ │ │
│ │ const dedupKey = `msg:dedup:${messageId}`; │ │
│ │ const exists = await redis.setnx(dedupKey, "1"); │ │
│ │ await redis.expire(dedupKey, 86400); // 24h 过期 │ │
│ │ │ │
│ │ if (!exists) { │ │
│ │ return; // 已处理过,跳过 │ │
│ │ } │ │
│ │ await process(message); │ │
│ │ │ │
│ │ 优点:高性能、不影响 DB │ │
│ │ 缺点:Redis 和 DB 不在同一个事务中 │ │
│ └─────────────────────────────────────────────────────┘ │
│ │
│ 策略3:版本号 / 状态机(适合更新操作) │
│ ┌─────────────────────────────────────────────────────┐ │
│ │ │ │
│ │ UPDATE orders SET status='paid', version=version+1│ │
│ │ WHERE order_id = ? AND version = ?; │ │
│ │ │ │
│ │ 重复执行 UPDATE 不会改变结果(天然幂等) │ │
│ └─────────────────────────────────────────────────────┘ │
│ │
└─────────────────────────────────────────────────────────────┘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
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
第6部分:消息顺序性
6.1 顺序性问题的本质
┌─────────────────────────────────────────────────────────────┐
│ 消息乱序是怎么发生的? │
├─────────────────────────────────────────────────────────────┤
│ │
│ 场景:用户的操作序列 │
│ ① 用户注册 → ② 修改昵称 → ③ 修改头像 │
│ │
│ 发送到 Topic,3 个 Partition: │
│ │
│ Partition 0: [注册] ────────────→ Consumer A │
│ Partition 1: [修改昵称] ────────→ Consumer B │
│ Partition 2: [修改头像] ────────→ Consumer C │
│ │
│ 三个 Consumer 并发消费,各自速度不同: │
│ Consumer B 先处理完"修改昵称"(5ms) │
│ Consumer C 接着处理完"修改头像"(10ms) │
│ Consumer A 最后处理完"注册"(因为连 DB,50ms) │
│ │
│ 结果:"修改头像"发生在"注册"之前 → 数据乱套了! │
│ │
└─────────────────────────────────────────────────────────────┘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
6.2 解决方案:分区有序
┌─────────────────────────────────────────────────────────────┐
│ 分区有序 —— 用 Key 保证同一实体有序 │
├─────────────────────────────────────────────────────────────┤
│ │
│ 核心思想:同一个业务 Key 的消息进同一个 Partition │
│ │
│ ┌─────────────────────────────────────────────────────┐ │
│ │ │ │
│ │ 消息 Key = userId(如 "user_123") │ │
│ │ │ │
│ │ hash("user_123") % 3 = 1 → Partition 1 │ │
│ │ │ │
│ │ Partition 1: │ │
│ │ [用户注册] [修改昵称] [修改头像] │ │
│ │ ↓ ↓ ↓ │ │
│ │ └───────────┴───────────┘ │ │
│ │ ↓ │ │
│ │ Consumer B 顺序消费 │ │
│ │ │ │
│ │ 注册 → 修改昵称 → 修改头像 严格有序! │ │
│ └─────────────────────────────────────────────────────┘ │
│ │
│ 代价和注意事项: │
│ ┌─────────────────────────────────────────────────────┐ │
│ │ ⚠️ 同一 Partition 的消息串行处理 → 影响吞吐 │ │
│ │ ⚠️ 热点 Key 问题(如一个大 V 的所有消息进同一分区)│ │
│ │ ⚠️ Partition 数不能动态减少(Kafka) │ │
│ │ ⚠️ Rebalance 时可能短暂乱序 │ │
│ │ │ │
│ │ 最佳实践: │ │
│ │ • 只对真正需要顺序的业务用 Key 路由 │ │
│ │ • 大多数场景不需要全局有序,分区有序就够了 │ │
│ │ • 评估业务能否容忍 Rebalance 期间的短暂乱序 │ │
│ └─────────────────────────────────────────────────────┘ │
│ │
└─────────────────────────────────────────────────────────────┘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
TypeScript 示例 —— 使用 Key 保证用户操作有序:
typescript
// Producer:同一用户的操作用 userId 做 Key,路由到同一分区
async function publishUserAction(
producer: KafkaProducer,
userId: string,
action: UserAction
): Promise<void> {
await producer.send({
topic: "user-actions",
messages: [{
// 关键:用 userId 做 Key,保证同一用户的消息有序
key: userId,
value: JSON.stringify({
userId,
action: action.type, // "register" | "updateProfile" | "uploadAvatar"
payload: action.payload,
timestamp: Date.now(),
}),
}],
});
}
// Consumer:单线程处理同一分区,保证消费顺序
//(Kafka 保证了同一 Partition 内的消息顺序投递)
async function processUserActions(consumer: KafkaConsumer): Promise<void> {
await consumer.subscribe({ topic: "user-actions" });
await consumer.run({
// 关键:一次只拉一条消息,手动确认,保证不乱序
eachMessage: async ({ message }) => {
const action = JSON.parse(message.value!.toString());
const userId = message.key!.toString();
console.log(`Processing [${action.action}] for user ${userId}`);
switch (action.action) {
case "register":
await userService.register(userId, action.payload);
break;
case "updateProfile":
await userService.updateProfile(userId, action.payload);
break;
case "uploadAvatar":
await userService.uploadAvatar(userId, action.payload);
break;
}
},
});
}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
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
第7部分:Kafka vs RabbitMQ vs RocketMQ vs Redis Streams 全景对比
7.1 全景对比表
┌─────────────────────────────────────────────────────────────┐
│ 四大 MQ 全景对比 │
├─────────────────────────────────────────────────────────────┤
│ │
│ ┌──────────┬──────────┬──────────┬──────────┬──────────┐ │
│ │ 维度 │ Kafka │ RabbitMQ │ RocketMQ │ Redis │ │
│ │ │ │ │ │ Streams │ │
│ ├──────────┼──────────┼──────────┼──────────┼──────────┤ │
│ │ 设计 │ 分布式 │ 经典 │ 分布式 │ 内存 │ │
│ │ 哲学 │ 流平台 │ 消息代理 │ 消息平台 │ 数据结构 │ │
│ ├──────────┼──────────┼──────────┼──────────┼──────────┤ │
│ │ 吞吐量 │ 百万/秒 │ 万/秒 │ 十万/秒 │ 十万/秒 │ │
│ ├──────────┼──────────┼──────────┼──────────┼──────────┤ │
│ │ 延迟 │ 毫秒级 │ 微秒级 │ 毫秒级 │ 微秒级 │ │
│ ├──────────┼──────────┼──────────┼──────────┼──────────┤ │
│ │ 消息 │ 持久化 │ 消费后 │ 持久化 │ 消费后 │ │
│ │ 保留 │ 可回溯 │ 删除 │ 可回溯 │ 删除 │ │
│ ├──────────┼──────────┼──────────┼──────────┼──────────┤ │
│ │ 消费 │ Pull │ Push │ Pull │ Pull │ │
│ │ 模式 │ │ │ │ (XREAD) │ │
│ ├──────────┼──────────┼──────────┼──────────┼──────────┤ │
│ │ 路由 │ Key→ │ Exchange │ Tag │ Stream │ │
│ │ 灵活性 │ Partition│ +Binding │ 过滤 │ +Group │ │
│ ├──────────┼──────────┼──────────┼──────────┼──────────┤ │
│ │ 事务 │ 支持 │ 不支持 │ 支持 │ 不支持 │ │
│ │ 消息 │ │ │(强) │ │ │
│ ├──────────┼──────────┼──────────┼──────────┼──────────┤ │
│ │ 延迟 │ 弱 │ 强 │ 强 │ 弱 │ │
│ │ 消息 │ │ │ │ │ │
│ ├──────────┼──────────┼──────────┼──────────┼──────────┤ │
│ │ 运维 │ 重 │ 轻 │ 中 │ 最轻 │ │
│ │ 复杂度 │ │ │ │ │ │
│ ├──────────┼──────────┼──────────┼──────────┼──────────┤ │
│ │ 典型 │ 大数据 │ 业务 │ 电商 │ 轻量级 │ │
│ │ 场景 │ 日志流 │ 任务队列 │ 金融 │ 实时 │ │
│ └──────────┴──────────┴──────────┴──────────┴──────────┘ │
│ │
└─────────────────────────────────────────────────────────────┘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
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
7.2 各 MQ 的核心差异
┌─────────────────────────────────────────────────────────────┐
│ 四大 MQ 的核心设计差异 │
├─────────────────────────────────────────────────────────────┤
│ │
│ Kafka —— 分布式流平台 │
│ ┌─────────────────────────────────────────────────────┐ │
│ │ • 核心抽象:不可变的有序日志(Immutable Log) │ │
│ │ • 设计目标:高吞吐、持久化、可回溯 │ │
│ │ • 独特优势: │ │
│ │ - 百万级 QPS(顺序磁盘 IO + zero-copy) │ │
│ │ - 消息可回溯重放(适合事件溯源、大数据) │ │
│ │ - Kafka Streams / Connect 生态 │ │
│ │ • 主要局限: │ │
│ │ - 运维重(ZooKeeper/KRaft、多 Broker 协调) │ │
│ │ - 延迟消息支持弱(没有内置延迟队列) │ │
│ │ - 消息数量多时性能下降(海量小消息) │ │
│ │ - 不适合需要灵活路由的场景 │ │
│ └─────────────────────────────────────────────────────┘ │
│ │
│ RabbitMQ —— 经典消息代理 │
│ ┌─────────────────────────────────────────────────────┐ │
│ │ • 核心抽象:Exchange + Queue + Binding │ │
│ │ • 设计目标:灵活路由、AMQP 标准、功能丰富 │ │
│ │ • 独特优势: │ │
│ │ - 灵活的路由规则(Direct/Topic/Fanout/Headers) │ │
│ │ - 成熟稳定(2007 年至今) │ │
│ │ - 延迟消息插件(rabbitmq_delayed_message_exchange)│ │
│ │ - 运维简单,Erlang/OTP 天然高可用 │ │
│ │ • 主要局限: │ │
│ │ - 吞吐量有限(单机万级 QPS) │ │
│ │ - 消息无法回溯 │ │
│ │ - 集群扩展不如 Kafka 方便 │ │
│ └─────────────────────────────────────────────────────┘ │
│ │
│ RocketMQ —— 阿里开源的分布式消息平台 │
│ ┌─────────────────────────────────────────────────────┐ │
│ │ • 核心抽象:Topic + Queue + Tag │ │
│ │ • 设计目标:万亿级消息容量、金融级可靠性 │ │
│ │ • 独特优势: │ │
│ │ - 事务消息(分布式事务场景的利器) │ │
│ │ - 延迟消息 18 个等级(1s → 2h) │ │
│ │ - 消息过滤(Tag + SQL 表达式) │ │
│ │ - 中国社区活跃,中文文档完善 │ │
│ │ • 主要局限: │ │
│ │ - 国际社区不如 Kafka 活跃 │ │
│ │ - 生态工具链不如 Kafka 丰富 │ │
│ │ - 中小项目太重 │ │
│ └─────────────────────────────────────────────────────┘ │
│ │
│ Redis Streams —— 轻量级流式数据结构 │
│ ┌─────────────────────────────────────────────────────┐ │
│ │ • 核心抽象:Stream + Consumer Group │ │
│ │ • 设计目标:Redis 生态内的轻量级 MQ │ │
│ │ • 独特优势: │ │
│ │ - 零部署成本(你已经有 Redis 了) │ │
│ │ - 微秒级延迟 │ │
│ │ - 和 Redis 其他数据结构无缝配合 │ │
│ │ • 主要局限: │ │
│ │ - 消息量受内存限制 │ │
│ │ - 持久化依赖 RDB/AOF(不保证不丢) │ │
│ │ - 没有死信队列、延迟消息等高级功能 │ │
│ └─────────────────────────────────────────────────────┘ │
│ │
└─────────────────────────────────────────────────────────────┘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
60
61
62
63
64
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
60
61
62
63
64
7.3 选型决策树
┌─────────────────────────────────────────────────────────────┐
│ MQ 选型决策树 │
├─────────────────────────────────────────────────────────────┤
│ │
│ 开始 │
│ │ │
│ ├── 你已经有 Redis 且消息量小? │
│ │ └── 是 → Redis Streams(零部署成本) │
│ │ │
│ ├── 需要灵活路由 / 延迟消息 / 优先级队列? │
│ │ └── 是 → RabbitMQ(经典、稳定、功能丰富) │
│ │ │
│ ├── 阿里云 / 国内环境 / 需要事务消息? │
│ │ └── 是 → RocketMQ(事务消息、中文生态) │
│ │ │
│ ├── 日活千万 / 大数据管道 / 事件溯源? │
│ │ └── 是 → Kafka(高吞吐、持久化、流处理) │
│ │ │
│ └── 不确定? │
│ └── 中小项目 → RabbitMQ(简单、够用) │
│ 大流量项目 → Kafka(扩容能力强) │
│ │
└─────────────────────────────────────────────────────────────┘1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
核心总结
总结1:MQ 的三大核心价值
异步解耦:Producer 和 Consumer 不直接依赖
削峰填谷:MQ 缓冲峰值流量,Consumer 匀速消费
最终一致性:通过消息驱动,各子系统最终达到一致状态1
2
3
2
3
总结2:消息可靠性四层保障
生产者确认 → Broker 持久化 → 消费者确认 → 死信队列
每一层都不可或缺,共同保证消息不丢1
2
2
总结3:消费语义务实选择
at-most-once: 丢了无所谓(监控日志)
at-least-once + 幂等: 绝大多数业务场景的最佳选择
exactly-once: 复杂且性能开销大,非必要不用1
2
3
2
3
总结4:顺序性
全局有序几乎不可能 → 分区有序(Key 哈希)
同一实体(用户/订单)的操作进同一分区 → 单线程消费 → 有序1
2
2
总结5:选型一句话
大数据管道 / 事件溯源 → Kafka
灵活路由 / 延迟消息 / 中小业务 → RabbitMQ
阿里云 / 事务消息 → RocketMQ
已有 Redis / 轻量场景 → Redis Streams1
2
3
4
2
3
4
章节测试
测试1:MQ 解决了同步调用的什么问题?
请至少列出三个。
测试2:Consumer Group 的核心规则是什么?
测试3:Kafka 中 Partition 数量与 Consumer 数量(同组)的关系是什么?
测试4:消息队列的三种消费语义分别是什么?各适用什么场景?
测试5:如何实现"实际上的 exactly-once"而不使用 Kafka 事务?
测试6:为什么消息会乱序?如何保证同一用户的操作有序?
测试7:以下场景分别适合哪种 MQ?
A. 收集千万级设备的日志,做实时分析 B. 电商系统的订单延迟关单(下单 30 分钟未支付自动取消) C. 你已经有 Redis 集群,需要做一个简单的异步任务队列 D. 需要事务消息保证"扣库存"和"发消息"的原子性
参考答案
测试1答案
答案:
- 异步解耦:Producer 不直接依赖 Consumer,下游挂了不影响上游
- 削峰填谷:MQ 缓冲峰值流量,下游按自身处理能力匀速消费
- 最终一致性:消息驱动各子系统异步处理,最终达到一致状态
- 此外还有:可扩展(加 Consumer 不用改 Producer)、可恢复(Consumer 挂了消息还在 MQ 里)
测试2答案
答案:
- 同一个 Consumer Group 内的 Consumer 协同消费 Topic
- 每个 Partition 只能被同组内一个 Consumer 消费
- 不同 Consumer Group 之间独立消费,互不影响
- Consumer 数量超过 Partition 数量时,多余的 Consumer 空闲
测试3答案
答案:
- Partition 数 >= Consumer 数时:每个 Consumer 分配到至少一个 Partition
- Partition 数 < Consumer 数时:多余的 Consumer 处于空闲状态
- 因此 Partition 数是同组内并行消费的上限
测试4答案
答案:
| 语义 | 含义 | 适用场景 |
|---|---|---|
| at-most-once | 可能丢,绝不重复 | 监控指标、日志采集 |
| at-least-once | 绝不丢,可能重复 | 大多数业务(配合幂等) |
| exactly-once | 既不丢也不重复 | 金融交易等强一致性场景 |
测试5答案
答案: 使用 at-least-once + 消费者幂等 的组合:
- 消息带唯一业务 ID
- 消费端通过数据库唯一约束 / Redis SETNX / 版本号等方式去重
- 效果等同于 exactly-once,但实现简单得多
测试6答案
答案:
- 乱序原因:同 Topic 的消息分布到多个 Partition,不同 Consumer 并发消费,处理速度不同导致乱序
- 解决方案:用业务 Key(如 userId)做哈希路由,保证同一实体的消息进入同一 Partition,由同一 Consumer 串行消费
测试7答案
答案:
- A:Kafka(高吞吐日志管道 + 流分析)
- B:RabbitMQ(延迟消息插件)或 RocketMQ(内置延迟消息)
- C:Redis Streams(零部署成本,已有 Redis)
- D:RocketMQ(内置事务消息)
相关笔记
- [[/03-web/07-middleware/02-mq-kafka]] - Kafka 核心机制深度解析
- [[/03-web/07-middleware/03-mq-others]] - 消息队列全景对比与选型
- [[/03-web/07-middleware/07-mq-advanced]] - 死信队列、延迟消息、消息积压
- [[/03-web/07-middleware/redis-guide/00-redis-overview]] - Redis 全景
- [[/03-web/06-databases-and-data-access/04-distributed-system]] - 分布式系统基础
下一步学习
- [ ] 阅读 [[/03-web/07-middleware/02-mq-kafka]] - 深入 Kafka 分区、偏移量、幂等性
- [ ] 阅读 [[03-kafka-streams-connect]] - Kafka Streams 流处理与 Connect 数据管道
- [ ] 阅读 [[04-rabbitmq-deep-dive]] - RabbitMQ 交换机和死信队列深度解析
学习状态:🟡 开始学习