消息队列学习

消息队列从零到实战:Kafka、RocketMQ、RabbitMQ 完整教程

面向读者:会一点 Python、FastAPI 和数据库,但从未系统接触消息队列。

学习目标:不仅“知道名词”,还能够解释原理、做技术选型、写出可靠代码,并知道系统出故障时该检查什么。

文档基准日期:2026-08-12。通用原理不依赖具体版本;涉及产品现状的内容以文末官方资料为准。


目录与学习路线


第 0 章:先看全局——你究竟要学会什么

0.1 一句话总纲

消息队列是一个位于生产者和消费者之间的、能够暂存并转交消息的中间系统。它让“发送方现在产生工作”和“接收方稍后完成工作”不必在同一时间、同一个进程、同一种速度下发生。

这句话包含了整本教程的主线:

  1. 时间解耦:发送方不必等接收方做完,所以可以异步。
  2. 速度解耦:流量来得快、处理得慢时,消息先积压,所以可以削峰填谷。
  3. 空间解耦:发送方不必知道每个接收方的地址和实现,所以可以降低系统耦合。
  4. 代价:一旦“当场调用”变成“以后处理”,你就必须面对消息丢失、重复、乱序、延迟、积压、数据不一致和排障困难。

因此,真正学会 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
2
3
4
5
6
7
用户下单
├─ 创建订单
├─ 扣减库存
├─ 增加积分
├─ 发送短信
├─ 生成物流任务
└─ 写入数据分析系统

我们会不断改造这个案例。每一章都不是孤立的:前一章暴露一个问题,后一章给出机制,再后一章讨论机制带来的新问题。

0.4 学习时牢记的三个层次

遇到任何 MQ 概念,按三个层次理解:

  1. 业务层:它解决什么业务问题?例如下单接口不应等待短信发送。
  2. 语义层:失败时允许丢、允许重、允许乱吗?一致性要到什么程度?
  3. 机制层:具体通过确认、重试、分区、副本、事务、幂等表等什么机制实现?

只背机制容易变成“会配参数但不会设计”;只谈业务容易变成“知道要可靠但不知道怎么做”。后续所有章节都同时回答这三层。


第一部分:为什么世界上需要消息队列

第 1 章:从一个同步下单接口开始

1.1 最直觉的写法

一个初学者很自然会写出:

1
2
3
4
5
6
7
8
9
@app.post("/orders")
async def create_order(request: CreateOrderRequest):
order = await order_db.create(request)
await inventory_service.deduct(order)
await points_service.add(order)
await sms_service.send(order)
await logistics_service.create_task(order)
await analytics_service.record(order)
return {"order_id": order.id}

这叫同步编排:调用者沿调用链等待每一步完成。注意,“同步”在这里描述业务等待关系,不等于 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
2
3
原来:订单服务 ─同步调用→ 短信服务

现在:订单服务 ─发送消息→ 消息队列 ─稍后交付→ 短信消费者

订单服务把“订单已创建”写成一条消息,Broker 先接住并存储;短信服务按自己的速度处理。至此,订单服务不再等待短信服务。

但问题也随之而来:Broker 是什么?消息放在哪里?谁来确认?Broker 宕机会不会丢?这些正是后续章节的内容。


第 2 章:消息队列到底是什么

2.1 用快递柜建立直觉

把 MQ 想象成小区快递柜:

  • 商家是 Producer(生产者),负责寄出包裹。
  • 包裹是 Message(消息)
  • 快递柜及其运营系统是 Broker(消息代理服务器)
  • 柜子的分类区域类似 Queue/Topic(队列/主题)
  • 取件人是 Consumer(消费者)
  • 取件码和取件记录类似消息 ID、确认和消费位点。

商家把包裹放进柜子后,不需要等住户回家;住户晚些时候自行取件。商家和住户不必同时在线,也不需要直接见面。这就是时间和空间解耦。

这个比喻也有边界:真实 MQ 会复制数据、批量传输、重试投递、保留历史、组织消费组;真实消息通常是小型业务事件,不应塞几百 MB 文件。

2.2 严格定义

消息队列通常是一类消息中间件,它提供:

  1. 接收生产者发送的消息;
  2. 按一定规则路由、持久化或缓存;
  3. 将消息交给一个或多个消费者;
  4. 记录消息是否被可靠处理,失败时重投;
  5. 在集群中通过副本、选主等机制提高可用性。

“队列”是通俗总称。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
2
3
4
5
6
7
8
9
10
11
12
13
{
"message_id": "01J...",
"event_type": "order.created",
"event_version": 1,
"occurred_at": "2026-08-12T10:20:30Z",
"aggregate_id": "order-1001",
"trace_id": "trace-abc",
"payload": {
"order_id": "order-1001",
"user_id": "user-9",
"amount": "99.00"
}
}
  • 元数据回答“这是谁、何时发生、如何追踪、按什么路由”。
  • 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
2
请求 → 创建订单(40ms) → 发短信(300ms) → 加积分(70ms) → 返回
总耗时约 410ms

改造后

1
2
3
请求 → 创建订单(40ms) → 可靠写入消息(20ms) → 返回“已受理”

短信/积分稍后消费

响应可能降到 60ms。关键不是“代码并发了”,而是用户请求的完成条件改变了:从“所有副作用完成”改为“核心事务完成且后续工作已被可靠接管”。

异步的代价

  • 用户立即查询时,积分可能还没增加;系统从强一致转向最终一致
  • 错误不再直接返回给原请求,必须通过日志、监控、补偿任务处理。
  • 调试从一条同步调用栈变成跨进程的消息链路。
  • 如果消息只“尽力发送”却未可靠落地,所谓异步可能只是把错误藏起来。

FastAPI BackgroundTasks 是 MQ 吗

不是。它只是在当前 Web 进程返回响应后执行函数:

1
2
3
4
5
6
7
8
from fastapi import BackgroundTasks, FastAPI

app = FastAPI()

@app.post("/notify")
async def notify(background_tasks: BackgroundTasks):
background_tasks.add_task(send_email)
return {"accepted": True}

进程崩溃、容器重启时任务可能消失;它没有独立持久化、跨实例协调、重试和死信。适合不重要的本地轻任务,不适合扣款、发货等关键业务。

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
2
3
4
5
6
订单服务必须知道:
- 各服务地址
- API 参数
- 超时时间
- 重试规则
- 每个服务是否在线

新增订阅者必须改订单服务。

事件解耦

订单服务只发布稳定业务事实 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
Producer → Exchange → Binding → Queue → Consumer

主要类型:

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
2
3
4
Topic: order-events
├─ Partition 0: offset 0,1,2...
├─ Partition 1: offset 0,1,2...
└─ Partition 2: offset 0,1,2...

消息通常按 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
2
Broker:有新消息,给你!
Consumer:收到,开始处理。

优点:使用简单、到达延迟低。难点:如果消费者很慢,Broker 或客户端必须做流量控制,否则消费者会被淹没。

拉模式 Pull

消费者主动向 Broker 请求消息,并决定何时拉、拉多少。

1
2
Consumer:我现在有能力,再给我 100 条。
Broker:这是下一批。

优点:消费者掌握节奏,便于批处理和背压;缺点是空轮询会浪费请求,客户端逻辑相对复杂。

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
2
prefetch 太小:消费者可能吃不饱,吞吐低
prefetch 太大:单个消费者囤积过多消息,内存高、分配不公平、崩溃后重投很多

经验上应结合单条处理耗时、并发数和内存测试,而不是盲抄固定值。

5.5 批量拉取:吞吐的来源之一

网络往返很贵。一次取 100 条通常比请求 100 次每次取 1 条高效。Kafka 的高吞吐很大程度来自顺序追加、页缓存、批处理、压缩和零拷贝等一整套设计,而不是单个“神奇参数”。

批处理也有代价:一批 100 条的第 37 条失败时如何提交进度?整批重试会重复前 36 条;跳过失败项又可能破坏顺序。所以批量越大,吞吐越高,失败语义越需要精心设计。

