消息队列学习
消息队列从零到实战:Kafka、RocketMQ、RabbitMQ 完整教程
面向读者:会一点 Python、FastAPI 和数据库,但从未系统接触消息队列。
学习目标:不仅“知道名词”,还能够解释原理、做技术选型、写出可靠代码,并知道系统出故障时该检查什么。
文档基准日期:2026-08-12。通用原理不依赖具体版本;涉及产品现状的内容以文末官方资料为准。
目录与学习路线
- 第 0 章:先看全局——你究竟要学会什么
- 第一部分:为什么世界上需要消息队列
- 第二部分:消息为什么会丢、会重、会乱
- 第三部分:深入理解三大主流 MQ
- 第四部分:高级场景与通用解法
- 第五部分:Python + FastAPI 完整项目
- 第六部分:生产工程、排障与能力检验
- 附录
第 0 章:先看全局——你究竟要学会什么
0.1 一句话总纲
消息队列是一个位于生产者和消费者之间的、能够暂存并转交消息的中间系统。它让“发送方现在产生工作”和“接收方稍后完成工作”不必在同一时间、同一个进程、同一种速度下发生。
这句话包含了整本教程的主线:
- 时间解耦:发送方不必等接收方做完,所以可以异步。
- 速度解耦:流量来得快、处理得慢时,消息先积压,所以可以削峰填谷。
- 空间解耦:发送方不必知道每个接收方的地址和实现,所以可以降低系统耦合。
- 代价:一旦“当场调用”变成“以后处理”,你就必须面对消息丢失、重复、乱序、延迟、积压、数据不一致和排障困难。
因此,真正学会 MQ 不是会写一个 send(),而是能回答:
- 消息发送成功到底表示什么?到客户端缓冲区、到 Broker 内存、写入磁盘,还是完成多副本复制?
- 消费成功到底表示什么?拿到了消息,还是业务数据库已经成功提交?
- 消费者处理完但来不及确认就宕机,会发生什么?
- 同一条消息来了两次,业务是否仍正确?
- 订单“创建→付款→发货”为什么有时会乱序,应该保证全局顺序还是同一订单的局部顺序?
- 数据库提交成功但消息发送失败,如何保证最终还能发出去?
- 队列积压一百万条时,应该扩消费者、限生产者、丢消息,还是降级?
0.2 一张知识地图
flowchart LR
A["同步调用的痛点"] --> B["引入 Broker"]
B --> C["异步处理"]
B --> D["削峰填谷"]
B --> E["服务解耦"]
C --> F["确认与重试"]
D --> G["积压与背压"]
E --> H["事件契约与版本"]
F --> I["至少一次 + 幂等"]
G --> J["分区与消费组"]
H --> K["事务消息 / Outbox"]
I --> L["顺序、延迟、死信"]
J --> M["Kafka / RocketMQ / RabbitMQ"]
K --> M
L --> M
M --> N["FastAPI 可靠消息项目"]
0.3 贯穿全书的案例
假设你正在开发一个电商下单接口:
1 | |
我们会不断改造这个案例。每一章都不是孤立的:前一章暴露一个问题,后一章给出机制,再后一章讨论机制带来的新问题。
0.4 学习时牢记的三个层次
遇到任何 MQ 概念,按三个层次理解:
- 业务层:它解决什么业务问题?例如下单接口不应等待短信发送。
- 语义层:失败时允许丢、允许重、允许乱吗?一致性要到什么程度?
- 机制层:具体通过确认、重试、分区、副本、事务、幂等表等什么机制实现?
只背机制容易变成“会配参数但不会设计”;只谈业务容易变成“知道要可靠但不知道怎么做”。后续所有章节都同时回答这三层。
第一部分:为什么世界上需要消息队列
第 1 章:从一个同步下单接口开始
1.1 最直觉的写法
一个初学者很自然会写出:
1 | |
这叫同步编排:调用者沿调用链等待每一步完成。注意,“同步”在这里描述业务等待关系,不等于 Python 代码是否使用 async def。即便用了 await,用户仍在等待所有下游完成。
假设耗时如下:
| 操作 | 耗时 |
|---|---|
| 创建订单 | 40 ms |
| 扣库存 | 80 ms |
| 增加积分 | 70 ms |
| 发短信 | 300 ms |
| 建物流任务 | 100 ms |
| 写分析系统 | 120 ms |
| 合计(忽略并行) | 710 ms |
用户真正关心的通常是“订单是否创建成功”,却为短信、积分、分析等非核心操作多等了几百毫秒。
1.2 同步调用的四个问题
问题一:响应慢
调用链越长,总耗时越大。串行时近似:
[
T{response}=T{order}+T{inventory}+T{points}+T_{sms}+\cdots
]
即使并行调用,总耗时也受最慢下游限制,还要处理并发异常。
问题二:可用性被“串联”
如果任一依赖失败就让下单失败,粗略地说,调用链总成功率是各环节成功率的乘积。六个环节各自可用性都是 99.9%,整体约为:
[
0.999^6 \approx 99.4\%
]
一个暂时不可用的短信服务,不应该阻止订单创建;但同步强依赖让它获得了“否决权”。
问题三:突发流量沿链路扩散
平时每秒 100 个订单,促销瞬间每秒 10,000 个。订单服务、积分服务、短信服务和分析数据库同时承受 100 倍压力。最弱的下游先崩,再通过超时和重试拖垮上游,形成级联故障。
问题四:改动牵一发动全身
新增“风控服务”时必须修改订单服务;短信接口变化也可能迫使订单服务重新发布。订单服务知道太多下游细节,这叫高耦合。
1.3 不是所有调用都应该异步
先划清边界:
- 必须立即拿到结果才能继续决策:检查登录、实时查询价格、支付授权,适合同步调用。
- 可以稍后完成:短信、邮件、积分、分析、搜索索引更新,适合异步消息。
- 结果虽重要但允许最终一致:支付成功后的记账、库存同步,需要可靠异步和一致性设计,而不是“随便异步”。
判断问题可以归纳为一句:当前请求是否必须依赖这个下游结果,才能向用户给出正确响应?
1.4 初步改造
1 | |
订单服务把“订单已创建”写成一条消息,Broker 先接住并存储;短信服务按自己的速度处理。至此,订单服务不再等待短信服务。
但问题也随之而来:Broker 是什么?消息放在哪里?谁来确认?Broker 宕机会不会丢?这些正是后续章节的内容。
第 2 章:消息队列到底是什么
2.1 用快递柜建立直觉
把 MQ 想象成小区快递柜:
- 商家是 Producer(生产者),负责寄出包裹。
- 包裹是 Message(消息)。
- 快递柜及其运营系统是 Broker(消息代理服务器)。
- 柜子的分类区域类似 Queue/Topic(队列/主题)。
- 取件人是 Consumer(消费者)。
- 取件码和取件记录类似消息 ID、确认和消费位点。
商家把包裹放进柜子后,不需要等住户回家;住户晚些时候自行取件。商家和住户不必同时在线,也不需要直接见面。这就是时间和空间解耦。
这个比喻也有边界:真实 MQ 会复制数据、批量传输、重试投递、保留历史、组织消费组;真实消息通常是小型业务事件,不应塞几百 MB 文件。
2.2 严格定义
消息队列通常是一类消息中间件,它提供:
- 接收生产者发送的消息;
- 按一定规则路由、持久化或缓存;
- 将消息交给一个或多个消费者;
- 记录消息是否被可靠处理,失败时重投;
- 在集群中通过副本、选主等机制提高可用性。
“队列”是通俗总称。Kafka 更准确地称为分布式事件流平台,RabbitMQ 擅长消息路由和任务队列,RocketMQ 介于业务消息平台和事件流之间。三者都能承担“生产者—Broker—消费者”的中介角色,但数据模型和优势不同。
2.3 最基本的角色
flowchart LR
P1["生产者 A"] --> B["Broker 集群"]
P2["生产者 B"] --> B
B --> C1["消费者组:库存"]
B --> C2["消费者组:积分"]
B --> C3["消费者组:短信"]
Producer
产生消息的应用。它需要决定:发到哪里、消息键是什么、失败是否重试、要等到何种确认。
Broker
接收、存储、路由和投递消息的服务器。生产环境通常是集群而不是单机。
Message
消息一般包括:
1 | |
- 元数据回答“这是谁、何时发生、如何追踪、按什么路由”。
payload回答“业务事实是什么”。
Queue / Topic
- Queue 更像待办任务集合,消息常由竞争消费者中的一个处理。
- Topic 更像事件分类或广播频道,同一事件可以被多个逻辑订阅方分别处理。
不要把名称当成绝对边界:RabbitMQ 可以用 exchange + 多队列做发布订阅,Kafka 可以用同一个 consumer group 做竞争消费。
Consumer
读取并处理消息的应用。可靠消费者必须决定何时确认、如何重试、如何幂等、失败多久进入死信。
Consumer Group
具有同一业务目的的一组消费者实例。组内通常分摊工作,组与组之间各自收到一份逻辑消息。
例如 order.created:
inventory-group有 4 个实例,共同分摊扣库存;points-group有 2 个实例,共同分摊加积分;- 两个组互不替代,所以库存和积分都能处理每笔订单。
2.4 消息是“事实”、命令还是任务
这三类很容易混淆:
| 类型 | 示例 | 含义 | 命名建议 |
|---|---|---|---|
| 事件 Event | order.created |
某个事实已经发生 | 过去式 |
| 命令 Command | inventory.reserve |
希望某个明确接收方做事 | 动词原形 |
| 任务 Task | image.resize |
可由任一 worker 完成的工作 | 动作名 |
事件发布者不应命令所有订阅者如何行动。订单服务只宣布“订单已创建”,积分、通知、分析系统自行决定是否响应。这比发布 send_sms_and_add_points 更解耦。
2.5 MQ 与内存队列、HTTP、任务框架的区别
与 Python asyncio.Queue 的区别
asyncio.Queue 只在单进程内存中协作:进程退出数据通常消失,不能天然跨机器,也没有集群复制、持久化、消费组和管理界面。它适合进程内生产者—消费者,不是分布式 Broker。
与 HTTP/RPC 的区别
HTTP 更像打电话:双方通常需要同时在线,调用方等待返回。MQ 更像发邮件:收件方可以稍后处理。HTTP 适合即时问答,MQ 适合事件通知和异步任务;两者常在一个系统中并存。
与 Celery 的区别
Celery 是分布式任务框架,提供任务定义、worker、重试、定时等上层抽象;它通常还需要 RabbitMQ 或 Redis 作为 Broker。RabbitMQ 是消息基础设施,Celery 是基于消息基础设施构建的任务执行体验。学习底层 MQ 能让你理解 Celery 为什么也会重复执行、积压和需要幂等。
第 3 章:异步、削峰与解耦
这是 MQ 最常被概括的三个作用,但背口诀远远不够。我们逐个看“没有 MQ 时怎样、加入后怎样、代价是什么”。
3.1 异步处理
改造前
1 | |
改造后
1 | |
响应可能降到 60ms。关键不是“代码并发了”,而是用户请求的完成条件改变了:从“所有副作用完成”改为“核心事务完成且后续工作已被可靠接管”。
异步的代价
- 用户立即查询时,积分可能还没增加;系统从强一致转向最终一致。
- 错误不再直接返回给原请求,必须通过日志、监控、补偿任务处理。
- 调试从一条同步调用栈变成跨进程的消息链路。
- 如果消息只“尽力发送”却未可靠落地,所谓异步可能只是把错误藏起来。
FastAPI BackgroundTasks 是 MQ 吗
不是。它只是在当前 Web 进程返回响应后执行函数:
1 | |
进程崩溃、容器重启时任务可能消失;它没有独立持久化、跨实例协调、重试和死信。适合不重要的本地轻任务,不适合扣款、发货等关键业务。
3.2 削峰填谷
水库模型
把流量想成河水:
- 生产速率 (P):每秒进入多少条消息;
- 消费能力 (C):每秒最多处理多少条;
- 积压量 (B):水库中已有多少消息。
在一段时间 (t) 内,如果 (P>C):
[
B{new}=B{old}+(P-C)\times t
]
促销 60 秒内每秒产生 10,000 个订单,而库存系统每秒处理 2,000 个:
[
(10000-2000)\times60=480000
]
会积压 48 万条。流量恢复为每秒 500 条后,库存系统可用每秒 1,500 条净能力消化积压,理论恢复约需:
[
480000/1500=320\text{ 秒}
]
MQ 做了什么
MQ 不是凭空提升库存数据库的速度,而是把尖峰摊平:Broker 在高峰先存住消息,消费者按可承受速率持续处理。这叫削峰填谷。
“削峰”不等于“限流”
- 削峰:允许请求进入,但将后续工作缓冲,代价是延迟和积压。
- 限流:超过系统预算时拒绝、排队等待或降级,保护入口和资源。
如果生产速度长期大于消费速度,队列只会成为越来越满的水库,最终磁盘耗尽。正确系统通常同时具备:入口限流、队列缓冲、消费者扩容、超载降级和容量预案。
3.3 降低系统耦合性
同步耦合
订单服务直接调用短信、积分、物流、推荐:
1 | |
新增订阅者必须改订单服务。
事件解耦
订单服务只发布稳定业务事实 order.created.v1。后来新增数据仓库消费者,不需要改订单服务。
但请注意:MQ 转移并降低了运行时耦合,没有消灭契约耦合。消费者仍依赖事件字段含义。随意删除 order_id 或改变金额单位,仍会破坏下游。因此第 9 章要专门讲 Schema 和版本兼容。
3.4 其他常见作用
广播与一对多通知
一个订单事件可被库存、积分、通知、分析四个逻辑订阅方独立消费。生产者无需调用四次。
失败隔离
短信服务宕机时,订单仍可创建,消息暂存在 Broker。恢复后继续处理。但前提是 Broker 容量足够、保留期足够,并且有积压告警。
数据复制和事件驱动
业务事实进入事件流后,可用于搜索索引、缓存更新、审计、实时大屏、机器学习特征等多个下游。
可回放
Kafka 等日志型系统可在保留期内重置 offset,重新计算数据。传统任务队列通常确认后删除或不可方便回放。这是 Kafka 在数据平台场景中的关键优势。
3.5 何时不该上 MQ
- 单体小系统、流量低、只有一个简单下游,同步调用已足够。
- 业务必须即时返回下游结果。
- 团队没有能力维护 Broker、监控积压和处理重复消息。
- 只是为了“显得微服务化”。
MQ 会引入基础设施、最终一致性、重复消费、链路追踪和运维成本。正确问题不是“MQ 好不好”,而是“异步、缓冲、广播、回放的收益是否大于复杂度”。
第 4 章:消息模型与拓扑
上一章说明了为什么需要 MQ。本章把业务意图翻译成拓扑:一条消息到底交给一个消费者还是多个?如何按条件路由?
4.1 点对点 / 工作队列
flowchart LR
P["生产者"] --> Q["任务队列"]
Q --> C1["Worker 1"]
Q --> C2["Worker 2"]
Q --> C3["Worker 3"]
同一任务由某一个 worker 处理。多个 worker 是为了分摊负载,而不是每个都执行一次。适合图片压缩、邮件发送、报表生成。
4.2 发布订阅
flowchart LR
P["订单服务"] --> T["order.created"]
T --> Q1["库存订阅"]
T --> Q2["积分订阅"]
T --> Q3["分析订阅"]
同一事件被多个逻辑订阅者各自处理。Kafka 中通常对应不同 consumer group;RabbitMQ 中通常是 exchange 路由到多个 queue;RocketMQ 中通常是不同 consumer group 订阅同一 topic。
4.3 RabbitMQ 式路由直觉
RabbitMQ 的生产者一般不是直接把消息塞入某个队列,而是发给 Exchange(交换机),交换机依据绑定规则路由到队列:
1 | |
主要类型:
| Exchange | 匹配方式 | 示例 |
|---|---|---|
| direct | routing key 精确匹配 | order.created |
| fanout | 忽略 key,广播到所有绑定队列 | 配置刷新通知 |
| topic | 单词模式匹配,* 一个词、# 多个词 |
order.*、order.# |
| headers | 按消息头匹配 | 多属性路由,较少使用 |
例如 routing key 是 order.cn.created:
order.*.created能匹配;order.#能匹配;payment.*不能匹配。
4.4 Kafka / RocketMQ 式主题直觉
生产者将消息写入 topic;topic 内部再划分为 partition 或 message queue,用于并行和顺序。消费者组维护读取进度。
1 | |
消息通常按 key(如 order_id)映射到某个分区。同一 key 保持在同一分区,才有机会保证同一订单的顺序。
4.5 队列模型和日志模型
这是理解 RabbitMQ 与 Kafka 差异的核心。
队列模型
消息等待被领取;确认完成后,Broker 可以删除它。关注“待处理工作有哪些”。RabbitMQ 经典队列/仲裁队列很符合这个直觉。
追加日志模型
消息顺序追加到日志,消费者用 offset 表示“读到哪”。消费并不会立即删除消息;是否保留由时间或容量策略决定。关注“发生过哪些事件”。Kafka 是典型日志模型。
不是非黑即白
RabbitMQ 还有 Streams,Kafka 也能承载异步任务,RocketMQ 既有消费进度和回溯,也提供丰富业务消息类型。选型应看主工作负载,不要只看产品标签。
4.6 一条消息的完整旅程
sequenceDiagram
participant App as 业务/生产者
participant P as Producer 客户端
participant B as Broker
participant C as Consumer
participant DB as 消费者数据库
App->>DB: 提交业务数据(可能与 Outbox 同事务)
App->>P: 构造并发送消息
P->>B: 网络发送
B->>B: 校验、写入、复制
B-->>P: 发送确认
B->>C: 投递或响应拉取
C->>DB: 执行业务并提交
C-->>B: ACK / 提交 offset
其中任何箭头都可能超时或失败。第 6 章将沿这条旅程逐段讨论可靠性。
第二部分:消息为什么会丢、会重、会乱
第 5 章:推模式、拉模式与长轮询
5.1 先定义“推”和“拉”
推模式 Push
Broker 主动把消息交给消费者。消费者注册回调,消息到达后回调被调用。
1 | |
优点:使用简单、到达延迟低。难点:如果消费者很慢,Broker 或客户端必须做流量控制,否则消费者会被淹没。
拉模式 Pull
消费者主动向 Broker 请求消息,并决定何时拉、拉多少。
1 | |
优点:消费者掌握节奏,便于批处理和背压;缺点是空轮询会浪费请求,客户端逻辑相对复杂。
5.2 短轮询、长轮询和真正服务器推送
“PushConsumer”这个名字不一定表示 Broker 通过永久主动推送实现。
- 短轮询:消费者不断问“有吗?有吗?”没有数据立即返回。请求多、空转严重。
- 长轮询:消费者发请求,若暂时无消息,Broker 挂起一段时间;有消息立即返回,超时才空返回。它在体验上接近推送,但底层仍是消费者发起请求。
- 真正推送:Broker 在已有连接上主动发送。需要很好的信用额度、prefetch 或流控机制。
许多 MQ 客户端所谓“推模式”其实是在 SDK 内部持续拉取,然后以回调形式“推给”用户代码。理解抽象和实现的区别,比死记产品宣传更重要。
5.3 三大 MQ 的大致风格
- Kafka:消费者主动 poll/fetch,天然是拉取模型;适合批量、吞吐和自行控制进度。
- RabbitMQ:常见用法是注册 consumer 后由 Broker 投递,结合 prefetch 控制尚未确认的消息数;也支持主动 basic.get,但不适合作为高吞吐常规消费模式。
- RocketMQ:提供 PushConsumer、SimpleConsumer、PullConsumer。PushConsumer 对业务代码呈现回调式推送,SimpleConsumer 给业务更强的拉取和确认控制;官方建议多数普通场景使用 PushConsumer 或 SimpleConsumer,PullConsumer 更偏流处理框架集成。
5.4 Prefetch:推模式的刹车
RabbitMQ 消费者设置 prefetch_count=20,可以直观理解为:在当前消费者还没有确认的消息达到 20 条后,Broker 暂停继续投递。
1 | |
经验上应结合单条处理耗时、并发数和内存测试,而不是盲抄固定值。
5.5 批量拉取:吞吐的来源之一
网络往返很贵。一次取 100 条通常比请求 100 次每次取 1 条高效。Kafka 的高吞吐很大程度来自顺序追加、页缓存、批处理、压缩和零拷贝等一整套设计,而不是单个“神奇参数”。
批处理也有代价:一批 100 条的第 37 条失败时如何提交进度?整批重试会重复前 36 条;跳过失败项又可能破坏顺序。所以批量越大,吞吐越高,失败语义越需要精心设计。
5.6 如何选择
- 任务短、希望接口简单、希望 Broker/SDK 管调度:回调式 PushConsumer。
- 需要批量、精确流控、长耗时或自己管理并发:Pull/SimpleConsumer。
- 数据流计算:通常拉取并基于 offset 管理进度。
无论选什么,最终都必须回答:“消费者什么时候表示这条消息归我负责,以及什么时候表示已经完成?”这就是下一章的确认和投递语义。
第 6 章:确认、偏移量与投递语义
6.1 先看消息丢失的三个区段
1 | |
- 发送端丢失:数据库已提交,但程序在调用 MQ 前崩溃;或者发送超时后直接放弃。
- Broker 丢失:消息只在内存或单副本,Broker 宕机;或选主时丢失未同步数据。
- 消费端丢失:消费者先确认,后处理;确认后进程崩溃,Broker 以为已完成。
可靠性不是打开一个 durable=true 就完成,而是端到端每一段都要闭环。
6.2 两种“确认”不要混淆
生产者确认
Broker 告诉生产者:“我已经按约定接管了消息。”
- RabbitMQ 叫 publisher confirms。
- Kafka 由
acks控制确认条件。 - RocketMQ 的同步发送会返回发送结果,生产者还需处理异常和不确定状态。
消费者确认
消费者告诉 Broker:“业务处理完成,可以推进消费进度或释放这条消息。”
- RabbitMQ 使用
ack/nack/reject。 - Kafka 提交 consumer offset。
- RocketMQ 消费者返回成功/失败或显式 ACK(依消费者类型而异)。
6.3 ACK 应该放在哪里
错误方式:
1 | |
正确原则:先让业务副作用成功提交,再 ACK。
1 | |
但第二段也有窗口:数据库提交成功,ACK 前进程宕机,Broker 会重新投递,于是业务执行两次。由此得到消息系统最重要的工程结论之一:
可靠消费通常选择“至少一次”,然后让消费者幂等;不要幻想普通网络系统天然只执行一次。
6.4 三种投递语义
至多一次 At-most-once
消息最多处理一次,可能一次也没有。
典型做法:先提交 offset/ACK,再处理;或发送失败不重试。适合可丢的指标、非关键日志。
至少一次 At-least-once
消息最终至少处理一次,但可能重复。
典型做法:处理成功后才确认;失败或超时重投。绝大多数可靠业务采用它,再配合消费幂等。
恰好一次 Exactly-once
从业务效果看,每个事件只产生一次结果。这是最容易被误解的概念。
“Broker 没有重复写入”不等于“数据库只扣一次钱”。端到端恰好一次需要把输入进度和输出结果纳入同一原子边界,或者通过幂等使重复执行的最终效果等价于一次。
Kafka 的事务可以在 Kafka 内实现 consume-process-produce:原子写多个 Kafka topic,并把输入 offset 一同提交;消费者用 read_committed 隔离级别只看已提交记录。但如果处理中还写一个普通 PostgreSQL 数据库,Kafka 事务不会自动把这个外部数据库拉进原子事务。
因此更准确的表达是:
“恰好一次”必须说明边界:Kafka 内部、消息存储层、某个数据库,还是端到端业务效果。
6.5 幂等:重复做,结果不变
数学上,幂等操作满足:
[
f(f(x))=f(x)
]
例如“把订单状态设置为 PAID”通常可幂等;“余额减 100”天然不幂等。
方案一:业务唯一约束
积分流水表增加唯一键 (event_id, action):
1 | |
同一 event_id 第二次插入触发唯一约束,消费者将它识别为“已处理”并 ACK。
方案二:Inbox / 消费记录表
在与业务更新同一个数据库事务中:
1 | |
为 (message_id, consumer_name) 建唯一键。重复消息的 INSERT 失败,整个事务不会再加积分。
方案三:状态机条件更新
1 | |
检查影响行数:第一次更新 1 行,重复处理更新 0 行。状态迁移必须合法。
Redis SETNX 的陷阱
如果先 SETNX message_id 成功,随后数据库写失败,重试会因为 key 已存在而跳过,反而丢失业务处理。除非 Redis 标记和业务数据具有可靠的原子或补偿设计,否则不要把它当通用幂等答案。数据库唯一约束通常更稳妥。
6.6 Kafka offset 到底是什么
每个 partition 是一个有序日志,offset 是记录在该分区中的位置:
1 | |
消费位置有两个相关概念:
- 当前 position:客户端已经拉到或准备读取哪里;
- committed offset:崩溃重启后从哪里恢复。
提交的通常是“下一条要读的 offset”。处理完 offset 42 后提交 43。
自动提交很方便,但可能在业务处理完成前推进进度。关键业务建议关闭自动提交,完成数据库事务后手动提交。仍然会有“数据库已提交、offset 未提交”的重复窗口,所以照样需要幂等。
6.7 发送超时为什么很棘手
生产者请求超时可能有两种真实情况:
- Broker 没收到,消息不存在;
- Broker 已写入,但确认响应在网络上丢了。
生产者只看到“超时”,无法判断是哪一种。若重试,第二种情况可能重复;若不重试,第一种情况会丢。因此需要:
- 生产者幂等机制(如 Kafka idempotent producer);
- 全局业务
message_id; - 消费端幂等;
- 对关键消息使用 Outbox,直到明确发布成功才标记。
6.8 可靠性闭环
一条重要消息至少要做到:
- 业务数据和“待发送消息”原子落库;
- 发布器持续重试,直到 Broker 确认;
- Broker 持久化并按需要复制;
- 消费者完成本地事务后确认;
- 消费者按
message_id幂等; - 暂时性失败重试,永久失败进死信;
- 对积压、失败率、死信量告警;
- 有人工或自动补偿、重放工具。
这比一句“我们用至少一次”具体得多。
第 7 章:分区、消费组、并发与顺序
7.1 为什么需要分区
一条单队列只能由一个线程串行处理,顺序简单但吞吐受限。将 topic 切成多个分区后可以并行:
1 | |
如果增加到 4 个消费者,每个大致分到一个分区;增加到 6 个,其中 2 个通常空闲。因此 Kafka 中一个 consumer group 的有效并行上限通常受分区数限制。
RabbitMQ 的一个队列可向多个消费者分发,模型不同,但也要用 prefetch、消费者并发数和队列数量控制吞吐与隔离。
7.2 分区键决定局部顺序和负载
假设按 order_id 作为 key:
1 | |
同一订单的创建、支付、取消进入同一分区,可以保持该订单局部顺序;不同订单落到不同分区并行处理。
错误 key 会产生热点。例如以 country=CN 为 key,绝大部分流量进入同一分区。好 key 应同时满足:
- 需要顺序的数据有相同 key;
- key 分布足够均匀;
- key 长期稳定;
- 不含敏感信息,或经过安全处理。
7.3 全局顺序为什么昂贵
全局顺序要求所有消息经过同一条串行通道,相当于把并行度降为 1。吞吐、可用性和扩展能力都受限。
多数业务并不需要“所有用户所有订单全局有序”,只需要“同一订单的事件有序”。这叫按业务键的局部顺序,是工程上更合理的目标。
7.4 消费组:负载均衡与广播的统一
Kafka 的一个简洁理解:
- 实例使用同一个 group id:组内竞争消费,分摊分区;
- 使用不同 group id:每个组各自维护 offset,每组都能读到全部 topic 数据。
1 | |
RocketMQ 也有集群消费与广播消费等模式。RabbitMQ 常通过“一个业务订阅对应一个队列,同队列多个 consumer 竞争”表达相似关系。
7.5 Rebalance:扩缩容的暂停与风险
Kafka consumer 加入、离开、超时,或分区数变化时,组需要重新分配分区,这叫 rebalance。期间可能短暂停顿;批处理尚未提交时还可能重复处理或提交失败。
设计建议:
- 不要在 poll 循环里做无限长阻塞任务;
- 控制单批大小和最长处理时间;
- 使用 rebalance listener 在分区被收回前处理或提交必要进度;
- 消费逻辑幂等;
- 避免实例频繁抖动;
- 监控 rebalance 次数和 consumer lag。
7.6 并发处理如何破坏顺序
即使 Broker 按 1、2、3 投递,你把它们扔进线程池后也可能按 2、3、1 完成。要保持同一 key 的顺序,处理和确认链路也必须有序。
常见方案:
- 同一分区由单线程串行处理;
- 按 key 分派到固定 worker(keyed executor);
- 每个 key 使用状态机和版本号,拒绝非法旧事件;
- 如果允许并发,等连续 offset 都完成后再提交“最高连续位点”,不能跨过未完成消息。
第 8 章:重试、死信、积压与背压
8.1 失败分为两类
暂时性失败
数据库短暂超时、下游限流、网络闪断。稍后重试有可能成功。
永久性失败
消息格式非法、必填字段缺失、业务状态不允许、目标用户永久不存在。反复重试只会浪费资源。
可靠消费者必须区分它们。不要用一个笼统的 except Exception: requeue() 无限重试。
8.2 指数退避与抖动
立即无限重试会形成重试风暴。常用退避:
[
delay_n=\min(base\times2^n,maxDelay)+jitter
]
例如 1s、2s、4s、8s、16s,最高 5min,并加入随机抖动,让大量消费者不要同时醒来冲击下游。
每条消息需要记录或能推断:重试次数、首次失败时间、最后错误类别。超过阈值进入死信。
8.3 死信队列 DLQ
死信是无法正常处理、需要隔离的消息。进入死信的常见原因:
- 达到最大重试次数;
- 消息过期;
- 被消费者拒绝且不重新入队;
- 队列超出限制或按策略转移。
DLQ 不是垃圾桶,而是异常工单箱。必须配套:
- 死信量告警;
- 查看原消息、headers、错误、重试次数的工具;
- 修复数据或代码后的重放能力;
- 重放限速和再次失败的保护;
- 敏感数据权限与保留策略。
8.4 毒消息与头部阻塞
一条每次都会失败的消息叫 poison message。若要求严格顺序,它会挡住后面的正常消息,形成头部阻塞。
需要在“顺序”和“可用性”间做业务决策:
- 严格顺序不可破坏:暂停该 key/分区,告警并人工处理;
- 可牺牲个别异常:有限重试后送 DLQ,让后续继续;
- 用业务版本校验和补偿任务修复跳过造成的状态缺口。
8.5 消息积压
积压是生产速度在一段时间内大于消费速度。原因可能是:
- 流量突增;
- 下游数据库变慢;
- 消费者异常退出或频繁 rebalance;
- 单条消息处理变慢;
- 分区/队列热点;
- poison message 重试风暴;
- 消费端部署了错误版本。
排查顺序:
- 看生产速率、消费速率、积压量和最老消息年龄;
- 看消费者存活数、错误率、处理耗时;
- 看分区是否倾斜;
- 看下游数据库、缓存、外部 API;
- 决定扩容、限流、降级、批处理或跳过毒消息。
8.6 背压 Backpressure
背压是下游向上游表达“我处理不过来”的机制。
- RabbitMQ 用 prefetch 限制未确认消息;
- Kafka 消费者控制 poll、pause/resume 和批量;
- 应用层可降低生产速率、入口限流或暂时关闭非核心功能;
- 队列长度/磁盘水位达到阈值时,Broker 也可能阻塞或拒绝发布。
缓冲、背压和限流形成完整闭环:队列吸收短期波动,背压传播压力,限流阻止长期超载。
8.7 重试 Topic 的设计
在 Kafka 这类没有传统“消息退回队头”语义的日志系统中,常见做法是:
1 | |
也可使用调度服务、时间轮、数据库任务表等实现。重点不是名称,而是避免当前分区被单条失败消息无限卡住,并保留重试上下文。
第 9 章:消息契约、序列化与主题设计
9.1 消息不是随便拼的 JSON
生产者和消费者通过消息契约协作。契约至少说明:
event_type的业务含义;- 字段名称、类型、单位、是否可空;
- 哪个字段是幂等键和分区键;
- 时间是 UTC 还是本地时区;
- 金额用十进制定点字符串还是最小货币单位整数;
- 哪些字段包含敏感数据;
- 兼容性与废弃策略。
9.2 推荐事件信封
1 | |
message_id 用于消息级去重;aggregate_id 常用于局部顺序;trace_id 用于观察链路;causation_id 能解释“这条消息由谁触发”。
9.3 事件通知与事件携带状态
只带 ID
1 | |
消费者再查订单服务。优点是消息小、数据来源单一;缺点是产生同步依赖、查询压力,回放旧消息时查到的可能是新状态。
携带需要的数据
1 | |
消费者可独立处理和回放,但契约更重,可能复制数据。实践中应携带消费者完成该事件所需的稳定最小快照,不要塞整个数据库对象。
9.4 向后兼容
生产者新增可选字段通常较安全;删除字段、改类型、改语义最危险。
建议:
- 消费者忽略未知字段;
- 新字段先设为可选并提供默认语义;
- 不复用旧字段表达新含义;
- 重大不兼容变化创建
v2契约或新 topic; - 使用 JSON Schema、Avro、Protobuf 等进行契约校验;
- 在 CI 中做生产者—消费者契约测试。
9.5 Topic / Queue 如何划分
过粗:所有事件都放 events,权限、保留期、容量和订阅困难。
过细:每种动作一个 topic,数量爆炸,运维困难。
常见维度:
- 按业务域:
order-events、payment-events; - 按数据敏感级别和权限隔离;
- 按保留期、吞吐和可靠性要求隔离;
- 按消息类型限制隔离(RocketMQ 5.x 的 FIFO、Delay、Transaction topic 有明确类型约束);
- 事件类型可放 header 或 payload,不一定每种事件一个 topic。
9.6 大消息怎么处理
MQ 不适合传大文件。推荐 Claim Check 模式:
- 文件写对象存储;
- 消息只包含对象 URI、校验和、大小和权限信息;
- 消费者按引用下载;
- 设置对象生命周期,避免消息还没消费文件已删除。
大消息会降低吞吐、放大重试成本、增加内存与网络压力,还可能触发 Broker/客户端消息大小限制。
9.7 安全与隐私
- 使用 TLS 保护传输;
- 生产者和消费者最小权限授权;
- 不在消息中放密码、令牌、完整银行卡信息;
- 敏感字段必要时加密或只放引用;
- 设置合理保留期和删除策略;
- 日志不要完整打印敏感 payload;
- 死信和重放工具同样需要审计权限。
至此,你已理解所有 MQ 共有的骨架。下一部分进入三种产品,观察它们如何以不同数据模型实现这些能力。
第三部分:深入理解三大主流 MQ
第 10 章:RabbitMQ
10.1 它最像什么
RabbitMQ 最像一座智能分拣中心:生产者把消息交给 exchange,exchange 根据 routing key 和 binding 把消息复制或路由到一个或多个 queue,消费者再从 queue 领取任务。
它尤其适合:
- 业务任务队列;
- 灵活路由和发布订阅;
- 低到中等延迟的在线业务;
- 每条消息需要明确 ACK、重试、死信;
- Python 项目希望有成熟 AMQP 客户端。
10.2 核心拓扑
flowchart LR
P["Publisher"] --> E["Exchange: domain.events"]
E -->|"order.created"| Q1["inventory.queue"]
E -->|"order.*"| Q2["analytics.queue"]
E -->|"#"| Q3["audit.queue"]
Q1 --> C1["库存消费者"]
Q2 --> C2["分析消费者"]
Q3 --> C3["审计消费者"]
重要事实:
- exchange 负责路由,本身通常不是保存消息的队列;
- queue 才保存等待消费的消息;
- binding 把 exchange 和 queue 连接起来;
- 一条消息可以路由到多个 queue,形成多个独立副本;
- 同一个 queue 上的多个 consumer 通常竞争消费。
10.3 四种 Exchange 应该会解释
Direct
binding key 与 routing key 精确相等才路由。适合明确事件名或任务类别。
1 | |
Fanout
忽略 routing key,发给所有绑定队列。适合配置广播、临时通知、日志广播。
Topic
按点分词并支持通配符:* 匹配一个词,# 匹配零个或多个词。
1 | |
Headers
按 headers 组合匹配,适合 routing key 难以表达的多属性路由,但复杂度和使用频率通常低于 topic。
10.4 Durable、Persistent 和 Confirm 是三件事
这是 RabbitMQ 初学者最常混淆的地方。
- durable exchange/queue:Broker 重启后拓扑定义仍存在。
- persistent message:消息声明为持久化,Broker 会按队列类型和策略将其持久保存。
- publisher confirm:Broker 告知发布者已经接管,发布者才知道是否要重试。
只声明 durable queue,却发送非持久消息,消息仍可能在重启中丢;只把消息设为 persistent,却不等待 confirm,发布者也无法可靠判断是否已被接管。生产可靠性需要三者和高可用队列类型一起考虑。
10.5 Consumer ACK、NACK 和 Reject
ack:成功处理,可删除或推进队列状态。nack(requeue=True):失败并重新入队,可能很快再次投递。nack(requeue=False)/reject(requeue=False):不重入原队列;若配置 DLX 则死信,否则可能丢弃。
不要无限 requeue=True,否则毒消息在消费者之间高速循环。更好的方式是带次数的延迟重试,达到上限进 DLQ。
10.6 Prefetch 与公平分发
如果一个慢消费者预取 1000 条,这些消息会处于 unacked 状态,其他空闲消费者拿不到。设置合适的 prefetch,让 Broker 只给每个消费者有限的未确认消息。
对于每个 worker 串行处理的任务,prefetch_count=1 最容易理解但吞吐未必最好;异步并发消费者可按并发上限设置,例如并发 20 时从 20 左右开始压测。最终数值取决于处理耗时、消息大小和故障恢复要求。
10.7 Classic Queue、Quorum Queue 与 Stream
Classic Queue
传统队列,适合非复制或临时类场景。历史上的 classic mirrored queues 已在 RabbitMQ 4.x 移除,不应把旧教程的镜像队列配置照搬到新部署。
Quorum Queue
基于 Raft 的持久复制 FIFO 队列,面向数据安全和可预测故障恢复。需要复制高可用业务队列时,应优先考虑它。代价是共识复制带来的资源和延迟成本,不适合大量频繁创建/删除的临时队列,也不一定适合超长积压。
生产者使用 publisher confirms,确认通常在消息达到所需复制安全条件后给出;消费者使用手动 ACK。三副本 quorum queue 能容忍的故障和可用性由多数派决定:三节点至少需要两节点形成多数。
Streams / Super Streams
RabbitMQ Streams 更像追加日志:支持保留和重复读取,适合大积压、回放和流场景;super stream 用分区扩展。若主需求已接近日志流,要把它与 Kafka 一起评估,而不是硬用传统 queue。
10.8 TTL、死信和延迟重试
RabbitMQ 可以配置消息 TTL 或队列 TTL。消息过期、被拒绝、不满足队列限制等情况可按规则进入 Dead Letter Exchange。
一种常见延迟重试拓扑:
1 | |
注意:用 TTL + DLX 做大量、超长、任意精度的定时调度并非总是最佳选择;要评估队头过期行为、资源规模和插件/版本能力。真正复杂的调度可使用独立调度服务或 RocketMQ 原生延迟消息。
10.9 RabbitMQ 中的顺序
队列在特定条件下按入队顺序投递,但以下情况都可能让业务完成顺序变化:
- 多个消费者并发;
- 消息失败后重新入队;
- 优先级队列;
- 消费者异步并发处理;
- 连接故障导致未确认消息重投。
要求严格业务顺序时,可以使用单活消费者、单队列串行、按业务 key 分多个队列,并使消费逻辑幂等。但顺序要求越强,并行吞吐越低。
10.10 RabbitMQ 的 Python 客户端
pika:官方教程常用,经典且容易理解;BlockingConnection 适合教学或同步 worker。aio-pika:基于 asyncio,提供 robust connection/channel 等抽象,更适合 FastAPI 异步生态。- Celery / Dramatiq 等:更上层的任务框架,可把 RabbitMQ 当 Broker;便利但会隐藏部分底层语义。
FastAPI 中要复用长连接与 channel,不要每个 HTTP 请求创建一个 TCP 连接。连接生命周期应在应用启动时创建、关闭时释放。
10.11 RabbitMQ 的优势与边界
优势:
- exchange 路由灵活;
- ACK、NACK、prefetch、TTL、DLX 等任务语义自然;
- Python 生态成熟;
- 业务队列上手快;
- 管理界面直观。
边界:
- 传统队列不是为超长事件历史和大规模任意回放而生;
- 大规模数据流吞吐场景通常 Kafka 更自然;
- 大量队列、复杂拓扑、长积压需要认真做容量设计;
- 延迟消息和事务消息不像 RocketMQ 那样以核心业务消息类型呈现。
第 11 章:Kafka
11.1 Kafka 不是“消费后删除的队列”
Kafka 的核心抽象是持久的、分区的追加日志。事件写入后不会因为某个消费者读过就立即删除,而是根据时间/容量保留策略清理;不同 consumer group 维护自己的读取进度,因此同一数据可被多次读取和回放。
1 | |
这使 Kafka 天然适合事件流、日志采集、CDC、实时计算和数据平台。
11.2 核心架构
- Broker:存储 topic partition 并处理生产/消费请求。
- Topic:事件类别。
- Partition:可并行的有序追加日志。
- Replica:分区副本;leader 处理读写,followers 复制。
- Producer:选择分区、批量、压缩并发送记录。
- Consumer Group:组内分摊 partition,各组独立读取。
- Controller quorum:管理集群元数据和控制面。现代 Kafka 使用 KRaft;Kafka 4.0 起已进入无 ZooKeeper 的时代,旧教程中的
--zookeeper管理方式不应继续照搬。
生产关键集群通常让 broker 和 controller 使用独立角色,并部署奇数个 controller;开发环境可合并角色。
11.3 为什么 Kafka 吞吐高
不是因为“Java 快”这一条,而是组合设计:
- 追加写和分段日志,减少随机 I/O;
- 依赖操作系统 page cache;
- producer/consumer 批处理;
- gzip、snappy、lz4、zstd 等批压缩;
- 网络传输和文件发送路径优化;
- partition 水平扩展;
- 消费者按 offset 顺序读取。
代价是为了吞吐而批量等待会增加一点延迟;分区数也带来元数据、文件句柄、复制和 rebalance 成本。
11.4 Producer 的发送过程
flowchart LR
A["业务 Record"] --> S["序列化"]
S --> P["分区器:key→partition"]
P --> B["按 partition 批缓存"]
B --> C["压缩与网络发送"]
C --> L["Partition Leader"]
L --> R["Follower 复制"]
R --> ACK["按 acks 条件确认"]
关键参数的思维方式:
acks=0:不等 Broker,延迟低但可能静默丢失。acks=1:leader 写入即确认;leader 未复制就故障时存在丢失窗口。acks=all:等所有当前 ISR 按规则确认,并与min.insync.replicas配合提高耐久性。enable.idempotence=true:通过 producer id 和序列号对协议级重试去重;现代官方客户端默认通常开启,但第三方 Python 客户端默认值和版本需单独检查。linger.ms:多等一点以凑批,提高吞吐但增加等待延迟。batch.size、压缩类型:影响吞吐、CPU、网络。
常见生产耐久配置思路是 replication factor 3、min.insync.replicas=2、producer acks=all,但它不是任何规模的万能常量,必须结合故障容忍与成本。
11.5 ISR 和副本语义
ISR 是与 leader 保持足够同步的副本集合。acks=all 不是等待所有配置的副本,而是依据 ISR 和 min.insync.replicas 的条件确认。
假设副本因子 3,min.insync.replicas=2:
- 三副本健康:写入满足要求后确认;
- 只剩两个同步副本:仍可写;
- 只剩一个同步副本:生产请求失败,牺牲可用性以保护数据安全。
这是 CAP/一致性取舍在具体参数中的表现。
11.6 Consumer、Poll 与 Offset
Kafka 消费者批量拉取。一个 group 内,同一 partition 在一个时刻通常只分配给一个 consumer 实例,因此分区内顺序容易推理。
关键参数与行为:
group.id:逻辑订阅者身份;乱改会产生新的消费进度。enable.auto.commit:自动还是手动提交;关键业务通常手动。auto.offset.reset=earliest/latest/...:只有没有有效已提交 offset 时才决定从哪里开始,不是每次启动都生效。- poll 间隔与 session/heartbeat:处理太久可能被认为失活,引发 rebalance。
max.poll.records:控制一批交给应用多少条。
11.7 Retention 与 Log Compaction
Delete retention
按时间或容量删除旧 segment。适合事件历史,例如保留 7 天订单事件。
Log compaction
按 key 保留较新的值,旧值最终被清理。适合构建可恢复状态,例如每个用户的最新配置。压缩是后台渐进过程,不保证 topic 任意时刻只有一个 key 的一个版本;消费者仍应能处理重复和历史更新。
墓碑记录(key 有值、value 为 null)常用于表达删除,并在保留规则后被清理。
11.8 Kafka 的事务与恰好一次边界
幂等生产者解决单个 producer 会话中协议重试造成的重复写。事务生产者进一步能原子写多个 partition/topic,并可把消费 offset 纳入同一 Kafka 事务。
1 | |
下游配置 isolation.level=read_committed,过滤未提交和已中止事务。Kafka Streams 提供更完整的 exactly-once processing 支持。
再次强调:写 PostgreSQL、调用支付网关、发短信不自动属于 Kafka 事务。跨外部系统仍需 Outbox、幂等、状态机或协调方案。
11.9 Kafka 如何做顺序
Kafka 只保证一个 partition 内的记录顺序,不保证跨 partition 全局顺序。同一业务 key 必须稳定映射到同一 partition;消费者也要按分区顺序处理。
生产者关闭幂等且允许多个 in-flight 请求时,失败批次重试可能造成重排;启用幂等并使用兼容配置可维持发送顺序。分区数变化还可能改变 hash(key) % N 的映射,因此扩分区后同一 key 的新旧事件可能位于不同分区,迁移设计要谨慎。
11.10 Kafka 的延迟消息
Kafka 核心没有像 RocketMQ 那样的通用原生任意延迟消息类型。常见方案:
- retry topics 按延迟等级分层;
- 独立调度服务将到期任务写回业务 topic;
- Kafka Streams 用时间窗口/状态存储实现;
- 短暂停可在消费者侧 pause,但不适合海量长延迟占用内存。
因此“复杂订单超时调度”若是核心诉求,RocketMQ 或专门调度系统更自然。
11.11 Kafka 的 Python 客户端
confluent-kafka:基于 librdkafka,性能和协议能力强,生产中常见。aiokafka:asyncio 风格,与 FastAPI 适配自然,支持消费者组、幂等生产和事务等能力;具体默认参数不能照搬 Java 客户端。kafka-python:纯 Python 生态中的常见客户端,选择前应核对维护状态、版本兼容和所需特性。
Python 做 Kafka 数据管道没问题;但 Kafka Streams 是 Java/Scala 库。Python 的复杂流处理通常使用 Flink/PyFlink、Bytewax、Faust 生态或独立计算框架,需按成熟度评估。
11.12 Kafka 的优势与边界
优势:
- 高吞吐和水平扩展;
- 事件保留、回放;
- 消费组与分区模型适合数据流;
- Kafka Connect、CDC、流处理生态强;
- 多订阅方读取同一事件历史成本自然。
边界:
- 复杂按条件路由不如 RabbitMQ exchange 直观;
- 单消息优先级、传统任务队列语义不是核心强项;
- 原生通用延时/死信需要应用架构补充;
- 分区规划、rebalance、offset 和集群运维有学习成本;
- 低流量简单任务使用 Kafka 可能过重。
第 12 章:RocketMQ
12.1 它最像什么
RocketMQ 是面向大规模业务消息的分布式消息与事件流平台。它既有 topic、consumer group、消费进度和回溯,又把 FIFO、延迟、事务、重试/死信、过滤等业务能力做成突出的一等功能。
它尤其适合:
- 电商、支付、交易等业务消息;
- 原生事务消息;
- 按业务 key 的 FIFO 顺序;
- 订单超时、定时触发等延迟消息;
- Java/云原生体系中的大规模消息平台。
12.2 5.x 核心角色
- Producer / Consumer:消息生产和消费客户端。
- NameServer:轻量路由发现,客户端从中获取 topic 路由。
- Broker:存储和传输消息。
- Controller(特定高可用模式):参与副本和主从切换管理。
- Proxy:5.x 云原生接入层,gRPC 客户端通过 Proxy 访问 Broker,便于多语言和无状态扩展。
- Topic / MessageQueue / ConsumerGroup:组织消息、并行和消费进度。
新系统使用 5.x gRPC SDK 时,需要服务端 5.0+ 并启用 Proxy;旧 Remoting SDK 和新 gRPC SDK API 不兼容,选型时先明确协议系列。
12.3 存储直觉:CommitLog 与逻辑队列
经典 RocketMQ 存储设计可用以下模型理解:
- 消息主体顺序写入 CommitLog;
- ConsumeQueue 为 topic/queue 建轻量逻辑索引,消费者据此读取;
- IndexFile 等结构支持按 key 查询。
这样既利用顺序写提高吞吐,又让不同 topic/queue 能定位消息。具体 5.x 部署还可能启用不同存储能力,但“主体顺序存储 + 逻辑消费索引”的直觉很重要。
12.4 消息类型
Normal
普通消息,适合一般异步通知和任务。
FIFO
通过 message group 规定顺序范围。同一 message group 的消息按 FIFO 处理,不同 group 可并行。例如 message_group=order_id。
Delay
消息在指定投递时间到达前处于定时状态,到期后才对消费者可见。5.x 使用毫秒级 Unix 时间戳表达投递时间,适合订单 30 分钟未支付自动取消等场景。
Transaction
用于协调“本地事务是否成功”和“消息是否应对消费者可见”,通过半消息、二次确认和事务状态回查实现最终一致。
RocketMQ 5.x 要求 topic 的消息类型与所发送高级消息类型一致。不要把普通、FIFO、Delay、Transaction 当作一个 topic 内随意混发的标签。
12.5 FIFO 的 message group
1 | |
把 order-1 的所有消息放入同一 message group,Broker 和客户端协议维持该 group 的顺序;order-2 可在另一 group 并行。
要真正有序还需:
- 生产端按序发送;
- 消费端遵循接收—处理—返回结果的串行链路;
- 不把消息异步扔到无序线程池后提前返回成功;
- 控制失败重试,明确毒消息策略;
- 避免一个超大 message group 形成热点。
12.6 延迟消息
5.x 定时/延迟消息在设定时间前存于定时存储,时间到后进入可消费状态。适合:
- 订单超时关闭;
- 延迟重试;
- 预约通知;
- 分布式定时触发。
需要注意:
- 它保证的是“不早于指定时间可见”,不等于业务代码在那个毫秒准时完成;积压和故障会造成晚到。
- 不要让海量消息集中在同一时刻到期,避免瞬时峰值。
- 最大延迟范围在不同服务版本、部署形态和文档页面间可能不同,必须按实际版本参数限制核验,不要把旧版固定 18 个延迟等级和 5.x 时间戳模型混为一谈。
12.7 事务消息:半消息与回查
流程先建立直觉,细节在第 14 章:
sequenceDiagram
participant P as 生产者
participant B as RocketMQ Broker
participant DB as 本地数据库
participant C as 消费者
P->>B: 1. 发送半消息
B-->>P: 2. 半消息成功
P->>DB: 3. 执行本地事务
alt 本地事务成功
P->>B: 4. Commit
B->>C: 5. 消息可见并投递
else 本地事务失败
P->>B: 4. Rollback
else 状态未知/响应丢失
B->>P: 回查本地事务状态
P-->>B: Commit / Rollback / Unknown
end
它解决“本地数据库提交与消息发送之间的原子一致性”问题,达到最终一致;它不保证下游消费者的业务数据库一定一次成功,所以消费者仍需重试和幂等。
12.8 过滤、重试与死信
RocketMQ 支持 tag 和 SQL 属性过滤。Tag 适合简单稳定分类;SQL 过滤更灵活但要评估 Broker 计算成本和支持配置。
消费失败或超时会按策略重投;超过最大次数进入死信队列。与 RabbitMQ 一样,返回成功前必须完成实际业务,不能把任务异步转给自建线程后提前告诉 SDK 成功。
12.9 PushConsumer、SimpleConsumer、PullConsumer
- PushConsumer:SDK 管取消息、并发、缓存和重试,通过 listener 回调给业务;易用,处理时间应可预期,回调内同步完成。
- SimpleConsumer:业务主动接收、处理并确认,对不可预测处理时长和自定义流控更友好;需要自行管理 receive/ack 和不可见时间等细节。
- PullConsumer:控制更底层,官方更倾向把它用于流处理框架集成;普通业务优先考虑前两者。
12.10 Python 生态必须现实评估
RocketMQ 的主要历史生态长期偏 Java。5.x 通过 gRPC 统一多语言 SDK;截至本文基准日期,Apache rocketmq-clients 仓库的能力矩阵已列出 Python 的普通、FIFO、延迟、事务生产以及 Simple/Push 消费等能力,但官方网站个别 SDK 概览页面的语言列表更新可能滞后。
实际采用前必须做兼容性验证:
- 目标 Broker/Proxy 版本;
- Python SDK 的发布包、成熟度和平台支持;
- 事务回查、FIFO、TLS、认证、可观测性是否齐全;
- 团队能否得到稳定运维支持。
学习中可以掌握其架构与语义;如果当前 Python 项目只需可靠任务队列,RabbitMQ 往往更省集成成本。如果公司基础设施已标准化 RocketMQ,则遵循平台规范通常比个人偏好重要。
12.11 RocketMQ 的优势与边界
优势:
- 业务消息类型丰富;
- FIFO、延迟、事务、重试/死信是核心能力;
- 大规模业务消息和电商场景经验深;
- 消息 key、tag 和查询运维对业务排障友好;
- 5.x gRPC/Proxy 改善多语言和云原生接入。
边界:
- 海外和 Python 社区生态通常不如 Kafka/RabbitMQ 广;
- 版本、协议、SDK 代际需要仔细辨别;
- 数据平台连接器和流计算生态通常 Kafka 更占优势;
- 简单小项目部署整套 RocketMQ 可能过重。
第 13 章:三者如何选型
13.1 先给结论,再解释条件
对于你当前“Python + FastAPI 初学者”的第一套实战,我推荐 RabbitMQ:AMQP 模型直观,Python 异步客户端成熟,ACK、prefetch、DLQ 能完整训练可靠任务消费。
学完 RabbitMQ 后用 Kafka 做第二个数据流实验:重点学习 partition、offset、consumer group、回放和 CDC/流处理。
RocketMQ 要重点掌握选型与高级语义,若你的企业环境使用它,再基于该环境的 5.x SDK 做生产编码。它的事务、顺序和延迟消息非常值得理解,但不必为了练 Python 而强行把第一个项目选成它。
13.2 对比矩阵
| 维度 | RabbitMQ | Kafka | RocketMQ |
|---|---|---|---|
| 核心直觉 | 智能路由 + 任务队列 | 分区追加日志 + 事件流 | 业务消息 + 事件流 |
| 典型优势 | 灵活路由、ACK、低延迟任务 | 高吞吐、保留回放、数据生态 | 事务/FIFO/延迟等业务能力 |
| 消费进度 | 队列投递/ACK | consumer group offset | consumer group 消费进度 |
| 消费后数据 | 队列消息通常确认后可删除 | 按保留策略保留,可重读 | 保留期内可回溯,逻辑标记进度 |
| 路由 | exchange 非常灵活 | 主要靠 topic/partition,应用侧分流 | topic + tag/SQL 过滤 |
| 顺序 | 单队列/单活等条件下 | partition 内 | message group / queue 内 |
| 延迟 | TTL+DLX、插件或调度设计 | 通常 retry topic/调度服务 | 原生 Delay 类型 |
| 事务消息 | 无同款核心模型;常用 Outbox | Kafka 内事务强,外部 DB 仍需设计 | 原生半消息 + 回查 |
| 大规模回放 | Stream 类型可做 | 核心能力 | 保留期内支持回溯 |
| Python 体验 | 很成熟,aio-pika/pika | 成熟,confluent-kafka/aiokafka | 5.x 官方多语言能力在发展,需验版本 |
| 常见场景 | 异步任务、业务路由、通知 | 日志、CDC、流计算、事件平台 | 电商交易、订单、金融业务消息 |
13.3 选型问题树
flowchart TD
A["最主要需求是什么?"] --> B{"事件保留、回放、超高吞吐、数据流?"}
B -->|"是"| K["优先 Kafka"]
B -->|"否"| C{"复杂路由、任务队列、Python 快速落地?"}
C -->|"是"| R["优先 RabbitMQ"]
C -->|"否"| D{"原生事务/FIFO/延迟是核心且组织已有生态?"}
D -->|"是"| M["优先 RocketMQ"]
D -->|"否"| E["按团队能力、托管服务、成本做 PoC"]
13.4 不要用单一指标选型
“Kafka 吞吐最高,所以都用 Kafka”是错误推理。技术选型至少评估:
- 业务语义:任务、事件流、事务、延迟、顺序、广播;
- 吞吐、峰值、消息大小、保留期、可接受延迟;
- 可丢/可重/可乱的边界;
- 客户端语言和功能成熟度;
- 团队已有平台和运维经验;
- 云托管服务、成本、跨地域和合规;
- 监控、重放、审计和故障恢复能力;
- 用代表性负载做 PoC,而不是相信脱离硬件和配置的 TPS 宣传数字。
13.5 常见场景推荐
| 场景 | 首选倾向 | 原因 |
|---|---|---|
| FastAPI 邮件/图片/通知任务 | RabbitMQ | 工作队列、ACK、prefetch 直观 |
| 复杂 routing key 分发 | RabbitMQ | topic/direct/fanout exchange |
| 数据库 CDC 到数仓 | Kafka | 事件保留、Connect/生态、回放 |
| 日志与埋点流 | Kafka | 批量高吞吐、多下游独立消费 |
| 实时聚合和窗口计算 | Kafka + 流计算 | 分区流、offset、状态处理生态 |
| 订单 30 分钟超时关闭 | RocketMQ 或调度系统 | 原生延迟消息自然 |
| 本地事务与发消息一致 | RocketMQ 事务消息或通用 Outbox | 两种可靠路径 |
| 同一订单严格局部顺序 | RocketMQ/Kafka 均可 | message group 或 partition key |
| 公司已有稳定统一 MQ 平台 | 优先公司平台 | 运维、权限、监控、成本更关键 |
13.6 多 MQ 共存是不是更先进
不一定。RabbitMQ 做在线任务、Kafka 做数据流是合理组合,但每增加一种中间件都增加部署、权限、监控、值班、SDK 规范和数据一致性成本。只有当第二种 MQ 带来的关键能力明显超过复杂度时才引入。
第四部分:高级场景与通用解法
第 14 章:分布式事务与最终一致性
14.1 真正的问题:两个系统无法用一个本地事务
订单数据库和 MQ 是两个独立系统:
1 | |
无论先后都有窗口。
先提交数据库,再发消息
数据库成功后进程崩溃,消息没发:订单存在,但库存/积分永远不知道。
先发消息,再提交数据库
消息发出后数据库失败:消费者看到一个根本不存在或已回滚的订单。
仅仅交换两行代码无法解决原子性。
14.2 先区分强一致与最终一致
- 强一致:事务完成的任意可见时刻,各系统状态满足一致规则。
- 最终一致:短时间允许不一致,只要没有新的失败,系统最终收敛到正确状态。
异步消息通常追求最终一致。用户可能先看到“订单已创建,积分处理中”,随后积分到账。关键是状态要透明、可追踪、可补偿,不能把暂时不一致伪装成已完成。
14.3 方案一:Transactional Outbox(通用首选)
在订单数据库中增加 outbox 表。本地事务同时写订单和待发消息:
1 | |
因为两次写在同一个数据库本地事务中,它们要么一起成功,要么一起失败。独立 relay 持续扫描 outbox 并发到 MQ:
flowchart LR
API["订单 API"] -->|"同一 DB 事务"| DB[("orders + outbox")]
DB --> Relay["Outbox Relay"]
Relay --> MQ["MQ"]
MQ --> Consumer["消费者:幂等处理"]
为什么仍可能重复
relay 发送成功后,在把 outbox 标记为 PUBLISHED 前崩溃。重启后会再次发送。因此 Outbox 解决“不丢”,消费者幂等解决“重复”。二者缺一不可。
Relay 两种实现
- 轮询发布:定时
SELECT ... FOR UPDATE SKIP LOCKED取待发布行。简单通用,但有轮询延迟和数据库压力。 - CDC 发布:Debezium 等读取数据库事务日志,把 outbox 变化写到 Kafka。吞吐和解耦更好,基础设施更复杂。
14.4 方案二:RocketMQ 事务消息
RocketMQ 事务消息把消息先存成消费者不可见的半消息,再执行本地事务,最后 commit/rollback。若最终确认丢失,Broker 回查生产者的本地事务状态。
生产者必须保存可查询的事务状态,例如通过订单表判断:
1 | |
关键限制:
- 它实现上游本地事务与消息可见性的最终一致,不是把所有下游数据库变成一个 ACID 大事务;
- 下游仍需幂等和重试;
- 回查函数必须快速、稳定、可重复;
- 不应长时间返回 unknown 或制造大量半消息;
- 只适合可以接受异步最终一致的场景。
14.5 方案三:Kafka 事务
如果输入和输出都在 Kafka:
1 | |
可以放入同一个 Kafka 事务,实现 Kafka 边界内的原子 consume-transform-produce。它非常适合流处理。
如果还要写普通业务数据库,常见选择仍是:
- 数据库幂等更新 + 手动提交 offset;
- 把 offset 与业务数据存入同一数据库事务(自行管理恢复);
- Outbox/CDC;
- 业务状态机和补偿。
14.6 Saga 与补偿
跨多个长事务时,不一定能“回滚一切”,而是执行补偿动作:
1 | |
补偿不是数据库回滚的完美逆操作:退款可能有手续费,短信无法“撤回”,物流可能已出库。因此 Saga 要明确每步的正向动作、补偿动作、幂等键、超时和人工介入状态。
14.7 状态机比“几个布尔值”可靠
不要设计:
1 | |
它可能出现同时 is_cancelled=true、is_shipped=true 的矛盾。使用明确状态机:
1 | |
每个消息只允许合法迁移,并使用版本号/条件更新防止旧事件覆盖新状态。
14.8 对账是最后一道防线
任何复杂分布式系统都应接受“在线机制可能有漏网异常”。定期对账:
- 订单已支付但账务无流水;
- outbox 长期 PENDING;
- 发货状态与物流状态不一致;
- 支付网关记录与本地记录不一致。
检测后自动补发、补偿或生成人工工单。可靠系统不是“绝不出错”,而是“错误可发现、可恢复、可审计”。
第 15 章:顺序消息
15.1 先问:要保证谁和谁的顺序
“保证顺序”是不完整需求。必须说清:
- 同一个生产者的发送顺序?
- Broker 存储顺序?
- 消费者收到顺序?
- 业务提交完成顺序?
- 全部消息全局顺序,还是同一订单局部顺序?
真正业务通常需要:同一聚合根(order_id/user_id/account_id)的状态变更有序。
15.2 端到端有序的条件
以订单事件为例:
- 生产者按业务顺序发;
order_id稳定映射到同一 partition/message group/queue;- Broker 在该顺序域内有序存储和交付;
- 消费端不无序并发;
- 失败重试不能悄悄让后续越过;
- 业务数据库用状态机/版本号防止旧写覆盖新写。
只满足第 3 条并不能叫业务有序。
15.3 Kafka 实现
1 | |
同一 key 到同一 partition;组内由同一 consumer 处理该 partition。消费端按 partition 串行或使用 keyed executor。
15.4 RocketMQ 实现
发送 FIFO 消息时将 message_group 设为 order_id。使用支持 FIFO 的消费方式,回调内同步处理后返回结果。不要把整个业务都放到一个 group,否则所有订单退化为串行。
15.5 RabbitMQ 实现
可选:
- 单队列 + Single Active Consumer;
- 一致性哈希/应用路由,把同一
order_id发到固定分片队列,每队列串行消费; - 消费端按 key 加锁或使用状态机。
第一种最简单但吞吐低,第二种能并行但扩分片时要处理映射变化。
15.6 版本号是顺序的第二保险
事件附带聚合版本:
1 | |
消费者保存已应用版本:
- 收到
version = current + 1:正常应用; - 收到
version <= current:重复或旧事件,忽略; - 收到
version > current + 1:发现缺口,暂停该 key、重试或回源修复。
这样即使基础设施发生意外乱序,业务也不会默默写坏数据。
第 16 章:延时消息与定时任务
16.1 两类时间需求
- 延时:从现在起 30 分钟后处理。
- 定时:在 2026-08-13 09:00:00 处理。
业务语义还应说明:允许晚多久、能否取消、能否修改时间、是否只执行一次、时区是什么。
16.2 订单超时取消的竞态
下单时发送“30 分钟后检查”消息。到期时不能直接取消,必须再次读取订单状态:
1 | |
用户可能在第 29 分 59 秒付款,而延迟消息第 30 分钟到达。状态机条件更新决定谁赢,而不是相信消息名称“cancel-order”就无条件取消。
16.3 各 MQ 的实现选择
- RocketMQ:原生 Delay 消息最自然,设置交付时间戳。
- RabbitMQ:TTL + DLX、延迟插件(若平台允许)或独立调度器。
- Kafka:延迟等级 retry topic、Kafka Streams 状态存储,或独立调度服务。
- 数据库任务表:
execute_at+FOR UPDATE SKIP LOCKED轮询,简单可审计,适合中小规模。
16.4 时间轮直觉
海量定时任务逐条创建系统定时器成本高。时间轮把未来时间划为槽:秒针移动到某槽时处理其中任务;超出一圈的任务记录剩余圈数或进入多级时间轮。很多延迟调度系统都可用这个直觉理解。
16.5 延迟消息不等于准时执行
MQ 最多保证到期后变为可投递。实际完成时间还受:
- Broker 调度精度;
- 到期消息瞬时数量;
- 队列积压;
- 消费者处理能力;
- 网络和下游数据库。
监控应看 actual_processed_at - scheduled_at 的调度延迟分布,而不仅看“消息有没有发”。
第 17 章:即时通讯
17.1 MQ 在 IM 中的位置
WebSocket 负责客户端长连接和实时双向传输;MQ 负责服务端内部的解耦、路由、缓冲和异步处理。二者不是替代关系。
flowchart LR
U1["用户 A"] <-->|"WebSocket"| G1["Gateway 1"]
G1 --> API["消息服务"]
API --> DB[("消息存储")]
API --> MQ["消息总线"]
MQ --> G2["Gateway 2"]
G2 <-->|"WebSocket"| U2["用户 B"]
MQ --> Push["离线推送"]
MQ --> Audit["审核/风控"]
17.2 一个聊天消息的旅程
- 客户端生成
client_msg_id并发送给 Gateway; - 消息服务鉴权,分配服务器序列号并持久化;
- 服务端向发送者返回“已接收”;
- 发布
chat.message.created到 MQ; - 目标用户所在 Gateway 订阅并通过 WebSocket 推送;
- 用户不在线则写离线队列/推送服务;
- 接收端 ACK,更新 delivered/read 状态。
17.3 去重、顺序和多端同步
- 客户端因网络超时会重发,服务端以
(sender_id, client_msg_id)去重; - 每个会话维护递增
conversation_seq,客户端检测缺口并拉历史; - 同一用户手机、网页、平板是多个订阅端,不应简单用“竞争消费只给一端”;
- 在线推送失败不代表消息丢失,持久消息仍可通过历史同步恢复。
17.4 Kafka、RabbitMQ、RocketMQ 在 IM 中的角色倾向
- Kafka:大规模消息事件流、审计、推荐/风控、多下游分析和回放。
- RocketMQ:业务消息、按会话/用户顺序、离线通知、延迟提醒。
- RabbitMQ:Gateway 内部任务、在线路由、通知工作队列;大量短生命周期队列要谨慎评估。
大型 IM 常同时有专用存储、路由层和消息总线,不会仅靠一个 MQ 解决所有在线状态和历史存储问题。
第 18 章:数据流、CDC 与实时计算
18.1 从“任务”转向“持续事件流”
任务队列关心“这件事谁来做完”;事件流关心“业务中持续发生了什么,多个下游如何独立读取、转换、聚合和重放”。
例如用户行为:
1 | |
推荐系统、风控、实时大屏、数仓都想读取这条流,但处理方式和进度不同。
18.2 CDC:把数据库变化变成事件
Change Data Capture 从数据库事务日志(如 MySQL binlog、PostgreSQL WAL)捕获增删改并写入事件流。
flowchart LR
DB[("业务数据库/WAL")] --> CDC["CDC Connector"]
CDC --> K["Kafka Topics"]
K --> ES["搜索索引"]
K --> DW["数据仓库"]
K --> Cache["缓存/物化视图"]
K --> Flink["实时计算"]
Kafka 在该场景优势明显:保留、回放、分区、多 consumer group 和 Connect 生态。
CDC 不是无脑复制:要处理表结构变更、快照与增量衔接、主键、删除事件、事务边界、敏感字段和下游幂等。
18.3 流处理的关键概念
无状态转换
过滤、映射:例如只保留支付成功事件,把金额分转换为元。
有状态计算
聚合、去重、join 需要保存状态。例如统计每分钟每个商品销量。
Event Time 与 Processing Time
- event time:事件实际发生时间;
- processing time:计算程序处理时间。
网络延迟会让旧事件晚到。按 processing time 统计可能把 10:59:59 的订单算进 11:00 窗口;严谨实时计算通常用 event time、watermark 和允许迟到策略。
Window
- Tumbling window:互不重叠固定窗口,如每 1 分钟。
- Sliding window:滑动重叠,如过去 5 分钟每 1 分钟计算一次。
- Session window:按用户活动间隔划分会话。
18.4 KTable / 物化视图直觉
事件流是“变化历史”,表是“某时刻最新状态”。将按 key 的更新流折叠,就形成物化视图:
1 | |
Kafka compacted topic 常用来支持这种可恢复状态。
18.5 数据流中的恰好一次仍有边界
Flink/Kafka Streams 可以把状态快照、输入进度和输出协调起来,提供很强的处理语义。但外部 sink 是否幂等、是否支持事务仍影响端到端结果。例如向第三方 HTTP API 发请求,checkpoint 无法撤回已经发出的请求,仍需业务幂等键。
第五部分:Python + FastAPI 完整项目
第 19 章:项目设计——可靠下单系统
19.1 项目目标
我们实现一个最小但可靠的订单—积分链路:
- FastAPI 接收创建订单请求;
- PostgreSQL 同一事务写
orders和outbox_events; - Outbox Relay 发布
order.created到 RabbitMQ; - 积分消费者收到事件,在同一数据库事务中写 Inbox 去重记录并增加积分;
- 业务成功后才 ACK;
- 暂时失败进入延迟重试,达到上限进入 DLQ;
- 可以故意杀进程,观察重复但不重复加积分。
这不是“生产万能模板”,但覆盖了真正重要的可靠性闭环。
19.2 架构
flowchart LR
Client["客户端"] --> API["FastAPI Order API"]
API -->|"一个本地事务"| PG[("PostgreSQL\norders + outbox")]
PG --> Relay["Outbox Relay"]
Relay -->|"publisher confirm"| EX["RabbitMQ topic exchange"]
EX --> Q["points queue"]
Q --> Worker["Points Worker"]
Worker -->|"Inbox + 积分同一事务"| PG
Q --> Retry["Retry Queue 10s"]
Retry --> Q
Q --> DLQ["Dead Letter Queue"]
19.3 为什么 API 不直接发 MQ
API 只做数据库本地事务。只要接口返回成功,outbox 行一定存在;即使 RabbitMQ 当时不可用,relay 恢复后仍会发送。返回语义是:
1 | |
建议 HTTP 返回 202 Accepted 或明确的订单状态,而不是谎称所有下游都已完成。
19.4 项目结构
1 | |
19.5 可靠性不变量
写代码前先写系统不变量:
- 订单存在,则对应 outbox 事件一定存在;
- outbox 发布可能重复,但不能永久静默丢失;
- 同一
message_id对积分服务最多产生一次业务效果; - 数据库事务成功后才 ACK;
- 达到最大重试次数的消息不会无限循环,而是进入 DLQ;
- 所有消息带
message_id、order_id、trace_id和版本。
如果代码不能证明这些不变量,就不能仅凭“测试发了一条成功”声称可靠。
第 20 章:RabbitMQ 项目代码
下面按文件给出可组装的教学项目。为了聚焦 MQ,开发环境用 SQLAlchemy
create_all;正式环境请使用 Alembic 迁移、密钥管理、TLS、连接池和完整监控。
20.1 依赖
requirements.txt:
1 | |
20.2 本地基础设施
compose.yaml:
1 | |
单节点容器只用于学习,quorum queue 在单节点上没有真正的节点级高可用。生产要部署奇数副本 RabbitMQ 集群并按故障域设计。
20.3 配置
app/config.py:
1 | |
20.4 数据库连接
app/db.py:
1 | |
20.5 数据模型
app/models.py:
1 | |
20.6 RabbitMQ 封装和拓扑
app/messaging.py:
1 | |
这里有几个故意体现可靠性的点:
connect_robust负责连接恢复,但恢复期间的业务语义仍要靠重试和幂等;publisher_confirms=True;mandatory=True,未路由消息不会被悄悄忽略;- 消息持久化;
- 队列为 durable quorum queue;
- 重试发布得到 confirm 后,消费者才能 ACK 原消息。
20.7 FastAPI:订单和 Outbox 同事务
app/main.py:
1 | |
注意:API 中没有 publisher.publish()。这是正确的——本地事务与 Outbox 已经把“订单存在”和“事件等待发送”绑定起来。
20.8 Outbox Relay
app/outbox_worker.py:
1 | |
这个教学实现为了让原理清晰,在网络发布期间持有一条数据库行锁。高吞吐生产版本应使用批量 claim:快速把一批行标记 PROCESSING 并提交,再发布;同时加入 claimed_at 租约和超时回收器,避免 worker 崩溃后永久卡在 PROCESSING。无论哪种实现,都必须接受“confirm 后、状态提交前崩溃”造成重复发布。
20.9 积分消费者:Inbox 幂等 + 手动 ACK
app/points_worker.py:
1 | |
一个并发细节
两个重复消息若第一次同时处理,同一用户账户“查不到后各自 INSERT”可能产生主键竞争。数据库会让其中一个事务失败并重试,Inbox 唯一键保证同一 message_id 不产生两次效果。生产中可预创建积分账户,或使用 PostgreSQL 原子 upsert:
1 | |
必须与 Inbox 插入放在同一事务中。
20.10 启动与验证
在项目目录依次运行:
1 | |
另开两个终端:
1 | |
1 | |
创建订单:
1 | |
查询积分:
1 | |
RabbitMQ 管理界面是 http://localhost:15672,用户名和密码均为 app。你应该能看到工作队列、重试队列、死信队列以及 ready/unacked 数量。
第 21 章:故障实验与验收
“能跑”不等于“可靠”。下面用故障主动验证不变量。
21.1 RabbitMQ 停机时创建订单
- 停止 RabbitMQ;
- 调用创建订单 API;
- API 应仍返回订单创建成功;
- 数据库中 outbox 为
PENDING; - 重启 RabbitMQ;
- relay 最终发布,积分增加。
验证的是:MQ 暂时不可用不会让已提交订单丢失事件。
21.2 Relay 在确认后、状态提交前崩溃
在 publish_event() 后、event.status = "PUBLISHED" 前临时加入抛异常或强制退出。消息已到 Broker,但 outbox 仍 PENDING,重启会重复发。
预期:消费者打印一次 processed、另一次 duplicate ignored,积分只增加 10。
验证的是:至少一次发布 + Inbox 幂等。
21.3 消费者在数据库提交后、ACK 前崩溃
在 apply_points_once 返回后、message.ack() 前强制退出消费者。Broker 会把 unacked 消息重新投递。
预期:Inbox 判断重复,积分不再增加,随后 ACK。
验证的是:业务提交与 ACK 之间重复窗口不会破坏结果。
21.4 暂时性数据库故障
消费时停止 PostgreSQL。消息应进入重试路径或保留,不能 ACK 后消失。恢复数据库后最终成功。
观察:
- retry queue 是否出现消息;
x-retry-count是否递增;- 重试是否存在 10 秒间隔;
- 恢复后积分是否只加一次。
21.5 永久坏消息
向 exchange 发布缺少 payload.user_id 的 JSON。消费者应识别契约错误,直接进入 DLQ,而不是重试五次。
验证的是:永久错误与暂时错误分类。
21.6 积压和扩容
停止 points worker,快速创建 1000 个订单,再观察 ready 消息增长。启动一个 worker,记录消费速率;再启动第二个 worker,观察是否分摊、吞吐是否提升以及数据库是否成为瓶颈。
不要只看队列条数,还看最老消息年龄。1 万条每条 1ms 和 100 条每条 10 分钟的严重性完全不同。
21.7 验收表
| 故障 | 允许结果 | 禁止结果 |
|---|---|---|
| MQ 短暂不可用 | outbox 积压,恢复后发送 | 订单存在但事件永久消失 |
| 发布响应丢失 | 可能重复发布 | 静默丢失 |
| 消费者 ACK 前宕机 | 重投、业务去重 | 重复加积分 |
| 消息格式错误 | DLQ + 告警 | 无限重试 |
| 消费者停机 | Broker 积压、恢复后追平 | API 级联崩溃 |
| 下游持续过载 | 限流/降级/扩容/告警 | 队列无限增长至磁盘耗尽 |
第 22 章:同一业务迁移到 Kafka
RabbitMQ 项目教会了任务投递和 ACK。迁移到 Kafka 时,业务不变量不变,机制映射发生变化:
| RabbitMQ | Kafka |
|---|---|
| exchange + routing key | topic(事件类型可放 value/header) |
| queue | consumer group 的逻辑订阅 |
| ACK | 提交 offset |
| prefetch | poll 批量、pause/resume、处理并发 |
| DLX/DLQ | 单独 retry/DLQ topic + 应用逻辑 |
| publisher confirm | acks、幂等 producer 和发送结果 |
22.1 异步 Producer
1 | |
Outbox Relay 只需把 Rabbit publisher 替换为这个 adapter。Outbox 仍然必要,因为 Kafka 幂等生产者解决的是 Kafka 协议重试重复,不解决“PostgreSQL 已提交但程序还没调用 producer 就崩溃”。
22.2 手动提交的幂等消费者
1 | |
生产版还要:
- rebalance listener;
- 有限重试和 retry/DLQ topics;
- 每 partition 批处理和连续 offset 提交;
- 处理耗时与 poll 超时协调;
- 监控 lag、rebalance 和 commit failure。
22.3 两个 consumer group 演示发布订阅
再启动 analytics-service-v1 group。它和 points-service-v1 都能读到所有订单事件;同一 group 内启动多个实例则分摊 partitions。
你应该亲自实验:
- topic 创建 3 partitions;
- points group 启 1、2、4 个实例,观察分配;
- 停掉一个实例,观察 rebalance;
- 将 group offset reset 到 earlier,重新消费历史;
- 使用相同
order_id发多条事件,观察同 partition 顺序; - 不带 key 发送,观察顺序域变化。
22.4 Kafka 内事务示例
若消费 raw-orders 并把清洗结果写 clean-orders,可用:
1 | |
每个并行 producer 实例要有稳定且唯一的 transactional_id,避免互相 fencing。真实代码需处理 rebalance、空批次和异常关闭。
第 23 章:RocketMQ 专项实践
本章代码基于 Apache
rocketmq-clients仓库当前 Python 5.x gRPC SDK 形态,用于理解 API 映射。该 SDK 与服务端/Proxy 版本耦合,实际运行前请以你安装版本的官方示例为准。安装包名在当前仓库setup.py中为rocketmq-python-client。
23.1 普通消息
1 | |
生产中仍应由 Outbox Relay 调用 producer,而不是在数据库提交后“顺手发一下”。
23.2 FIFO 消息
1 | |
topic 必须按服务端要求创建为 FIFO 类型。消费端也必须使用 FIFO listener/配置并同步完成处理。
23.3 延迟消息
1 | |
到期消费者仍要执行 cancel_if_status_is_pending 条件更新,不能无条件取消。
23.4 SimpleConsumer:处理后再 ACK
1 | |
不可见时间必须大于正常处理耗时,并处理消费超时导致的重复。如果业务时长不可预测,需要续期/修改不可见时间的相应 SDK 能力或调整消费模式。
23.5 事务消息代码骨架
1 | |
本地事务状态必须持久化并可供回查,不能把结果只存在内存变量里。还要认真处理“数据库提交结果不确定”的情况,不能见异常就武断 rollback。
23.6 三个项目练习如何串起来
- RabbitMQ 项目练可靠任务语义:confirm、ACK、prefetch、retry、DLQ、Inbox。
- Kafka 改造练日志语义:partition、key、offset、group、lag、回放、Kafka 内事务。
- RocketMQ 专项练业务高级消息:message group、delivery timestamp、半消息和回查。
学完后,你面对新 MQ 也能从数据模型、确认边界、进度、顺序、重试和副本六个方面快速理解,而不再依赖背 API。
第六部分:生产工程、排障与能力检验
第 24 章:生产环境设计清单
24.1 先写 SLO,而不是先调参数
上线前把“可靠、实时、高吞吐”翻译成数字:
| 目标 | 示例 |
|---|---|
| 峰值生产速率 | 20,000 msg/s |
| 平均/最大消息大小 | 2 KB / 256 KB |
| 端到端延迟 | p99 < 3 s |
| 可接受丢失 | 核心订单 0;非核心埋点可少量丢 |
| 可接受重复 | 基础设施可重复,业务效果不可重复 |
| 保留期 | 订单事件 7 天;审计事件 180 天 |
| RTO | Broker 故障后 5 分钟恢复服务 |
| RPO | 已确认订单事件不丢 |
| 可容忍积压 | 峰值 30 分钟并在 2 小时内追平 |
没有这些数字,无法判断副本数、磁盘、分区数、消费者数量和报警阈值。
24.2 容量估算
原始写入带宽
[
Bandwidth=messageRate\times averageMessageSize
]
20,000 msg/s、平均 2 KB,原始约 40 MB/s。还要加协议开销、索引、复制、消费者出流量和安全余量。
保留存储
粗略估算:
[
Storage=rate\times size\times retentionSeconds\times replicationFactor\times overhead
]
若 2,000 msg/s、1 KB、保留 7 天、副本 3,忽略压缩和额外开销:
[
2000\times1KB\times604800\times3\approx3.63TB
]
真实结果受压缩率、segment、索引和副本策略影响。至少留 30%—50% 安全水位,并进行真实消息压测。
积压恢复时间
[
T_{recover}=\frac{backlog}{consumeRate-produceRate}
]
如果消费只能刚好追上生产,积压永远不会下降。容量设计必须留“追债能力”。
24.3 分区和队列数量
Kafka/RocketMQ 的并行度常受 partition/message queue 数限制。估算:
[
Partitions\ge\max(\frac{targetProduce}{singlePartitionProduce},
\frac{targetConsume}{singlePartitionConsume})
]
单分区能力必须在目标硬件、消息大小、压缩、acks 和消费者逻辑下压测。分区太少限制并行;太多增加元数据、文件、副本、选主和 rebalance 开销。扩分区还可能影响 key 映射和顺序。
RabbitMQ 不应靠无限创建 queue 解决所有并行问题;队列数量、每队列 backlog、consumer 数、quorum 副本都会消耗资源。
24.4 高可用不是“集群三个节点”一句话
要逐层问:
- 生产者能否发现并重连健康节点?
- 数据是否复制,而不只是元数据复制?
- ACK/confirm 在何种复制条件下返回?
- 多数派丢失时选择停止写还是冒险写?
- 节点是否跨故障域/可用区,而不是三台虚拟机在同一宿主机?
- 负载均衡器是否正确处理长连接?
- DNS、证书、时钟、磁盘和网络分区是否演练过?
24.5 生产者清单
- 连接长时间复用,配置合理超时;
- 关键消息等待 publisher confirm /
acks=all/ 同步发送结果; - 处理超时的不确定状态,带稳定
message_id; - 使用 Outbox 或事务消息解决数据库—MQ 双写;
- 重试有退避和总超时,不在无限循环中压垮 Broker;
- 记录 publish latency、失败类型、topic、message_id、trace_id;
- 设置最大消息大小,发送前校验;
- 应用关闭时停止接收新请求、flush 生产缓冲、等待在途确认;
- 不把凭据、隐私数据打印进日志。
24.6 消费者清单
- 业务成功提交后才 ACK/commit;
- 使用数据库唯一约束/Inbox/状态机实现幂等;
- 暂时错误有界重试,永久错误直接 DLQ;
- 并发、prefetch、批量和数据库连接池相匹配;
- 消息处理有超时,长任务可拆分或续期;
- 顺序域明确,不在无序线程池破坏顺序;
- 正确处理 rebalance、连接恢复和优雅关闭;
- 停机前停止拉新消息,等待在途业务事务完成,再提交进度;
- 对 DLQ 有负责人、告警和重放流程;
- 消费者版本发布支持灰度和回滚。
24.7 Broker 和主题清单
- 生产和开发隔离,关键业务与非关键流量隔离;
- 副本、最小同步副本、quorum 成员符合 RPO;
- 磁盘容量、IOPS、网络和文件句柄有预算;
- retention、TTL、最大长度和消息大小有明确策略;
- 禁止或谨慎使用自动创建 topic/queue,避免拼写错误生成孤儿资源;
- topic/queue 命名、owner、数据分级、保留期登记;
- TLS、认证、ACL/vhost/namespace 最小权限;
- 备份的是配置还是消息数据要说清;
- 升级、滚动重启、副本重建和灾备切换做过演练。
24.8 连接和协程管理
在 FastAPI 中:
- producer connection 放应用 lifespan 中创建并复用;
- 不在每个请求中连接/关闭 Broker;
- 不让无限制
asyncio.create_task逃逸; - 消费者通常作为独立进程部署,不和 Web worker 绑在一起,否则 Web 自动扩缩容会引发消费实例频繁变化;
- CPU 密集任务移至进程池/专门 worker,不要阻塞 asyncio 事件循环。
24.9 Schema 治理
一个严谨流程:
- Schema 进入版本控制;
- CI 校验生产者变化是否向后兼容;
- 消费者先兼容新旧字段并上线;
- 生产者再开始发送新字段;
- 观察所有消费者升级;
- 经过弃用窗口后才移除旧字段;
- 保留示例消息和契约测试。
24.10 多机房和灾备
跨地域同步复制有高延迟和网络分区问题。要明确:
- 主备还是双活;
- 是否允许两个区域同时写同一顺序域;
- 故障切换时如何避免双写和 offset 冲突;
- 异步复制的 RPO 是多少;
- DNS/客户端切换多久;
- 回切如何处理重复事件。
Kafka 可使用 MirrorMaker 类跨集群复制方案;其他 MQ 也有 Federation、Shovel、集群/复制等不同能力。跨集群复制通常提供灾备与分发,不自动为业务创造全局强一致。
第 25 章:监控、容量与排障
25.1 四个黄金信号
流量
生产/消费 msg/s、byte/s、批量大小、各 topic/queue 分布。
延迟
- 发布确认延迟;
- Broker 内等待时间;
- 消费处理耗时;
- 端到端延迟
processed_at - occurred_at; - 最老未处理消息年龄。
错误
发送失败、消费失败、重试率、DLQ 速率、反序列化错误、权限错误、offset commit/ACK 失败。
饱和度
CPU、内存、磁盘使用率/IOPS、网络、文件句柄、连接/通道数、线程/协程、数据库连接池。
25.2 为什么“队列长度”不够
队列 10 万条可能只是正常批量,也可能已经延迟 2 小时。至少同时看:
- ready/backlog 数量;
- unacked/inflight 数量;
- 最老消息年龄;
- 生产与消费速率差;
- 按分区/队列的倾斜;
- 预计清空时间。
25.3 产品侧重点
RabbitMQ
- messages ready / unacked;
- publish、deliver、ack、redeliver rate;
- connection/channel/consumer 数;
- memory/disk alarm;
- quorum queue leader/member 健康;
- unroutable returned messages;
- DLQ 和 retry queue 增长。
Kafka
- consumer lag 和 lag 时间;
- under-replicated/offline partitions;
- ISR shrink/expand;
- request latency、produce/fetch error;
- broker disk、network、page cache 相关表现;
- controller/quorum 健康;
- rebalance 和 commit failure;
- 热 partition。
RocketMQ
- send/consume TPS 与延迟;
- consumer lag、消费时间差;
- Broker/Proxy/NameServer 可用性;
- CommitLog/存储磁盘水位;
- 重试和 DLQ;
- transaction half message 回查与堆积;
- delay/FIFO 消息处理延迟;
- queue/broker 分布倾斜。
25.4 分布式追踪
将 trace_id 放消息 header 或信封,生产 span 记录发送,消费 span 使用消息上下文建立关联。异步链路通常不是传统父子调用栈,追踪系统可能用 link 表达因果关系。
每条关键日志至少包含:
1 | |
不要只打印“消费失败”。
25.5 积压排障剧本
flowchart TD
A["发现 lag/最老消息年龄升高"] --> B{"生产速率突增?"}
B -->|"是"| C["检查是否预期峰值;扩容/限流/降级"]
B -->|"否"| D{"消费实例与错误率正常?"}
D -->|"否"| E["恢复实例;查看部署/权限/反序列化/重试"]
D -->|"是"| F{"处理耗时或下游变慢?"}
F -->|"是"| G["数据库/API/连接池/锁/索引排障"]
F -->|"否"| H{"某分区/队列热点?"}
H -->|"是"| I["检查 key 分布和毒消息"]
H -->|"否"| J["检查 Broker 磁盘/网络/复制/流控"]
处理积压时不要立刻无限扩容消费者。若数据库已经满载,扩容会让它更快崩溃。先定位瓶颈,再决定扩容、批量优化、索引优化、限流或降级。
25.6 消息丢失排障
沿消息旅程逐段寻找证据:
- 业务数据库是否提交?outbox 是否有行?
- relay 是否 claim?publish 是否收到 confirm?
- 是否发错 topic/routing key?是否 unroutable?
- Broker 是否按预期持久化和复制?是否过期/达到长度限制?
- 消费组是否订阅正确?offset 是否被错误重置?
- 消费者是否提前 ACK/commit?
- 是否进入 retry/DLQ?
- 是否其实处理成功,只是日志/查询读到旧缓存?
没有 message_id 和审计记录时,这类排障会非常痛苦。
25.7 重复消息排障
重复本身在至少一次系统中不是异常。重点是来源和业务是否幂等:
- 发送超时后的生产者重试;
- relay confirm 后未标记;
- 消费者数据库提交后 ACK 前崩溃;
- ACK/offset commit 失败;
- rebalance 中批次重复;
- 人工重放;
- 上游本就发布两条不同
message_id的重复业务事件。
前五种应由消息级幂等处理;最后一种需要业务唯一键(例如支付单号)进一步防重。
25.8 乱序排障
检查:生产端是否多实例并发无协调、key 是否一致、分区数是否变化、是否跨 partition、重试是否越过、消费者是否并发完成、状态版本是否校验。不要只盯 Broker,因为业务乱序经常发生在生产和消费线程模型中。
第 26 章:反模式与常见误区
26.1 把 MQ 当数据库
MQ 可以持久化,不等于适合任意查询、关系约束和永久业务真相。核心业务状态仍应进入合适数据库;MQ 保存的是传递中的消息或事件历史。
26.2 先 ACK 再处理
进程崩溃会永久丢业务。关键消费必须事务成功后 ACK。
26.3 认为“至少一次”无需幂等
恰恰相反,至少一次的定义就允许重复。没有幂等,重试机制越可靠,重复副作用越可靠。
26.4 在消息中塞完整大对象或文件
放大网络、内存、复制和重试。使用小事件或 Claim Check。
26.5 无限立即重试
毒消息和下游故障会形成 CPU/网络风暴。使用错误分类、有界次数、指数退避、DLQ。
26.6 把 DLQ 当终点
无人告警、无人查看、无法重放的 DLQ 只是更隐蔽的丢失。
26.7 一个 topic/queue 装全世界
不同保留期、权限、吞吐和可靠性互相干扰,契约难治理。
26.8 每个事件一个 topic
另一极端会让资源和运维爆炸。按业务域和运行属性划分,而不是机械按类名划分。
26.9 随意修改 consumer group 名
Kafka/RocketMQ 中 group 名往往代表一套消费进度。改名可能从 earliest 重放全部,或从 latest 跳过历史,取决于配置。
26.10 相信“开启持久化就绝不丢”
端到端还包含发布确认、副本、生产端 Outbox、消费者 ACK 和幂等。任何一段缺失都可能丢。
26.11 追求全局顺序
通常是不必要的吞吐自杀。先缩小到订单/用户/账户级局部顺序。
26.12 在消费者里调用慢 API 却不管超时
消息长期 inflight,触发重投、占满 worker。设置超时、熔断、批量、隔离池;必要时拆成下一阶段事件。
26.13 每个 HTTP 请求新建 MQ 连接
握手昂贵、连接数爆炸。复用长连接,合理使用 channel/session。
26.14 只监控 Broker 存活
Broker 活着不代表业务在前进。consumer lag、最老年龄、DLQ 和端到端成功率更接近用户体验。
26.15 把异步当成性能免费午餐
异步只是把等待移到后台。总工作量没消失,还增加序列化、网络、存储和一致性成本。
26.16 把 MQ 当同步 RPC
Request/Reply over MQ 可以实现,但会继承相关 ID、临时回复、超时、重复响应和资源清理复杂度。如果业务天然是立即问答,HTTP/gRPC 往往更简单。
26.17 双写时靠 try/except
数据库和 MQ 不是一个事务,异常捕获无法消除崩溃窗口。用 Outbox 或产品事务消息。
26.18 用内存集合永久去重
进程重启丢失、无法跨实例、无限增长。关键幂等状态应持久化,并有与消息保留/重放窗口匹配的生命周期。
26.19 为了追吞吐而无界并发
最终耗尽数据库连接、内存和下游配额。使用 semaphore、prefetch、批量、连接池和背压。
26.20 混用不同版本教程
典型例子:Kafka 新集群还照抄 ZooKeeper 命令;RabbitMQ 4.x 还配置已移除的 classic mirroring;RocketMQ 5.x 时间戳延迟与 4.x 固定延迟等级混用;旧 Python RocketMQ 客户端和新 gRPC SDK 混用。学习原理要稳定,API 和部署必须查对应版本官方文档。
第 27 章:面试表达与自测题
27.1 一分钟解释“什么是消息队列”
消息队列是位于生产者与消费者之间的消息中间件。生产者把业务事件或任务交给 Broker,Broker 负责路由、存储和投递,消费者按自己的节奏处理。它主要解决同步调用链响应慢、突发流量冲击下游和服务之间耦合过高的问题,因此常用于异步处理、削峰填谷和事件驱动。但引入后系统会从同步强一致转向异步最终一致,并带来丢失、重复、乱序、积压和排障问题。工程上通常用发布确认、Broker 持久化与副本、处理后 ACK、至少一次投递加消费幂等、有限重试和死信,以及 Outbox/事务消息形成可靠闭环。
这段回答先定义、再价值、再代价、最后机制,逻辑完整。
27.2 解释异步、削峰、解耦
异步是让非核心后续工作不阻塞当前请求,例如订单创建后短信稍后发送;削峰是让 Broker 暂存瞬时超出下游处理能力的消息,把尖峰摊到后续时间,但它不能解决长期生产大于消费,所以还要限流和扩容;解耦是生产者只发布稳定业务事件而不直接依赖所有下游,新增订阅者无需修改生产者,不过消息 Schema 仍是一种契约耦合。
27.3 解释消息不丢
不要回答“开启持久化”。推荐回答:
要端到端分析。生产端用 Outbox 或事务消息避免数据库已提交但消息没发;发送端等待 Broker confirm/acks,并对不确定结果重试;Broker 使用持久化和多副本;消费者在业务事务提交后才 ACK/提交 offset;由于 ACK 前宕机会重复投递,消费者用唯一约束/Inbox/状态机幂等;失败有界重试并进入 DLQ;最后用积压、死信和对账监控兜底。
27.4 解释重复消费
消费者可能处理并提交数据库后,在 ACK 前宕机,Broker 不知道已处理就会重投;生产者发送超时也可能因确认丢失而重发。因此至少一次系统中重复是正常现象。解决重点不是阻止一切重复,而是用稳定 message_id 和数据库唯一约束、Inbox 表或条件状态更新,让重复消息的业务效果等价于一次。
27.5 解释顺序消息
先定义顺序范围。全局顺序会把并行度降为 1,通常只需同一订单局部有序。Kafka 让同一 order_id 作为 key 进入同一 partition,RocketMQ 用 message group,RabbitMQ 可用单活或按 key 分片队列;消费者还必须串行处理同一顺序域,并用聚合版本和状态机防止旧事件覆盖新状态。Broker 有序不等于业务完成有序。
27.6 解释 Kafka、RabbitMQ、RocketMQ 选型
Kafka 的核心是分区追加日志,擅长高吞吐、事件保留回放、CDC 和流处理;RabbitMQ 的 exchange/queue 模型擅长灵活路由、任务队列、ACK/prefetch/DLX,Python 在线业务集成方便;RocketMQ 面向大规模业务消息,原生 FIFO、延迟和事务消息突出。选型不能只看 TPS,要看业务语义、保留回放、延迟/顺序/事务需求、客户端语言、现有平台和团队运维能力。
27.7 基础自测
先自己回答,再看本章末答案。
- Python
asyncio.Queue为什么不能替代 RabbitMQ? - FastAPI
BackgroundTasks为什么不适合扣款? - queue 和 topic 的直觉差异是什么?
- “推消费者”底层是否一定是真正服务器推送?
- 削峰和限流有什么区别?
- 为什么 MQ 降低运行时耦合却没有消灭契约耦合?
- 一个 consumer group 中消费者数超过 Kafka partition 数会怎样?
- 为什么消费后 ACK 仍然会重复?
- 为什么
SETNX不是万能幂等? - DLQ 为什么不是最终解决方案?
27.8 深入自测
- 订单数据库提交成功、消息未发,至少给出两种解决方案。
- Kafka
acks=all为什么还要配合min.insync.replicas? - Kafka exactly-once 为什么不自动覆盖 PostgreSQL?
- RabbitMQ durable queue、persistent message、publisher confirm 分别解决什么?
- RocketMQ 事务消息为什么仍要求下游幂等?
- 为什么同一 partition 有序仍可能业务乱序?
- 积压 100 万条时为什么不能立刻无限扩 consumer?
- 延迟消息到期为什么不能无条件取消订单?
- Outbox relay 为什么会重复发布?
- 如何判断某条消息是暂时性错误还是永久性错误?
27.9 设计题
题一:秒杀下单
要求入口每秒 50,000 请求,数据库只能处理每秒 3,000。请说明入口限流、库存预扣、防超卖、MQ 缓冲、消费者扩容、积压告警和用户状态反馈,不要只说“用 MQ 削峰”。
题二:支付成功事件
要求订单、账务、积分、通知最终一致。画出 Outbox、四个 consumer group、各自 Inbox、重试/DLQ、对账流程,并说明支付消息重复和乱序时如何保证结果。
题三:聊天系统
说明 WebSocket 与 MQ 的分工、会话顺序键、多设备订阅、离线消息、客户端重发去重、服务器序列号和历史补洞。
题四:实时数据平台
说明 PostgreSQL CDC → Kafka → Flink → 搜索/数仓的数据流,以及 Schema 演进、分区键、event time、迟到数据、checkpoint 和 sink 幂等。
27.10 自测答案要点
asyncio.Queue是单进程内存结构,缺少跨机器、持久化、副本、消费组和管理能力。- Web 进程崩溃会丢任务,没有独立持久化和重试闭环。
- queue 偏待办任务和确认后删除;topic/日志偏事件分类、多订阅和保留回放。
- 不一定,常由 SDK 长轮询后回调业务代码。
- 削峰接收并缓冲,限流拒绝/等待/降级以保护资源。
- 生产者不依赖消费者在线和地址,但双方仍依赖消息字段语义。
- 多余实例通常无 partition 可分而空闲。
- 数据库已提交、ACK 未到 Broker 时宕机会重投。
- 去重标记与业务写不原子时,标记成功后业务失败会造成永久跳过。
- 必须告警、分析、修复、重放和审计,否则只是隐藏丢失。
- Transactional Outbox;RocketMQ 事务消息;特定边界还可用数据库日志 CDC。
acks=all基于当前 ISR;若 ISR 可缩到 1,仍可能单副本确认,最小同步副本限制安全下限。- Kafka 事务协调 Kafka 记录和 offset,不控制外部数据库事务管理器。
- 拓扑重启存在;消息持久存储;生产者知道 Broker 是否接管。
- 事务消息只协调上游本地事务和消息可见性,消费者执行仍可能失败或重复。
- 消费端并发线程可能 2 先于 1 完成,重试也会改变完成顺序。
- 瓶颈可能是数据库;无限扩容会耗尽连接并加剧故障。
- 到期时可能已付款,必须用状态机条件更新消解竞态。
- Broker confirm 后、outbox 标记 PUBLISHED 前宕机,重启会再发。
- 网络/超时/限流通常暂时;Schema 缺字段、非法状态通常永久,但要基于明确错误分类而非只按异常类名猜测。
第 28 章:继续学习的路线
28.1 四阶段计划
第一阶段:建立直觉(2—3 天)
- 重读第 1—9 章;
- 能画出 producer → Broker → consumer;
- 能口述异步、削峰、解耦和代价;
- 手写 at-most/at-least/exactly-once 的差异;
- 用一张纸推演“数据库提交后 ACK 前宕机”。
第二阶段:RabbitMQ 可靠项目(4—7 天)
- 跑通第 20 章;
- 完成第 21 章所有故障实验;
- 自己加入邮件 consumer group;
- 给 DLQ 做查询和重放接口;
- 把固定 10 秒重试改成 10s、1m、10m 多级退避;
- 添加 Prometheus 指标和 trace_id 日志。
第三阶段:Kafka 数据流(4—7 天)
- 运行 3 partition、多 group、多实例实验;
- 手动提交和重置 offset;
- 验证同 key 顺序与无 key 分布;
- 制造 consumer crash 和 rebalance;
- 将 Outbox 发布器换成 Kafka;
- 可进一步用 Debezium 把 outbox 通过 CDC 发布。
第四阶段:高级业务与生产化(1—2 周)
- 用 RocketMQ FIFO/Delay/Transaction 做专项 PoC;
- 设计订单状态机和补偿 Saga;
- 做 10 倍峰值压测与积压恢复测试;
- 写 SLO、容量表、告警、运行手册;
- 让另一位开发者故意制造故障,你只看指标和日志排查。
28.2 你真正“学会”的验收标准
你能独立做到以下事项,才算从入门进入实战:
- 不看资料解释 MQ 的价值与代价;
- 从业务一致性要求选择同步或异步;
- 为消息设计 event envelope、幂等键、顺序键和版本;
- 解释 ACK/offset 的提交位置和崩溃窗口;
- 用 Outbox + Inbox 写可靠链路;
- 设计有限重试、指数退避、DLQ 和重放;
- 根据回放、路由、事务、延迟和生态选择三种 MQ;
- 估算吞吐、存储、积压恢复时间和并行度;
- 从 lag、最老年龄、错误率和下游指标定位故障;
- 在面试中不只背“异步削峰解耦”,还能说出它们引入的最终一致性和恢复机制。
28.3 最终心智模型
当你下次看到任何消息系统,依次问十个问题:
- 数据模型是任务队列还是追加日志?
- 生产者成功的确认边界在哪里?
- Broker 如何持久化和复制?
- 消费者如何获得消息,推、拉还是长轮询?
- 进度由 ACK、offset 还是不可见时间管理?
- 失败如何重试,何时进入死信?
- 重复如何幂等?
- 顺序在哪个范围内保证?
- 数据库与消息如何保持最终一致?
- 如何监控、重放、对账和灾备?
能回答这十问,你就不只会 Kafka、RabbitMQ 或 RocketMQ 的 API,而是理解了消息系统本身。
附录
附录 A:术语表
| 术语 | 含义 |
|---|---|
| Producer / Publisher | 生产、发送消息的客户端 |
| Consumer / Subscriber | 读取并处理消息的客户端 |
| Broker | 接收、存储、路由和投递消息的服务器 |
| Topic | 事件逻辑分类,常可分区并被多个组订阅 |
| Queue | 保存待处理消息的队列 |
| Exchange | RabbitMQ 中按规则将消息路由到队列的组件 |
| Binding | RabbitMQ exchange 到 queue 的路由关系 |
| Routing Key | RabbitMQ 路由匹配键 |
| Partition | Topic 的并行有序子日志 |
| Offset | 某 partition 中记录的位置/消费进度 |
| Consumer Group | 共同承担一个逻辑订阅的一组消费者 |
| ACK | 消费者确认业务处理成功 |
| Publisher Confirm | Broker 对生产者的接管确认 |
| At-most-once | 至多一次,可能丢但不重 |
| At-least-once | 至少一次,不丢倾向但可能重 |
| Exactly-once | 在明确边界内业务效果恰好一次 |
| Idempotency | 重复执行和执行一次结果等价 |
| Retry | 暂时失败后的再次处理 |
| DLQ | 死信队列,隔离最终失败消息 |
| Backlog / Lag | 尚未消费的积压或进度差 |
| Backpressure | 下游处理不过来时向上游传播压力 |
| Rebalance | 消费组成员变化后的分区重分配 |
| Retention | 消息保留时间/容量策略 |
| Log Compaction | 按 key 保留较新状态的日志清理策略 |
| Outbox | 与业务数据同事务写入的待发事件表 |
| Inbox | 消费者持久化的已处理消息去重表 |
| CDC | 从数据库事务日志捕获变更 |
| Saga | 一系列本地事务及其补偿动作 |
| FIFO | 先进先出;必须说明顺序范围和条件 |
| ISR | Kafka 中与 leader 足够同步的副本集合 |
| Quorum | 多数派,共识复制常依赖它 |
| KRaft | Kafka 自管理元数据的 Raft 控制面模式 |
| Message Group | RocketMQ FIFO 消息的局部顺序标识 |
| Half Message | RocketMQ 事务消息提交决定前的暂不可见消息 |
| TTL | 消息或队列的生存时间 |
| Poison Message | 持续导致消费失败的毒消息 |
附录 B:一页选型速查
1 | |
附录 C:可靠消息伪代码模板
生产端
1 | |
消费端
1 | |
附录 D:官方资料与延伸阅读
以下链接用于核对产品事实和 API;阅读时选择与你部署版本一致的文档。
Apache Kafka
- Kafka Introduction:事件、Topic、Partition、Consumer Group
- Kafka Design:持久化、复制、投递与 Exactly-once
- Kafka Producer Configuration
- Kafka Consumer Configuration
- Kafka KRaft Operations
- Kafka Streams
- aiokafka Producer:确认、幂等与事务
- aiokafka Consumer:Offset 与手动提交
RabbitMQ
- RabbitMQ Exchanges
- RabbitMQ Reliability Guide
- RabbitMQ Consumer Acknowledgements and Publisher Confirms
- RabbitMQ Consumer Prefetch
- RabbitMQ Dead Letter Exchanges
- RabbitMQ TTL
- RabbitMQ Quorum Queues
- RabbitMQ Streams
- RabbitMQ Python Tutorials
- aio-pika Documentation
Apache RocketMQ
- RocketMQ 5.0 Domain Model
- RocketMQ Message
- RocketMQ FIFO Message
- RocketMQ Delay Message
- RocketMQ Transaction Message
- RocketMQ Consumer Types
- RocketMQ Consumption Retry
- RocketMQ SDK Overview
- Apache RocketMQ 5.x 多语言客户端仓库
- 当前 Python 5.x SDK 示例目录
通用架构
- Microservices.io:Transactional Outbox Pattern
- Debezium Outbox Event Router
- CloudEvents Specification
- FastAPI Lifespan Events
附录 E:版本相关事实说明
- Kafka、RabbitMQ、RocketMQ 仍在持续演进。本文将稳定原理与当前产品能力分开描述。
- Kafka 4.x 的运维方式已以 KRaft 为中心,旧 ZooKeeper 教程只适合历史版本。
- RabbitMQ 4.x 已移除 classic mirrored queues;复制队列优先看 quorum queues 或 streams。
- RocketMQ 4.x 的固定延迟等级示例与 5.x Delay topic/时间戳模型不是同一套 API。
- RocketMQ 官方网站 SDK 概览与
rocketmq-clients仓库的语言能力矩阵可能存在更新节奏差异;Python 接入应以目标发布版本、Proxy 与集成测试共同确认。 aiokafka、aio-pika等第三方客户端的默认配置可能和 Java 官方客户端不同。教程中显式配置关键语义,就是为了避免把一种客户端的默认值误认为产品永恒规则。
结语
消息队列最迷人的地方不是“把一条 JSON 从 A 搬到 B”,而是它迫使我们认真面对分布式系统的现实:网络会超时,进程会在任意一行代码后崩溃,消息会重复,机器速度不同,强一致有成本,恢复能力比完美幻想更重要。
从现在起,看到一条消息时,不要只问“怎么发”,而要问:它代表什么事实、成功边界在哪里、失败后谁负责、重复是否安全、顺序范围是什么、多久必须完成、如何发现并恢复异常。能把这些问题回答清楚,你就已经具备将消息队列用于真实项目的核心能力。