Skip to content
Gains Summary
Main Navigation 首页 / Home
C++ 编程 / C++ Programming
系统与高性能 / Systems & Performance
Web 开发 / Web Development
人工智能 / Artificial Intelligence
工业软件 / Industrial Software
其他内容 / Other Topics
C++ 编程 / C++系统与性能 / SystemsWeb 开发 / Web人工智能 / AI工业软件 / Industrial

外观

Sidebar Navigation

← Web 开发 / Web Development

中间件 / Middleware

1. Web 中间件全景 / The Web Middleware Landscape

2. 中间件——流量洪峰下的系统保护 / Middleware for Protecting Systems Under Traffic Spikes

cache layer

1. 缓存架构全景 / Cache Architecture Overview

2. Redis 深入——为什么你的缓存总是出问题 / Redis Internals and Cache Failure Modes

3. Redis 数据结构场景应用——什么时候用什么 / Choosing Redis Data Structures for Real Applications

4. 缓存策略——库存变了,缓存怎么处理 / Cache Strategies and Inventory Consistency

5. Redis 缓存三剑客——穿透/击穿/雪崩 + 一致性策略 / Redis Cache Penetration, Breakdown, Avalanche, and Consistency

6. Redis 分布式锁——从 SETNX 到 Redisson / Redis Distributed Locks from SETNX to Redisson

7. Redis Cluster 与 Sentinel 高可用架构 / Redis Cluster and Sentinel High-Availability Architecture

8. Redis 高级特性——Stream / PubSub / Module / LLM 应用

9. Redis 架构深度分析——为什么 Redis 能这么快 / Redis Architecture and the Sources of Its Performance

10. Redis 全景——为什么你的系统需要一个缓存层 / The Redis Landscape and Why Systems Need a Cache Layer

message queue

1. 消息队列全景——为什么你的系统需要一个中间人 / Message Queue Overview and Why Your System Needs a Middleman

2. Kafka 核心——为什么你的消息总是"丢"了 / Kafka Fundamentals and Message Delivery Semantics

3. 消息队列高级——死信队列、延迟消息、消息积压 / Advanced Messaging with Dead Letters, Delays, and Backlogs

4. Kafka Streams 与 Connect —— 让数据自己流动起来 / Kafka Streams and Connect for Streaming Data Pipelines

5. RabbitMQ 深度解析 —— 灵活路由与消息可靠性 / RabbitMQ Deep Dive into Flexible Routing and Reliability

6. 消息队列对比——为什么最终选了 Kafka / Comparing Message Queues and Choosing Kafka

search engine

1. 搜索引擎知识体系 / Search Engine Knowledge System

2. Elasticsearch——为什么 Like 查询总是那么慢 / Elasticsearch for Full-Text Search at Scale

3. Elasticsearch 查询 DSL 深入——为什么你的搜索总是不准 / Elasticsearch Query DSL Deep Dive

4. Elasticsearch 集群规划与运维——为什么你的集群总是"黄" / Elasticsearch Cluster Planning and Operations

5. Meilisearch 与轻量搜索替代方案——当 ES 太重时 / Meilisearch and Lightweight Search Alternatives

infrastructure

1. 基础设施组件全景 / Infrastructure Components Landscape

2. Nginx 与反向代理 / Nginx and Reverse Proxy

3. 服务发现 / Service Discovery

4. 配置中心 / Configuration Center

本页目录

消息队列全景——为什么你的系统需要一个中间人 / 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

这就是消息队列要解决的核心问题:请求不能丢,也不能让下游崩掉。


第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

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

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部分: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.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

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

第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

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

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

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

第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

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

第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

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

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

第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

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

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

核心总结 ​

总结1:MQ 的三大核心价值 ​

异步解耦:Producer 和 Consumer 不直接依赖
削峰填谷:MQ 缓冲峰值流量,Consumer 匀速消费
最终一致性:通过消息驱动,各子系统最终达到一致状态
1
2
3

总结2:消息可靠性四层保障 ​

生产者确认 → Broker 持久化 → 消费者确认 → 死信队列
每一层都不可或缺,共同保证消息不丢
1
2

总结3:消费语义务实选择 ​

at-most-once: 丢了无所谓(监控日志)
at-least-once + 幂等: 绝大多数业务场景的最佳选择
exactly-once: 复杂且性能开销大,非必要不用
1
2
3

总结4:顺序性 ​

全局有序几乎不可能 → 分区有序(Key 哈希)
同一实体(用户/订单)的操作进同一分区 → 单线程消费 → 有序
1
2

总结5:选型一句话 ​

大数据管道 / 事件溯源 → Kafka
灵活路由 / 延迟消息 / 中小业务 → RabbitMQ
阿里云 / 事务消息 → RocketMQ
已有 Redis / 轻量场景 → Redis Streams
1
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答案 ​

答案:

  1. 异步解耦:Producer 不直接依赖 Consumer,下游挂了不影响上游
  2. 削峰填谷:MQ 缓冲峰值流量,下游按自身处理能力匀速消费
  3. 最终一致性:消息驱动各子系统异步处理,最终达到一致状态
  4. 此外还有:可扩展(加 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 + 消费者幂等 的组合:

  1. 消息带唯一业务 ID
  2. 消费端通过数据库唯一约束 / Redis SETNX / 版本号等方式去重
  3. 效果等同于 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 交换机和死信队列深度解析

学习状态:🟡 开始学习

最后更新于:

Pager
上一篇10. Redis 全景——为什么你的系统需要一个缓存层 / The Redis Landscape and Why Systems Need a Cache Layer
下一篇2. Kafka 核心——为什么你的消息总是"丢"了 / Kafka Fundamentals and Message Delivery Semantics

持续记录,持续成长

Copyright © Tidenflow