5.6 如何选择

  • 任务短、希望接口简单、希望 Broker/SDK 管调度:回调式 PushConsumer。
  • 需要批量、精确流控、长耗时或自己管理并发:Pull/SimpleConsumer。
  • 数据流计算:通常拉取并基于 offset 管理进度。

无论选什么,最终都必须回答:“消费者什么时候表示这条消息归我负责,以及什么时候表示已经完成?”这就是下一章的确认和投递语义。


第 6 章:确认、偏移量与投递语义

6.1 先看消息丢失的三个区段

1
生产者业务代码 ──①──> Producer/Broker ──②──> Broker 持久化/副本 ──③──> 消费者业务
  1. 发送端丢失:数据库已提交,但程序在调用 MQ 前崩溃;或者发送超时后直接放弃。
  2. Broker 丢失:消息只在内存或单副本,Broker 宕机;或选主时丢失未同步数据。
  3. 消费端丢失:消费者先确认,后处理;确认后进程崩溃,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
2
ack(message)
await database.save(message) # 此处宕机,消息已被认为完成

正确原则:先让业务副作用成功提交,再 ACK。

1
2
await database.save(message)
ack(message)

但第二段也有窗口:数据库提交成功,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
2
3
4
5
6
7
CREATE TABLE points_ledger (
id BIGSERIAL PRIMARY KEY,
event_id VARCHAR(64) NOT NULL,
user_id VARCHAR(64) NOT NULL,
points INTEGER NOT NULL,
UNIQUE (event_id)
);

同一 event_id 第二次插入触发唯一约束,消费者将它识别为“已处理”并 ACK。

方案二:Inbox / 消费记录表

在与业务更新同一个数据库事务中:

1
2
3
4
5
BEGIN;
INSERT INTO consumer_inbox(message_id, consumer_name)
VALUES ('m-100', 'points-service');
UPDATE user_points SET points = points + 10 WHERE user_id = 'u-1';
COMMIT;

(message_id, consumer_name) 建唯一键。重复消息的 INSERT 失败,整个事务不会再加积分。

方案三:状态机条件更新

1
2
3
UPDATE orders
SET status = 'PAID'
WHERE id = :order_id AND status = 'PENDING_PAYMENT';

检查影响行数:第一次更新 1 行,重复处理更新 0 行。状态迁移必须合法。

Redis SETNX 的陷阱

如果先 SETNX message_id 成功,随后数据库写失败,重试会因为 key 已存在而跳过,反而丢失业务处理。除非 Redis 标记和业务数据具有可靠的原子或补偿设计,否则不要把它当通用幂等答案。数据库唯一约束通常更稳妥。

6.6 Kafka offset 到底是什么

每个 partition 是一个有序日志,offset 是记录在该分区中的位置:

1
2
3
partition-0: [offset 0][offset 1][offset 2][offset 3]...

group-A 下一条读 3

消费位置有两个相关概念:

  • 当前 position:客户端已经拉到或准备读取哪里;
  • committed offset:崩溃重启后从哪里恢复。

提交的通常是“下一条要读的 offset”。处理完 offset 42 后提交 43。

自动提交很方便,但可能在业务处理完成前推进进度。关键业务建议关闭自动提交,完成数据库事务后手动提交。仍然会有“数据库已提交、offset 未提交”的重复窗口,所以照样需要幂等。

6.7 发送超时为什么很棘手

生产者请求超时可能有两种真实情况:

  1. Broker 没收到,消息不存在;
  2. Broker 已写入,但确认响应在网络上丢了。

生产者只看到“超时”,无法判断是哪一种。若重试,第二种情况可能重复;若不重试,第一种情况会丢。因此需要:

  • 生产者幂等机制(如 Kafka idempotent producer);
  • 全局业务 message_id
  • 消费端幂等;
  • 对关键消息使用 Outbox,直到明确发布成功才标记。

6.8 可靠性闭环

一条重要消息至少要做到:

  1. 业务数据和“待发送消息”原子落库;
  2. 发布器持续重试,直到 Broker 确认;
  3. Broker 持久化并按需要复制;
  4. 消费者完成本地事务后确认;
  5. 消费者按 message_id 幂等;
  6. 暂时性失败重试,永久失败进死信;
  7. 对积压、失败率、死信量告警;
  8. 有人工或自动补偿、重放工具。

这比一句“我们用至少一次”具体得多。


第 7 章:分区、消费组、并发与顺序

7.1 为什么需要分区

一条单队列只能由一个线程串行处理,顺序简单但吞吐受限。将 topic 切成多个分区后可以并行:

1
2
3
4
5
Topic 有 4 个分区
Consumer Group 有 2 个实例

C1 ← P0, P1
C2 ← P2, P3

如果增加到 4 个消费者,每个大致分到一个分区;增加到 6 个,其中 2 个通常空闲。因此 Kafka 中一个 consumer group 的有效并行上限通常受分区数限制。

RabbitMQ 的一个队列可向多个消费者分发,模型不同,但也要用 prefetch、消费者并发数和队列数量控制吞吐与隔离。

7.2 分区键决定局部顺序和负载

假设按 order_id 作为 key:

1
partition = hash(order_id) % partition_count

同一订单的创建、支付、取消进入同一分区,可以保持该订单局部顺序;不同订单落到不同分区并行处理。

错误 key 会产生热点。例如以 country=CN 为 key,绝大部分流量进入同一分区。好 key 应同时满足:

  • 需要顺序的数据有相同 key;
  • key 分布足够均匀;
  • key 长期稳定;
  • 不含敏感信息,或经过安全处理。

7.3 全局顺序为什么昂贵

全局顺序要求所有消息经过同一条串行通道,相当于把并行度降为 1。吞吐、可用性和扩展能力都受限。

多数业务并不需要“所有用户所有订单全局有序”,只需要“同一订单的事件有序”。这叫按业务键的局部顺序,是工程上更合理的目标。

7.4 消费组:负载均衡与广播的统一

Kafka 的一个简洁理解:

  • 实例使用同一个 group id:组内竞争消费,分摊分区;
  • 使用不同 group id:每个组各自维护 offset,每组都能读到全部 topic 数据。
1
2
3
4
order-events
├─ inventory-group: I1, I2, I3(组内分摊)
├─ notification-group: N1, N2(组内分摊)
└─ analytics-group: A1, A2(组内分摊)

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 不是垃圾桶,而是异常工单箱。必须配套:

  1. 死信量告警;
  2. 查看原消息、headers、错误、重试次数的工具;
  3. 修复数据或代码后的重放能力;
  4. 重放限速和再次失败的保护;
  5. 敏感数据权限与保留策略。

8.4 毒消息与头部阻塞

一条每次都会失败的消息叫 poison message。若要求严格顺序,它会挡住后面的正常消息,形成头部阻塞。

需要在“顺序”和“可用性”间做业务决策:

  • 严格顺序不可破坏:暂停该 key/分区,告警并人工处理;
  • 可牺牲个别异常:有限重试后送 DLQ,让后续继续;
  • 用业务版本校验和补偿任务修复跳过造成的状态缺口。

8.5 消息积压

积压是生产速度在一段时间内大于消费速度。原因可能是:

  • 流量突增;
  • 下游数据库变慢;
  • 消费者异常退出或频繁 rebalance;
  • 单条消息处理变慢;
  • 分区/队列热点;
  • poison message 重试风暴;
  • 消费端部署了错误版本。

排查顺序:

  1. 看生产速率、消费速率、积压量和最老消息年龄;
  2. 看消费者存活数、错误率、处理耗时;
  3. 看分区是否倾斜;
  4. 看下游数据库、缓存、外部 API;
  5. 决定扩容、限流、降级、批处理或跳过毒消息。

8.6 背压 Backpressure

背压是下游向上游表达“我处理不过来”的机制。

  • RabbitMQ 用 prefetch 限制未确认消息;
  • Kafka 消费者控制 poll、pause/resume 和批量;
  • 应用层可降低生产速率、入口限流或暂时关闭非核心功能;
  • 队列长度/磁盘水位达到阈值时,Broker 也可能阻塞或拒绝发布。

缓冲、背压和限流形成完整闭环:队列吸收短期波动,背压传播压力,限流阻止长期超载。

8.7 重试 Topic 的设计

在 Kafka 这类没有传统“消息退回队头”语义的日志系统中,常见做法是:

1
2
3
orders → 处理失败 → orders.retry.10s
→ 到期转发 → orders
多次失败 → orders.dlq

也可使用调度服务、时间轮、数据库任务表等实现。重点不是名称,而是避免当前分区被单条失败消息无限卡住,并保留重试上下文。


第 9 章:消息契约、序列化与主题设计

9.1 消息不是随便拼的 JSON

生产者和消费者通过消息契约协作。契约至少说明:

  • event_type 的业务含义;
  • 字段名称、类型、单位、是否可空;
  • 哪个字段是幂等键和分区键;
  • 时间是 UTC 还是本地时区;
  • 金额用十进制定点字符串还是最小货币单位整数;
  • 哪些字段包含敏感数据;
  • 兼容性与废弃策略。

9.2 推荐事件信封

1
2
3
4
5
6
7
8
9
10
11
12
13
{
"message_id": "唯一消息 ID",
"event_type": "order.created",
"event_version": 1,
"source": "order-service",
"occurred_at": "2026-08-12T10:20:30.123Z",
"aggregate_type": "order",
"aggregate_id": "order-1001",
"trace_id": "跨服务追踪 ID",
"correlation_id": "一条业务流程的关联 ID",
"causation_id": "导致本事件的上游消息 ID",
"payload": {}
}

message_id 用于消息级去重;aggregate_id 常用于局部顺序;trace_id 用于观察链路;causation_id 能解释“这条消息由谁触发”。

9.3 事件通知与事件携带状态

只带 ID

1
{"event_type": "order.created", "order_id": "o-1"}

消费者再查订单服务。优点是消息小、数据来源单一;缺点是产生同步依赖、查询压力,回放旧消息时查到的可能是新状态。

携带需要的数据

1
2
3
4
{
"event_type": "order.created",
"payload": {"order_id": "o-1", "user_id": "u-1", "amount_cent": 9900}
}

消费者可独立处理和回放,但契约更重,可能复制数据。实践中应携带消费者完成该事件所需的稳定最小快照,不要塞整个数据库对象。

9.4 向后兼容

生产者新增可选字段通常较安全;删除字段、改类型、改语义最危险。

建议:

  • 消费者忽略未知字段;
  • 新字段先设为可选并提供默认语义;
  • 不复用旧字段表达新含义;
  • 重大不兼容变化创建 v2 契约或新 topic;
  • 使用 JSON Schema、Avro、Protobuf 等进行契约校验;
  • 在 CI 中做生产者—消费者契约测试。

9.5 Topic / Queue 如何划分

过粗:所有事件都放 events,权限、保留期、容量和订阅困难。

过细:每种动作一个 topic,数量爆炸,运维困难。

常见维度:

  • 按业务域:order-eventspayment-events
  • 按数据敏感级别和权限隔离;
  • 按保留期、吞吐和可靠性要求隔离;
  • 按消息类型限制隔离(RocketMQ 5.x 的 FIFO、Delay、Transaction topic 有明确类型约束);
  • 事件类型可放 header 或 payload,不一定每种事件一个 topic。

9.6 大消息怎么处理

MQ 不适合传大文件。推荐 Claim Check 模式:

  1. 文件写对象存储;
  2. 消息只包含对象 URI、校验和、大小和权限信息;
  3. 消费者按引用下载;
  4. 设置对象生命周期,避免消息还没消费文件已删除。

大消息会降低吞吐、放大重试成本、增加内存与网络压力,还可能触发 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
2
3
binding: order.created
message: order.created → 匹配
message: order.paid → 不匹配

Fanout

忽略 routing key,发给所有绑定队列。适合配置广播、临时通知、日志广播。

Topic

按点分词并支持通配符:* 匹配一个词,# 匹配零个或多个词。

1
2
3
order.*       匹配 order.created,不匹配 order.cn.created
order.# 两者都匹配
*.failed 匹配 payment.failed

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
2
3
4
5
6
7
8
9
work.queue
│ 处理失败,拒绝且不回原队列

retry.exchange → retry.10s.queue(TTL 10 秒,无消费者)
│ 过期后通过 DLX

work.exchange → work.queue

超过次数 → dead.exchange → work.dlq

注意:用 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
2
3
Topic: payment-events
Partition 0:
[0:create][1:authorized][2:captured][3:refunded]...

这使 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
2
3
4
5
6
7
读取 input topic

处理

在同一 Kafka 事务中:
- 写 output topic
- 提交 input offsets

下游配置 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
2
order-1: create → pay → ship
order-2: create → cancel

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 概览页面的语言列表更新可能滞后。

实际采用前必须做兼容性验证:

  1. 目标 Broker/Proxy 版本;
  2. Python SDK 的发布包、成熟度和平台支持;
  3. 事务回查、FIFO、TLS、认证、可观测性是否齐全;
  4. 团队能否得到稳定运维支持。

学习中可以掌握其架构与语义;如果当前 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”是错误推理。技术选型至少评估:

  1. 业务语义:任务、事件流、事务、延迟、顺序、广播;
  2. 吞吐、峰值、消息大小、保留期、可接受延迟;
  3. 可丢/可重/可乱的边界;
  4. 客户端语言和功能成熟度;
  5. 团队已有平台和运维经验;
  6. 云托管服务、成本、跨地域和合规;
  7. 监控、重放、审计和故障恢复能力;
  8. 用代表性负载做 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
2
await db.commit_order()
await mq.publish("order.created")

无论先后都有窗口。

先提交数据库,再发消息

数据库成功后进程崩溃,消息没发:订单存在,但库存/积分永远不知道。

先发消息,再提交数据库

消息发出后数据库失败:消费者看到一个根本不存在或已回滚的订单。

仅仅交换两行代码无法解决原子性。

14.2 先区分强一致与最终一致

  • 强一致:事务完成的任意可见时刻,各系统状态满足一致规则。
  • 最终一致:短时间允许不一致,只要没有新的失败,系统最终收敛到正确状态。

异步消息通常追求最终一致。用户可能先看到“订单已创建,积分处理中”,随后积分到账。关键是状态要透明、可追踪、可补偿,不能把暂时不一致伪装成已完成。

14.3 方案一:Transactional Outbox(通用首选)

在订单数据库中增加 outbox 表。本地事务同时写订单和待发消息:

1
2
3
4
5
6
7
8
9
10
11
12
BEGIN;

INSERT INTO orders(id, user_id, status, amount_cent)
VALUES ('o-1', 'u-1', 'CREATED', 9900);

INSERT INTO outbox_events(
event_id, aggregate_id, event_type, payload, status
) VALUES (
'm-1', 'o-1', 'order.created', '{...}', 'PENDING'
);

COMMIT;

因为两次写在同一个数据库本地事务中,它们要么一起成功,要么一起失败。独立 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 两种实现

  1. 轮询发布:定时 SELECT ... FOR UPDATE SKIP LOCKED 取待发布行。简单通用,但有轮询延迟和数据库压力。
  2. CDC 发布:Debezium 等读取数据库事务日志,把 outbox 变化写到 Kafka。吞吐和解耦更好,基础设施更复杂。

14.4 方案二:RocketMQ 事务消息

RocketMQ 事务消息把消息先存成消费者不可见的半消息,再执行本地事务,最后 commit/rollback。若最终确认丢失,Broker 回查生产者的本地事务状态。

生产者必须保存可查询的事务状态,例如通过订单表判断:

1
2
3
订单 o-1 已存在且状态有效 → COMMIT
明确不存在且不可能仍在执行 → ROLLBACK
本地事务仍进行中/暂不能判断 → UNKNOWN

关键限制:

  • 它实现上游本地事务与消息可见性的最终一致,不是把所有下游数据库变成一个 ACID 大事务;
  • 下游仍需幂等和重试;
  • 回查函数必须快速、稳定、可重复;
  • 不应长时间返回 unknown 或制造大量半消息;
  • 只适合可以接受异步最终一致的场景。

14.5 方案三:Kafka 事务

如果输入和输出都在 Kafka:

1
消费 topic-A → 处理 → 写 topic-B + 提交 topic-A offset

可以放入同一个 Kafka 事务,实现 Kafka 边界内的原子 consume-transform-produce。它非常适合流处理。

如果还要写普通业务数据库,常见选择仍是:

  • 数据库幂等更新 + 手动提交 offset;
  • 把 offset 与业务数据存入同一数据库事务(自行管理恢复);
  • Outbox/CDC;
  • 业务状态机和补偿。

14.6 Saga 与补偿

跨多个长事务时,不一定能“回滚一切”,而是执行补偿动作:

1
2
3
4
5
6
创建订单成功
→ 预留库存成功
→ 支付失败
→ 发布 payment.failed
→ 释放库存
→ 取消订单

补偿不是数据库回滚的完美逆操作:退款可能有手续费,短信无法“撤回”,物流可能已出库。因此 Saga 要明确每步的正向动作、补偿动作、幂等键、超时和人工介入状态。

14.7 状态机比“几个布尔值”可靠

不要设计:

1
is_paid, is_cancelled, is_shipped

它可能出现同时 is_cancelled=trueis_shipped=true 的矛盾。使用明确状态机:

1
2
CREATED → PENDING_PAYMENT → PAID → FULFILLING → SHIPPED
└────────→ CANCELLED

每个消息只允许合法迁移,并使用版本号/条件更新防止旧事件覆盖新状态。

14.8 对账是最后一道防线

任何复杂分布式系统都应接受“在线机制可能有漏网异常”。定期对账:

  • 订单已支付但账务无流水;
  • outbox 长期 PENDING;
  • 发货状态与物流状态不一致;
  • 支付网关记录与本地记录不一致。

检测后自动补发、补偿或生成人工工单。可靠系统不是“绝不出错”,而是“错误可发现、可恢复、可审计”。


第 15 章:顺序消息

15.1 先问:要保证谁和谁的顺序

“保证顺序”是不完整需求。必须说清:

  • 同一个生产者的发送顺序?
  • Broker 存储顺序?
  • 消费者收到顺序?
  • 业务提交完成顺序?
  • 全部消息全局顺序,还是同一订单局部顺序?

真正业务通常需要:同一聚合根(order_id/user_id/account_id)的状态变更有序。

15.2 端到端有序的条件

以订单事件为例:

  1. 生产者按业务顺序发;
  2. order_id 稳定映射到同一 partition/message group/queue;
  3. Broker 在该顺序域内有序存储和交付;
  4. 消费端不无序并发;
  5. 失败重试不能悄悄让后续越过;
  6. 业务数据库用状态机/版本号防止旧写覆盖新写。

只满足第 3 条并不能叫业务有序。

15.3 Kafka 实现

1
2
3
4
5
await producer.send_and_wait(
"order-events",
key=order_id.encode(),
value=payload,
)

同一 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
{"order_id": "o-1", "status": "PAID", "aggregate_version": 3}

消费者保存已应用版本:

  • 收到 version = current + 1:正常应用;
  • 收到 version <= current:重复或旧事件,忽略;
  • 收到 version > current + 1:发现缺口,暂停该 key、重试或回源修复。

这样即使基础设施发生意外乱序,业务也不会默默写坏数据。


第 16 章:延时消息与定时任务

16.1 两类时间需求

  • 延时:从现在起 30 分钟后处理。
  • 定时:在 2026-08-13 09:00:00 处理。

业务语义还应说明:允许晚多久、能否取消、能否修改时间、是否只执行一次、时区是什么。

16.2 订单超时取消的竞态

下单时发送“30 分钟后检查”消息。到期时不能直接取消,必须再次读取订单状态:

1
2
3
order = await repo.get(order_id)
if order.status == "PENDING_PAYMENT":
await repo.cancel_if_still_pending(order_id)

用户可能在第 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 一个聊天消息的旅程

  1. 客户端生成 client_msg_id 并发送给 Gateway;
  2. 消息服务鉴权,分配服务器序列号并持久化;
  3. 服务端向发送者返回“已接收”;
  4. 发布 chat.message.created 到 MQ;
  5. 目标用户所在 Gateway 订阅并通过 WebSocket 推送;
  6. 用户不在线则写离线队列/推送服务;
  7. 接收端 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
page_view → add_to_cart → order_created → payment_succeeded

推荐系统、风控、实时大屏、数仓都想读取这条流,但处理方式和进度不同。

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
2
3
4
5
6
(user-1, level=1)
(user-1, level=2)
(user-2, level=1)
↓ 按 key 聚合最新值
user-1 → level 2
user-2 → level 1

Kafka compacted topic 常用来支持这种可恢复状态。

18.5 数据流中的恰好一次仍有边界

Flink/Kafka Streams 可以把状态快照、输入进度和输出协调起来,提供很强的处理语义。但外部 sink 是否幂等、是否支持事务仍影响端到端结果。例如向第三方 HTTP API 发请求,checkpoint 无法撤回已经发出的请求,仍需业务幂等键。


第五部分:Python + FastAPI 完整项目

第 19 章:项目设计——可靠下单系统

19.1 项目目标

我们实现一个最小但可靠的订单—积分链路:

  1. FastAPI 接收创建订单请求;
  2. PostgreSQL 同一事务写 ordersoutbox_events
  3. Outbox Relay 发布 order.created 到 RabbitMQ;
  4. 积分消费者收到事件,在同一数据库事务中写 Inbox 去重记录并增加积分;
  5. 业务成功后才 ACK;
  6. 暂时失败进入延迟重试,达到上限进入 DLQ;
  7. 可以故意杀进程,观察重复但不重复加积分。

这不是“生产万能模板”,但覆盖了真正重要的可靠性闭环。

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
2
3
4
5
6
7
8
9
10
11
12
mq_order_lab/
├── compose.yaml
├── requirements.txt
└── app/
├── __init__.py
├── config.py
├── db.py
├── models.py
├── messaging.py
├── main.py
├── outbox_worker.py
└── points_worker.py

19.5 可靠性不变量

写代码前先写系统不变量:

  • 订单存在,则对应 outbox 事件一定存在;
  • outbox 发布可能重复,但不能永久静默丢失;
  • 同一 message_id 对积分服务最多产生一次业务效果;
  • 数据库事务成功后才 ACK;
  • 达到最大重试次数的消息不会无限循环,而是进入 DLQ;
  • 所有消息带 message_idorder_idtrace_id 和版本。

如果代码不能证明这些不变量,就不能仅凭“测试发了一条成功”声称可靠。


第 20 章:RabbitMQ 项目代码

下面按文件给出可组装的教学项目。为了聚焦 MQ,开发环境用 SQLAlchemy create_all;正式环境请使用 Alembic 迁移、密钥管理、TLS、连接池和完整监控。

20.1 依赖

requirements.txt

1
2
3
4
5
6
fastapi
uvicorn[standard]
sqlalchemy>=2
asyncpg
aio-pika
pydantic-settings

20.2 本地基础设施

compose.yaml

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
services:
postgres:
image: postgres:16
environment:
POSTGRES_USER: app
POSTGRES_PASSWORD: app
POSTGRES_DB: orders
ports:
- "5432:5432"
healthcheck:
test: ["CMD-SHELL", "pg_isready -U app -d orders"]
interval: 5s
timeout: 3s
retries: 20
volumes:
- pg_data:/var/lib/postgresql/data

rabbitmq:
image: rabbitmq:4-management
environment:
RABBITMQ_DEFAULT_USER: app
RABBITMQ_DEFAULT_PASS: app
ports:
- "5672:5672"
- "15672:15672"
healthcheck:
test: ["CMD", "rabbitmq-diagnostics", "-q", "ping"]
interval: 5s
timeout: 5s
retries: 20
volumes:
- rabbit_data:/var/lib/rabbitmq

volumes:
pg_data:
rabbit_data:

单节点容器只用于学习,quorum queue 在单节点上没有真正的节点级高可用。生产要部署奇数副本 RabbitMQ 集群并按故障域设计。

20.3 配置

app/config.py

1
2
3
4
5
6
7
8
9
10
11
12
13
14
from pydantic_settings import BaseSettings, SettingsConfigDict


class Settings(BaseSettings):
database_url: str = "postgresql+asyncpg://app:app@localhost:5432/orders"
amqp_url: str = "amqp://app:app@localhost:5672/"
exchange_name: str = "domain.events"
retry_exchange_name: str = "domain.retry"
dead_exchange_name: str = "domain.dead"

model_config = SettingsConfigDict(env_file=".env", extra="ignore")


settings = Settings()

20.4 数据库连接

app/db.py

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
from sqlalchemy.ext.asyncio import (
AsyncSession,
async_sessionmaker,
create_async_engine,
)

from .config import settings


engine = create_async_engine(
settings.database_url,
pool_pre_ping=True,
)

SessionFactory = async_sessionmaker(
engine,
class_=AsyncSession,
expire_on_commit=False,
)

20.5 数据模型

app/models.py

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
import uuid
from datetime import datetime, timezone

from sqlalchemy import DateTime, Integer, JSON, String, Text, UniqueConstraint
from sqlalchemy.orm import DeclarativeBase, Mapped, mapped_column


def utcnow() -> datetime:
return datetime.now(timezone.utc)


class Base(DeclarativeBase):
pass


class Order(Base):
__tablename__ = "orders"

id: Mapped[str] = mapped_column(String(36), primary_key=True)
user_id: Mapped[str] = mapped_column(String(64), index=True)
amount_cent: Mapped[int] = mapped_column(Integer)
status: Mapped[str] = mapped_column(String(32), default="CREATED")
created_at: Mapped[datetime] = mapped_column(
DateTime(timezone=True), default=utcnow
)


class OutboxEvent(Base):
__tablename__ = "outbox_events"

event_id: Mapped[str] = mapped_column(String(36), primary_key=True)
aggregate_id: Mapped[str] = mapped_column(String(36), index=True)
event_type: Mapped[str] = mapped_column(String(100), index=True)
payload: Mapped[dict] = mapped_column(JSON)
status: Mapped[str] = mapped_column(String(20), default="PENDING", index=True)
attempts: Mapped[int] = mapped_column(Integer, default=0)
last_error: Mapped[str | None] = mapped_column(Text, nullable=True)
created_at: Mapped[datetime] = mapped_column(
DateTime(timezone=True), default=utcnow
)
published_at: Mapped[datetime | None] = mapped_column(
DateTime(timezone=True), nullable=True
)


class ConsumerInbox(Base):
__tablename__ = "consumer_inbox"
__table_args__ = (
UniqueConstraint("message_id", "consumer_name", name="uq_inbox_message"),
)

id: Mapped[int] = mapped_column(primary_key=True, autoincrement=True)
message_id: Mapped[str] = mapped_column(String(64))
consumer_name: Mapped[str] = mapped_column(String(100))
processed_at: Mapped[datetime] = mapped_column(
DateTime(timezone=True), default=utcnow
)


class PointsAccount(Base):
__tablename__ = "points_accounts"

user_id: Mapped[str] = mapped_column(String(64), primary_key=True)
points: Mapped[int] = mapped_column(Integer, default=0)

20.6 RabbitMQ 封装和拓扑

app/messaging.py

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
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
import json

import aio_pika
from aio_pika import DeliveryMode, ExchangeType, Message

from .config import settings


class RabbitPublisher:
def __init__(self) -> None:
self.connection: aio_pika.RobustConnection | None = None
self.channel: aio_pika.abc.AbstractRobustChannel | None = None
self.events: aio_pika.abc.AbstractRobustExchange | None = None
self.retry: aio_pika.abc.AbstractRobustExchange | None = None
self.dead: aio_pika.abc.AbstractRobustExchange | None = None

async def connect(self) -> None:
self.connection = await aio_pika.connect_robust(settings.amqp_url)
self.channel = await self.connection.channel(
publisher_confirms=True,
on_return_raises=True,
)
await self.channel.set_qos(prefetch_count=20)

self.events = await self.channel.declare_exchange(
settings.exchange_name, ExchangeType.TOPIC, durable=True
)
self.retry = await self.channel.declare_exchange(
settings.retry_exchange_name, ExchangeType.DIRECT, durable=True
)
self.dead = await self.channel.declare_exchange(
settings.dead_exchange_name, ExchangeType.DIRECT, durable=True
)

points_q = await self.channel.declare_queue(
"points.order-created.q",
durable=True,
arguments={"x-queue-type": "quorum"},
)
await points_q.bind(self.events, routing_key="order.created")

retry_q = await self.channel.declare_queue(
"points.order-created.retry.10s.q",
durable=True,
arguments={
"x-queue-type": "quorum",
"x-message-ttl": 10_000,
"x-dead-letter-exchange": settings.exchange_name,
"x-dead-letter-routing-key": "order.created",
},
)
await retry_q.bind(self.retry, routing_key="points.retry.10s")

dead_q = await self.channel.declare_queue(
"points.order-created.dlq",
durable=True,
arguments={"x-queue-type": "quorum"},
)
await dead_q.bind(self.dead, routing_key="points.dead")

async def publish_event(
self,
*,
routing_key: str,
body: dict,
message_id: str,
headers: dict | None = None,
) -> None:
assert self.events is not None
await self.events.publish(
Message(
body=json.dumps(body, ensure_ascii=False).encode(),
content_type="application/json",
delivery_mode=DeliveryMode.PERSISTENT,
message_id=message_id,
headers=headers or {},
),
routing_key=routing_key,
mandatory=True,
)

async def publish_retry(
self, *, body: bytes, message_id: str, retry_count: int
) -> None:
assert self.retry is not None
await self.retry.publish(
Message(
body=body,
content_type="application/json",
delivery_mode=DeliveryMode.PERSISTENT,
message_id=message_id,
headers={"x-retry-count": retry_count},
),
routing_key="points.retry.10s",
mandatory=True,
)

async def publish_dead(
self, *, body: bytes, message_id: str, retry_count: int, error: str
) -> None:
assert self.dead is not None
await self.dead.publish(
Message(
body=body,
content_type="application/json",
delivery_mode=DeliveryMode.PERSISTENT,
message_id=message_id,
headers={
"x-retry-count": retry_count,
"x-last-error": error[:500],
},
),
routing_key="points.dead",
mandatory=True,
)

async def close(self) -> None:
if self.connection is not None:
await self.connection.close()

这里有几个故意体现可靠性的点:

  • connect_robust 负责连接恢复,但恢复期间的业务语义仍要靠重试和幂等;
  • publisher_confirms=True
  • mandatory=True,未路由消息不会被悄悄忽略;
  • 消息持久化;
  • 队列为 durable quorum queue;
  • 重试发布得到 confirm 后,消费者才能 ACK 原消息。

20.7 FastAPI:订单和 Outbox 同事务

app/main.py

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
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
import uuid
from contextlib import asynccontextmanager
from datetime import datetime, timezone

from fastapi import FastAPI, HTTPException, status
from pydantic import BaseModel, Field
from sqlalchemy import select

from .db import SessionFactory, engine
from .models import Base, Order, OutboxEvent, PointsAccount


class CreateOrderRequest(BaseModel):
user_id: str = Field(min_length=1, max_length=64)
amount_cent: int = Field(gt=0)


@asynccontextmanager
async def lifespan(_: FastAPI):
# 仅供教学。生产环境改用 Alembic migration。
async with engine.begin() as conn:
await conn.run_sync(Base.metadata.create_all)
yield
await engine.dispose()


app = FastAPI(title="Reliable MQ Order Lab", lifespan=lifespan)


@app.post("/orders", status_code=status.HTTP_202_ACCEPTED)
async def create_order(request: CreateOrderRequest):
order_id = str(uuid.uuid4())
event_id = str(uuid.uuid4())
trace_id = str(uuid.uuid4())
now = datetime.now(timezone.utc)

event = {
"message_id": event_id,
"event_type": "order.created",
"event_version": 1,
"source": "order-service",
"occurred_at": now.isoformat(),
"aggregate_type": "order",
"aggregate_id": order_id,
"trace_id": trace_id,
"payload": {
"order_id": order_id,
"user_id": request.user_id,
"amount_cent": request.amount_cent,
},
}

async with SessionFactory() as session:
async with session.begin():
session.add(
Order(
id=order_id,
user_id=request.user_id,
amount_cent=request.amount_cent,
status="CREATED",
)
)
session.add(
OutboxEvent(
event_id=event_id,
aggregate_id=order_id,
event_type="order.created",
payload=event,
status="PENDING",
)
)

return {
"order_id": order_id,
"event_id": event_id,
"status": "CREATED",
"downstream_status": "PENDING",
}


@app.get("/users/{user_id}/points")
async def get_points(user_id: str):
async with SessionFactory() as session:
account = await session.get(PointsAccount, user_id)
if account is None:
return {"user_id": user_id, "points": 0}
return {"user_id": user_id, "points": account.points}


@app.get("/orders/{order_id}")
async def get_order(order_id: str):
async with SessionFactory() as session:
result = await session.execute(select(Order).where(Order.id == order_id))
order = result.scalar_one_or_none()
if order is None:
raise HTTPException(404, "order not found")
return {
"id": order.id,
"user_id": order.user_id,
"amount_cent": order.amount_cent,
"status": order.status,
}

注意:API 中没有 publisher.publish()。这是正确的——本地事务与 Outbox 已经把“订单存在”和“事件等待发送”绑定起来。

20.8 Outbox Relay

app/outbox_worker.py

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
65
import asyncio
from datetime import datetime, timezone

from sqlalchemy import select

from .db import SessionFactory
from .messaging import RabbitPublisher
from .models import OutboxEvent


async def publish_one(publisher: RabbitPublisher) -> bool:
publish_error: Exception | None = None
async with SessionFactory() as session:
async with session.begin():
result = await session.execute(
select(OutboxEvent)
.where(OutboxEvent.status == "PENDING")
.order_by(OutboxEvent.created_at)
.limit(1)
.with_for_update(skip_locked=True)
)
event = result.scalar_one_or_none()
if event is None:
return False

try:
await publisher.publish_event(
routing_key=event.event_type,
body=event.payload,
message_id=event.event_id,
)
# confirm 成功后才在同一事务中标已发布。
event.status = "PUBLISHED"
event.published_at = datetime.now(timezone.utc)
event.attempts += 1
event.last_error = None
except Exception as exc:
event.attempts += 1
event.last_error = repr(exc)[:2000]
# 保持 PENDING,下一轮继续重试。
# 先让 attempts/last_error 提交,再在事务外抛出。
publish_error = exc
if publish_error is not None:
raise publish_error
return True


async def main() -> None:
publisher = RabbitPublisher()
await publisher.connect()
try:
while True:
try:
found = await publish_one(publisher)
if not found:
await asyncio.sleep(0.5)
except Exception as exc:
print("outbox publish failed:", repr(exc))
await asyncio.sleep(2)
finally:
await publisher.close()


if __name__ == "__main__":
asyncio.run(main())

这个教学实现为了让原理清晰,在网络发布期间持有一条数据库行锁。高吞吐生产版本应使用批量 claim:快速把一批行标记 PROCESSING 并提交,再发布;同时加入 claimed_at 租约和超时回收器,避免 worker 崩溃后永久卡在 PROCESSING。无论哪种实现,都必须接受“confirm 后、状态提交前崩溃”造成重复发布。

20.9 积分消费者:Inbox 幂等 + 手动 ACK

app/points_worker.py

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
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
import asyncio
import json

import aio_pika
from sqlalchemy.dialects.postgresql import insert as pg_insert

from .config import settings
from .db import SessionFactory
from .messaging import RabbitPublisher
from .models import ConsumerInbox, PointsAccount

CONSUMER_NAME = "points-service"
MAX_RETRIES = 5


async def apply_points_once(event: dict) -> bool:
"""返回 True 表示首次处理;False 表示重复消息。"""
message_id = event["message_id"]
user_id = event["payload"]["user_id"]

async with SessionFactory() as session:
async with session.begin():
statement = (
pg_insert(ConsumerInbox)
.values(message_id=message_id, consumer_name=CONSUMER_NAME)
.on_conflict_do_nothing(
index_elements=["message_id", "consumer_name"]
)
.returning(ConsumerInbox.id)
)
inbox_id = (await session.execute(statement)).scalar_one_or_none()
if inbox_id is None:
return False

account = await session.get(
PointsAccount, user_id, with_for_update=True
)
if account is None:
account = PointsAccount(user_id=user_id, points=0)
session.add(account)
account.points += 10
return True


async def main() -> None:
publisher = RabbitPublisher()
await publisher.connect()
assert publisher.channel is not None

queue = await publisher.channel.get_queue("points.order-created.q")

async def on_message(message: aio_pika.IncomingMessage) -> None:
message_id = message.message_id
if not message_id:
# 无稳定 ID 无法安全幂等,视为不可恢复格式错误。
await publisher.publish_dead(
body=message.body,
message_id="missing-message-id",
retry_count=0,
error="missing message_id",
)
await message.ack()
return

try:
event = json.loads(message.body)
first_time = await apply_points_once(event)
print("processed" if first_time else "duplicate ignored", message_id)
# 数据库事务已经提交,才 ACK。
await message.ack()
except (json.JSONDecodeError, KeyError, TypeError) as exc:
# 格式/契约错误不会因等待而变好,直接死信。
await publisher.publish_dead(
body=message.body,
message_id=message_id,
retry_count=int(message.headers.get("x-retry-count", 0)),
error=f"permanent error: {exc!r}",
)
await message.ack()
except Exception as exc:
retry_count = int(message.headers.get("x-retry-count", 0)) + 1
try:
if retry_count <= MAX_RETRIES:
await publisher.publish_retry(
body=message.body,
message_id=message_id,
retry_count=retry_count,
)
else:
await publisher.publish_dead(
body=message.body,
message_id=message_id,
retry_count=retry_count,
error=repr(exc),
)
# 重试/DLQ 发布已 confirm,才移除原消息。
await message.ack()
except Exception:
# 连重试消息也发布失败,原消息必须保留。
await message.nack(requeue=True)

await queue.consume(on_message, no_ack=False)
print("points worker is running")
try:
await asyncio.Future()
finally:
await publisher.close()


if __name__ == "__main__":
asyncio.run(main())

一个并发细节

两个重复消息若第一次同时处理,同一用户账户“查不到后各自 INSERT”可能产生主键竞争。数据库会让其中一个事务失败并重试,Inbox 唯一键保证同一 message_id 不产生两次效果。生产中可预创建积分账户,或使用 PostgreSQL 原子 upsert:

1
2
3
4
INSERT INTO points_accounts(user_id, points)
VALUES (:user_id, 10)
ON CONFLICT (user_id)
DO UPDATE SET points = points_accounts.points + 10;

必须与 Inbox 插入放在同一事务中。

20.10 启动与验证

在项目目录依次运行:

1
2
3
4
5
docker compose up -d
python -m venv .venv
.\.venv\Scripts\Activate.ps1
pip install -r requirements.txt
uvicorn app.main:app --reload

另开两个终端:

1
python -m app.outbox_worker
1
python -m app.points_worker

创建订单:

1
2
3
4
5
Invoke-RestMethod `
-Method Post `
-Uri http://127.0.0.1:8000/orders `
-ContentType application/json `
-Body '{"user_id":"u-100","amount_cent":9900}'

查询积分:

1
Invoke-RestMethod http://127.0.0.1:8000/users/u-100/points

RabbitMQ 管理界面是 http://localhost:15672,用户名和密码均为 app。你应该能看到工作队列、重试队列、死信队列以及 ready/unacked 数量。


第 21 章:故障实验与验收

“能跑”不等于“可靠”。下面用故障主动验证不变量。

21.1 RabbitMQ 停机时创建订单

  1. 停止 RabbitMQ;
  2. 调用创建订单 API;
  3. API 应仍返回订单创建成功;
  4. 数据库中 outbox 为 PENDING
  5. 重启 RabbitMQ;
  6. 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
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
import json

from aiokafka import AIOKafkaProducer


class KafkaPublisher:
def __init__(self, bootstrap_servers: str):
self.producer = AIOKafkaProducer(
bootstrap_servers=bootstrap_servers,
acks="all",
enable_idempotence=True,
value_serializer=lambda value: json.dumps(value).encode(),
)

async def start(self):
await self.producer.start()

async def publish(self, event: dict):
await self.producer.send_and_wait(
"order-events",
key=event["aggregate_id"].encode(),
value=event,
headers=[
("event_type", event["event_type"].encode()),
("message_id", event["message_id"].encode()),
],
)

async def stop(self):
await self.producer.stop()

Outbox Relay 只需把 Rabbit publisher 替换为这个 adapter。Outbox 仍然必要,因为 Kafka 幂等生产者解决的是 Kafka 协议重试重复,不解决“PostgreSQL 已提交但程序还没调用 producer 就崩溃”。

22.2 手动提交的幂等消费者

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
import asyncio
import json

from aiokafka import AIOKafkaConsumer, TopicPartition


async def run_consumer():
consumer = AIOKafkaConsumer(
"order-events",
bootstrap_servers="localhost:9092",
group_id="points-service-v1",
enable_auto_commit=False,
auto_offset_reset="earliest",
value_deserializer=json.loads,
)
await consumer.start()
try:
async for message in consumer:
tp = TopicPartition(message.topic, message.partition)
try:
# 与 RabbitMQ 版共用 Inbox + 积分事务。
await apply_points_once(message.value)
# 提交“下一条要读的 offset”。
await consumer.commit({tp: message.offset + 1})
except Exception as exc:
print("processing failed", exc)
# 教学版让同一分区回到失败消息,避免越过它。
consumer.seek(tp, message.offset)
await asyncio.sleep(2)
finally:
await consumer.stop()

生产版还要:

  • 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。

你应该亲自实验:

  1. topic 创建 3 partitions;
  2. points group 启 1、2、4 个实例,观察分配;
  3. 停掉一个实例,观察 rebalance;
  4. 将 group offset reset 到 earlier,重新消费历史;
  5. 使用相同 order_id 发多条事件,观察同 partition 顺序;
  6. 不带 key 发送,观察顺序域变化。

22.4 Kafka 内事务示例

若消费 raw-orders 并把清洗结果写 clean-orders,可用:

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
consumer = AIOKafkaConsumer(
"raw-orders",
bootstrap_servers="localhost:9092",
group_id="cleaner-v1",
enable_auto_commit=False,
isolation_level="read_committed",
)

producer = AIOKafkaProducer(
bootstrap_servers="localhost:9092",
transactional_id="cleaner-shard-0",
)

# 启动 producer/consumer 后:
batch = await consumer.getmany(timeout_ms=1000, max_records=100)
async with producer.transaction():
offsets = {}
for tp, messages in batch.items():
for message in messages:
output = transform(message.value)
await producer.send("clean-orders", output, key=message.key)
if messages:
offsets[tp] = messages[-1].offset + 1
await producer.send_offsets_to_transaction(offsets, "cleaner-v1")

每个并行 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
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
from rocketmq import ClientConfiguration, Credentials, Message, Producer

endpoints = "127.0.0.1:8081" # RocketMQ Proxy endpoint,按实际环境修改
topic = "order-events"

config = ClientConfiguration(endpoints, Credentials())
producer = Producer(config, (topic,))
producer.startup()
try:
message = Message()
message.topic = topic
message.tag = "order-created"
message.keys = "order-o-100"
message.body = b'{"message_id":"m-100","order_id":"o-100"}'
receipt = producer.send(message)
print(receipt)
finally:
producer.shutdown()

生产中仍应由 Outbox Relay 调用 producer,而不是在数据库提交后“顺手发一下”。

23.2 FIFO 消息

1
2
3
4
5
6
7
for state in ["CREATED", "PAID", "SHIPPED"]:
message = Message()
message.topic = "order-fifo-events"
message.tag = "order-state-changed"
message.message_group = "order-o-100" # 同一订单的顺序域
message.body = state.encode()
producer.send(message)

topic 必须按服务端要求创建为 FIFO 类型。消费端也必须使用 FIFO listener/配置并同步完成处理。

23.3 延迟消息

1
2
3
4
5
6
7
8
9
10
11
12
13
import time

message = Message()
message.topic = "order-timeout-events"
message.tag = "check-unpaid-order"
message.keys = "order-o-100"
message.body = b'{"order_id":"o-100"}'

# 当前 Apache Python 仓库示例使用 epoch seconds;
# RocketMQ 概念文档常以毫秒时间戳描述协议字段。
# 务必以你所用 SDK 版本的字段单位为准,做一次 10 秒集成测试。
message.delivery_timestamp = int(time.time()) + 30 * 60
producer.send(message)

到期消费者仍要执行 cancel_if_status_is_pending 条件更新,不能无条件取消。

23.4 SimpleConsumer:处理后再 ACK

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
from rocketmq import (
ClientConfiguration,
Credentials,
FilterExpression,
SimpleConsumer,
)

config = ClientConfiguration("127.0.0.1:8081", Credentials())
consumer = SimpleConsumer(
config,
"points-service-v1",
{"order-events": FilterExpression("order-created")},
)
consumer.startup()
try:
while True:
# 当前 SDK 示例:最多 32 条;不可见时长参数按该 SDK 单位核验。
messages = consumer.receive(32, 15)
for message in messages or []:
try:
event = decode(message.body)
apply_points_idempotently(event)
consumer.ack(message) # 业务事务成功后再确认
except Exception as exc:
print("leave unacked for retry", exc)
finally:
consumer.shutdown()

不可见时间必须大于正常处理耗时,并处理消费超时导致的重复。如果业务时长不可预测,需要续期/修改不可见时间的相应 SDK 能力或调整消费模式。

23.5 事务消息代码骨架

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
from rocketmq import (
ClientConfiguration,
Credentials,
Message,
Producer,
TransactionChecker,
TransactionResolution,
)


class OrderTransactionChecker(TransactionChecker):
def check(self, message: Message) -> TransactionResolution:
order_id = message.keys
state = query_local_transaction_state(order_id)
if state == "COMMITTED":
return TransactionResolution.COMMIT
if state == "ROLLED_BACK":
return TransactionResolution.ROLLBACK
return TransactionResolution.UNKNOWN


config = ClientConfiguration("127.0.0.1:8081", Credentials())
producer = Producer(
config,
("order-transaction-events",),
checker=OrderTransactionChecker(),
)
producer.startup()
try:
transaction = producer.begin_transaction()
message = Message()
message.topic = "order-transaction-events"
message.keys = "order-o-100"
message.body = b'{"order_id":"o-100"}'
producer.send(message, transaction)

try:
commit_local_order_transaction("o-100")
transaction.commit()
except Exception:
rollback_local_order_transaction("o-100")
transaction.rollback()
raise
finally:
producer.shutdown()

本地事务状态必须持久化并可供回查,不能把结果只存在内存变量里。还要认真处理“数据库提交结果不确定”的情况,不能见异常就武断 rollback。

23.6 三个项目练习如何串起来

  1. RabbitMQ 项目练可靠任务语义:confirm、ACK、prefetch、retry、DLQ、Inbox。
  2. Kafka 改造练日志语义:partition、key、offset、group、lag、回放、Kafka 内事务。
  3. 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 治理

一个严谨流程:

  1. Schema 进入版本控制;
  2. CI 校验生产者变化是否向后兼容;
  3. 消费者先兼容新旧字段并上线;
  4. 生产者再开始发送新字段;
  5. 观察所有消费者升级;
  6. 经过弃用窗口后才移除旧字段;
  7. 保留示例消息和契约测试。

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
2
message_id, event_type, topic/queue, partition/offset(若有), consumer_group,
aggregate_id, trace_id, retry_count, processing_result

不要只打印“消费失败”。

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 消息丢失排障

沿消息旅程逐段寻找证据:

  1. 业务数据库是否提交?outbox 是否有行?
  2. relay 是否 claim?publish 是否收到 confirm?
  3. 是否发错 topic/routing key?是否 unroutable?
  4. Broker 是否按预期持久化和复制?是否过期/达到长度限制?
  5. 消费组是否订阅正确?offset 是否被错误重置?
  6. 消费者是否提前 ACK/commit?
  7. 是否进入 retry/DLQ?
  8. 是否其实处理成功,只是日志/查询读到旧缓存?

没有 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 基础自测

先自己回答,再看本章末答案。

  1. Python asyncio.Queue 为什么不能替代 RabbitMQ?
  2. FastAPI BackgroundTasks 为什么不适合扣款?
  3. queue 和 topic 的直觉差异是什么?
  4. “推消费者”底层是否一定是真正服务器推送?
  5. 削峰和限流有什么区别?
  6. 为什么 MQ 降低运行时耦合却没有消灭契约耦合?
  7. 一个 consumer group 中消费者数超过 Kafka partition 数会怎样?
  8. 为什么消费后 ACK 仍然会重复?
  9. 为什么 SETNX 不是万能幂等?
  10. DLQ 为什么不是最终解决方案?

27.8 深入自测

  1. 订单数据库提交成功、消息未发,至少给出两种解决方案。
  2. Kafka acks=all 为什么还要配合 min.insync.replicas
  3. Kafka exactly-once 为什么不自动覆盖 PostgreSQL?
  4. RabbitMQ durable queue、persistent message、publisher confirm 分别解决什么?
  5. RocketMQ 事务消息为什么仍要求下游幂等?
  6. 为什么同一 partition 有序仍可能业务乱序?
  7. 积压 100 万条时为什么不能立刻无限扩 consumer?
  8. 延迟消息到期为什么不能无条件取消订单?
  9. Outbox relay 为什么会重复发布?
  10. 如何判断某条消息是暂时性错误还是永久性错误?

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 自测答案要点

  1. asyncio.Queue 是单进程内存结构,缺少跨机器、持久化、副本、消费组和管理能力。
  2. Web 进程崩溃会丢任务,没有独立持久化和重试闭环。
  3. queue 偏待办任务和确认后删除;topic/日志偏事件分类、多订阅和保留回放。
  4. 不一定,常由 SDK 长轮询后回调业务代码。
  5. 削峰接收并缓冲,限流拒绝/等待/降级以保护资源。
  6. 生产者不依赖消费者在线和地址,但双方仍依赖消息字段语义。
  7. 多余实例通常无 partition 可分而空闲。
  8. 数据库已提交、ACK 未到 Broker 时宕机会重投。
  9. 去重标记与业务写不原子时,标记成功后业务失败会造成永久跳过。
  10. 必须告警、分析、修复、重放和审计,否则只是隐藏丢失。
  11. Transactional Outbox;RocketMQ 事务消息;特定边界还可用数据库日志 CDC。
  12. acks=all 基于当前 ISR;若 ISR 可缩到 1,仍可能单副本确认,最小同步副本限制安全下限。
  13. Kafka 事务协调 Kafka 记录和 offset,不控制外部数据库事务管理器。
  14. 拓扑重启存在;消息持久存储;生产者知道 Broker 是否接管。
  15. 事务消息只协调上游本地事务和消息可见性,消费者执行仍可能失败或重复。
  16. 消费端并发线程可能 2 先于 1 完成,重试也会改变完成顺序。
  17. 瓶颈可能是数据库;无限扩容会耗尽连接并加剧故障。
  18. 到期时可能已付款,必须用状态机条件更新消解竞态。
  19. Broker confirm 后、outbox 标记 PUBLISHED 前宕机,重启会再发。
  20. 网络/超时/限流通常暂时;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 最终心智模型

当你下次看到任何消息系统,依次问十个问题:

  1. 数据模型是任务队列还是追加日志?
  2. 生产者成功的确认边界在哪里?
  3. Broker 如何持久化和复制?
  4. 消费者如何获得消息,推、拉还是长轮询?
  5. 进度由 ACK、offset 还是不可见时间管理?
  6. 失败如何重试,何时进入死信?
  7. 重复如何幂等?
  8. 顺序在哪个范围内保证?
  9. 数据库与消息如何保持最终一致?
  10. 如何监控、重放、对账和灾备?

能回答这十问,你就不只会 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
2
3
4
5
6
7
8
9
10
11
12
13
14
需要任务队列、灵活路由、Python 快速落地
→ RabbitMQ

需要大规模事件流、长期保留、回放、CDC、流处理
→ Kafka

需要原生业务事务消息、FIFO、定时/延迟,并有 RocketMQ 平台能力
→ RocketMQ

只需一个很小的进程内异步任务
→ asyncio.Queue / FastAPI BackgroundTasks(仅限允许进程崩溃丢失)

需要跨进程可靠任务,但不想直接管理底层 API
→ 评估 Celery/Dramatiq + 合适 Broker

附录 C:可靠消息伪代码模板

生产端

1
2
3
4
5
6
7
8
9
10
11
async with database.transaction():
save_business_state()
save_outbox_event(message_id, payload)

# 独立 relay
while event := claim_pending_outbox():
try:
publish_and_wait_for_confirm(event)
mark_published(event)
except TemporaryError:
release_for_retry(event)

消费端

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
async def consume(message):
try:
async with database.transaction():
if inbox_already_contains(message.id):
return ACK
validate_contract(message)
apply_business_change(message)
insert_inbox(message.id)
return ACK
except PermanentError:
publish_to_dlq_with_confirm(message)
return ACK
except TemporaryError:
schedule_bounded_retry_or_dlq(message)
return ACK_ONLY_AFTER_RETRY_PUBLISH_CONFIRMED

附录 D:官方资料与延伸阅读

以下链接用于核对产品事实和 API;阅读时选择与你部署版本一致的文档。

Apache Kafka

RabbitMQ

Apache RocketMQ

通用架构

附录 E:版本相关事实说明

  1. Kafka、RabbitMQ、RocketMQ 仍在持续演进。本文将稳定原理与当前产品能力分开描述。
  2. Kafka 4.x 的运维方式已以 KRaft 为中心,旧 ZooKeeper 教程只适合历史版本。
  3. RabbitMQ 4.x 已移除 classic mirrored queues;复制队列优先看 quorum queues 或 streams。
  4. RocketMQ 4.x 的固定延迟等级示例与 5.x Delay topic/时间戳模型不是同一套 API。
  5. RocketMQ 官方网站 SDK 概览与 rocketmq-clients 仓库的语言能力矩阵可能存在更新节奏差异;Python 接入应以目标发布版本、Proxy 与集成测试共同确认。
  6. aiokafkaaio-pika 等第三方客户端的默认配置可能和 Java 官方客户端不同。教程中显式配置关键语义,就是为了避免把一种客户端的默认值误认为产品永恒规则。

结语

消息队列最迷人的地方不是“把一条 JSON 从 A 搬到 B”,而是它迫使我们认真面对分布式系统的现实:网络会超时,进程会在任意一行代码后崩溃,消息会重复,机器速度不同,强一致有成本,恢复能力比完美幻想更重要。

从现在起,看到一条消息时,不要只问“怎么发”,而要问:它代表什么事实、成功边界在哪里、失败后谁负责、重复是否安全、顺序范围是什么、多久必须完成、如何发现并恢复异常。能把这些问题回答清楚,你就已经具备将消息队列用于真实项目的核心能力。


消息队列学习
http://jack-constantine.github.io/2026/08/12/消息队列学习/
作者
JackConstantine
发布于
2026年8月12日
许可协议