MQ 从零到精通验收清单
这页用来回答:消息队列专栏是不是只讲了异步、解耦、削峰这些表面词,还是能让零基础真正理解生产级 MQ?
学完 MQ 不能只会说“加个队列”。你要能解释消息从业务代码发出,到 Broker 存储,再到消费者处理成功的每一步;能知道哪里会丢、哪里会重复、哪里会乱序、哪里会堆积;能在面试和生产排障中给出证据链。
最终目标
学完 MQ 专栏,你至少要能做到:
- 画出生产者、Broker、Topic/Queue、ConsumerGroup、消费者之间的关系。
- 解释同步调用为什么会被 MQ 改成异步链路,以及代价是什么。
- 解释消息从发送、存储、投递、消费、ACK/Offset 的全过程。
- 解释消息为什么会丢、为什么会重复、为什么会乱序。
- 设计消费端幂等、重试、死信、补偿和对账。
- 排查消息堆积、Lag 分布、Ready/Unacked、扩容反弹。
- 解释 Kafka、RabbitMQ、RocketMQ 的模型差异和选型边界。
- 在订单、支付、库存、ES 同步、短信通知、采集任务里落地 MQ。
- 面试时标准回答在面试页,追问原理能跳回知识点页展开。
总路线
flowchart TD
A["阶段 1<br/>理解为什么需要 MQ"] --> B["阶段 2<br/>消息发送全过程"]
B --> C["阶段 3<br/>Broker 存储与投递"]
C --> D["阶段 4<br/>消费者 ACK / Offset"]
D --> E["阶段 5<br/>可靠性、重复、幂等"]
E --> F["阶段 6<br/>顺序、事务、延迟、死信"]
F --> G["阶段 7<br/>堆积、Lag、背压"]
G --> H["阶段 8<br/>三大 MQ 选型"]
H --> I["阶段 9<br/>商业场景和面试"]为什么要这样学?
| 阶段 | 原因 | 如果跳过 |
|---|---|---|
| 先学作用 | 知道 MQ 解决什么,不解决什么 | 什么场景都想加 MQ |
| 再学发送 | 丢消息首先可能发生在生产者 | 只会怪 Broker |
| 再学存储 | Broker 不是黑盒,可靠性依赖持久化和复制 | 不会解释刷盘、副本、保留 |
| 再学消费 | 消费成功不等于消息投递到了 | 提前 ACK、自动提交 offset 导致漏处理 |
| 再学幂等 | 重复消费是常态风险 | 订单、库存、短信重复执行 |
| 再学高级能力 | 顺序、事务、延迟、死信都有边界 | 以为 MQ 能保证全链路事务 |
| 再学堆积 | 生产问题最高频 | 只会加消费者或清队列 |
| 最后选型 | 不同产品模型不同 | 用错产品或用错模型 |
阶段 1:为什么需要 MQ
MQ 的核心不是“炫技”,而是把同步链路拆成异步事件链路。
flowchart TD
A["订单服务创建订单"] --> B["写订单库"]
B --> C["发送订单创建事件"]
C --> D["MQ Broker"]
D --> E["库存消费者"]
D --> F["积分消费者"]
D --> G["短信消费者"]
D --> H["ES 同步消费者"]必须理解
| 能力 | 解释 | 不这样会怎样 |
|---|---|---|
| 异步 | 主链路不用等所有下游完成 | 下游慢会拖慢下单接口 |
| 解耦 | 新增消费者不用改生产者核心流程 | 每加一个下游都要改订单服务 |
| 削峰 | 高峰消息先进入 Broker | 流量直接打爆数据库或接口 |
| 最终一致 | 主业务先成功,下游靠重试补偿追上 | 分布式强一致成本过高 |
| 可回放 | Kafka 等日志模型可按 offset 重放 | 数据同步失败后无法重建 |
边界
MQ 不是万能的。
| 误用 | 问题 |
|---|---|
| 把核心强一致操作全丢给 MQ | 用户可能看到中间状态 |
| 没有幂等就用 MQ | 重复消息会重复扣库存、发短信 |
| 没有补偿就用 MQ | 下游失败后数据长期不一致 |
| 没有监控就用 MQ | 堆积几小时才发现 |
| 把 MQ 当数据库 | 查询、事务、复杂条件检索都不适合 |
阶段 2:消息发送全过程
生产者发送消息不是一行 send 就结束。真正可靠发送至少要考虑消息构造、路由、序列化、发送确认、失败重试和本地事务边界。
flowchart TD
A["业务事务开始"] --> B["写业务表"]
B --> C["构建事件消息"]
C --> D["序列化和设置 Key"]
D --> E["发送到 Broker"]
E --> F{"Broker 是否确认"}
F -- "是" --> G["记录发送成功"]
F -- "否" --> H["重试或写入本地消息表"]发送阶段可能丢在哪里
| 位置 | 例子 | 解决方式 |
|---|---|---|
| 业务已提交但消息没发 | 订单已支付,通知消息没出去 | Outbox 本地消息表、事务消息、CDC |
| 发送到 Broker 超时 | 网络抖动,生产者不知道成功没成功 | 发送确认、幂等 key、查询发送结果 |
| Producer 进程宕机 | 内存消息未发出 | 先落库再异步投递 |
| 序列化失败 | 消息结构错误 | 参数校验、统一协议、失败告警 |
Outbox Demo
create table outbox_message (
id bigint primary key auto_increment,
biz_key varchar(128) not null,
topic varchar(128) not null,
payload text not null,
status varchar(32) not null,
retry_count int not null default 0,
created_at datetime not null,
updated_at datetime not null,
unique key uk_biz_key (biz_key)
);业务事务内同时写订单和消息表:
@Transactional(rollbackFor = Exception.class)
public void paySuccess(PaySuccessCommand command) {
orderRepository.markPaid(command.orderNo());
OutboxMessage message = OutboxMessage.create(
"ORDER_PAID:" + command.orderNo(),
"order.paid",
toJson(new OrderPaidEvent(command.orderNo()))
);
outboxRepository.save(message);
}后台任务扫描 outbox_message 发送到 MQ,成功后改状态。这样即使发送服务宕机,也能从数据库继续补发。
阶段 3:Broker 存储与投递
Broker 负责接收、存储和投递消息,但不同 MQ 模型差异很大。
| 产品 | 核心模型 | 你必须理解 |
|---|---|---|
| RabbitMQ | Exchange、Queue、Binding | 路由、ACK、Ready、Unacked、DLX |
| Kafka | Topic、Partition、Log、Offset | 分区、日志追加、消费组、Lag |
| RocketMQ | Topic、MessageQueue、CommitLog、ConsumeQueue | NameServer、Broker、队列、重试、事务 |
Broker 不是只“转发”
flowchart TD
A["Producer 发送消息"] --> B["Broker 接收"]
B --> C["写入内存和日志文件"]
C --> D["复制或刷盘策略"]
D --> E["更新队列或索引结构"]
E --> F["等待消费者拉取或推送"]如果 Broker 没有持久化、没有副本、磁盘水位失控,消息可靠性和堆积能力都会出问题。
阶段 4:消费者 ACK 和 Offset
消费者阶段最容易出现两个相反错误:
- 业务没成功就确认,造成漏处理。
- 业务成功但确认失败,造成重复处理。
flowchart TD
A["消费者拿到消息"] --> B["解析和幂等检查"]
B --> C["执行业务处理"]
C --> D{"业务是否真正成功"}
D -- "成功" --> E["ACK 或提交 Offset"]
D -- "失败" --> F["重试、延迟队列或死信"]Kafka Offset 为什么难
Kafka 的 offset 表示某个消费组在某个 Partition 里已经安全处理到哪里。
并发处理时不能随便提交最大 offset。
offset 100 成功
offset 101 成功
offset 102 失败
offset 103 成功此时不能提交到 104,因为 102 还没处理成功。否则消费者重启后会从 104 继续,102 就丢了。
正确思路是:每个分区只提交“连续成功”的最大位置。
RabbitMQ ACK 为什么重要
| ACK 策略 | 风险 |
|---|---|
| 自动 ACK | 消息刚投递就确认,业务失败会丢 |
| 业务成功后手动 ACK | 更可靠,但可能重复,需要幂等 |
| 失败立即 requeue | 可能形成重试风暴 |
| 失败进死信 | 需要告警、补偿和重放工具 |
阶段 5:重复消费和幂等
MQ 生产系统要默认“消息可能重复”。
重复来源:
| 来源 | 例子 |
|---|---|
| 生产者重试 | Broker 已收到,但响应超时,生产者又发一次 |
| 消费者宕机 | 业务成功后,还没 ACK 或提交 offset |
| Broker 重投 | Broker 没收到确认 |
| Rebalance | Kafka 分区重新分配后从旧 offset 继续 |
| 手工重放 | 运维为了补数据重新投递 |
幂等 Demo
create table mq_consume_log (
id bigint primary key auto_increment,
message_id varchar(128) not null,
consumer_group varchar(128) not null,
biz_key varchar(128) not null,
status varchar(32) not null,
created_at datetime not null,
unique key uk_msg_group (message_id, consumer_group),
unique key uk_biz_group (biz_key, consumer_group)
);@Transactional(rollbackFor = Exception.class)
public void consume(OrderPaidEvent event) {
boolean first = consumeLogRepository.tryInsert(
event.messageId(),
"point-service",
event.orderNo()
);
if (!first) {
return;
}
pointRepository.addPoint(event.userId(), event.orderNo(), event.amount());
consumeLogRepository.markSuccess(event.messageId(), "point-service");
}为什么要用唯一索引?
| 做法 | 问题 |
|---|---|
| 先查再插 | 并发下两个消费者可能都查不到 |
| Redis 只做临时去重 | 过期后重放可能再次执行 |
| 数据库唯一索引 | 能用数据库原子约束兜底 |
阶段 6:顺序、事务、延迟、死信
顺序消息
顺序通常只能保证局部顺序,例如同一个订单的消息进入同一个分区或队列。
flowchart TD
A["orderId=1001 创建"] --> B["按 orderId 路由到 Partition 1"]
C["orderId=1001 支付"] --> B
D["orderId=1001 取消"] --> B
B --> E["同一个消费者按顺序处理"]全局顺序意味着所有消息进同一队列串行消费,吞吐会很低。商业项目通常只要求“同一个业务 key 局部有序”。
事务消息
事务消息解决的是:本地事务和消息发送之间的一致性。
它不解决:
- 消费端业务一定成功。
- 消费端不会重复。
- 下游数据库和上游数据库强一致。
- 所有服务同时提交或回滚。
消费端仍然要幂等、重试、死信、补偿和对账。
死信队列
死信队列不是垃圾桶,而是异常消息隔离区。
flowchart TD
A["业务队列"] --> B["消费者处理"]
B --> C{"是否成功"}
C -- "成功" --> D["确认消费"]
C -- "失败但可重试" --> E["延迟重试"]
E --> B
C -- "多次失败" --> F["死信队列"]
F --> G["告警、修复、重放或人工补偿"]阶段 7:堆积、Lag、背压
消息堆积的本质公式:
堆积增长速度 = 生产 TPS - 成功消费 TPS如果成功消费 TPS 小于生产 TPS,堆积一定会增长。拉取 TPS 很高没有用,关键是业务真正成功并确认的 TPS。
Lag 分布
总 Lag 只能说明积压规模,Lag 分布才能说明病灶。
| 分布 | 说明 | 排查方向 |
|---|---|---|
| 所有分区都高 | 整体消费能力不足或下游慢 | 扩容、批处理、优化 DB/ES/接口 |
| 少数分区高 | 热点 key、慢消息、异常消费者 | 查 key、线程栈、错误日志 |
| Unacked 高 | RabbitMQ 消息已投递但未确认 | 查消费者耗时、ACK 分支、prefetch |
| Ready 高 | RabbitMQ 消息还在队列等待 | 查消费者在线、投递速率、队列配置 |
扩容后再次堆积
如果扩容消费者后短暂有效,后面又开始堆积,说明扩容打出了更深层瓶颈。
flowchart TD
A["扩容消费者"] --> B["拉取 TPS 上升"]
B --> C["短期 Lag 下降"]
C --> D["下游 DB / ES / HTTP 压力升高"]
D --> E["单条处理耗时变长"]
E --> F["成功消费 TPS 下降"]
F --> G["Lag 再次上升"]排查必须看:
- 生产 TPS。
- 拉取 TPS。
- 成功 TPS。
- P95/P99 消费耗时。
- 错误率和重试量。
- Lag 分布。
- 消费者线程池和连接池。
- 下游慢 SQL、ES Bulk、HTTP 超时。
详细看:消息堆积与背压、扩容后再次堆积、Lag 分布与堆积定位。
阶段 8:RabbitMQ、Kafka、RocketMQ 选型
| 维度 | RabbitMQ | Kafka | RocketMQ |
|---|---|---|---|
| 模型 | Exchange + Queue | Topic + Partition + Log | Topic + MessageQueue |
| 典型优势 | 路由灵活、ACK、死信、任务队列 | 高吞吐、日志流、回放、生态 | 业务消息、事务、顺序、重试 |
| 常见场景 | 企业集成、延迟任务、通知 | 日志、埋点、CDC、实时计算 | 订单、支付、库存、交易链路 |
| 排查重点 | Ready、Unacked、prefetch | Lag、Partition、Rebalance | Diff、队列、重试、死信 |
| 误用风险 | 队列太多、Unacked 堆积 | 分区 key 设计错误、offset 提交错误 | 事务消息边界误解、顺序阻塞 |
选型不是“哪个更强”,而是看业务:
- 复杂路由、传统任务分发,优先考虑 RabbitMQ。
- 高吞吐事件流、日志、CDC、可回放,优先考虑 Kafka。
- Java 业务消息、事务消息、顺序消息、消费重试,优先考虑 RocketMQ。
阶段 9:商业场景验收
订单支付成功事件
必须做到:
- 支付状态更新和消息发送要可靠衔接。
- 消费端用订单号幂等。
- 积分、短信、ES 同步互不影响。
- 下游失败进入重试或死信。
- 有对账任务发现漏处理。
MySQL 同步 ES
必须做到:
- MySQL 是事实源,ES 是搜索视图。
- 消费消息后按主键回查 MySQL 最新快照。
- ES upsert 幂等。
- ES 写失败进入重试和补偿。
- Lag 超 SLA 告警。
医疗采集任务异步入库
必须做到:
- 原始采集数据不能只放内存。
- 消息包含采集批次号、医院编码、数据类型。
- 消费端按批次和唯一业务键幂等。
- 入库失败记录异常表。
- 支持补采、重放、对账和人工处理。
面试闭环
| 高频问题 | 标准回答 | 深入原理 |
|---|---|---|
| MQ 有什么用 | MQ 面试题 | MQ 从零到生产级掌握 |
| MQ 会不会丢 | MQ 面试题 | 可靠性链路 |
| 为什么重复消费 | MQ 面试题 | 幂等消费 |
| 消息堆积怎么办 | MQ 面试题 | 消息堆积与背压 |
| Lag 分布是什么 | MQ 面试题 | Lag 分布与定位 |
| 扩容后再次堆积 | MQ 面试题 | 扩容后再次堆积 |
| 三大 MQ 怎么选 | MQ 面试题 | RabbitMQ/Kafka/RocketMQ 选型 |
最终验收题
| 问题 | 合格标准 |
|---|---|
| 画 MQ 全链路 | 能画出生产、Broker、消费、ACK/Offset、重试、死信 |
| 解释可靠性 | 能分生产者、Broker、消费者三段分析 |
| 解释重复消费 | 能说出来源,并给幂等 Demo |
| 解释顺序消息 | 能说清局部顺序和吞吐代价 |
| 解释事务消息 | 能说清解决什么、不解决什么 |
| 排查堆积 | 能按生产 TPS、成功 TPS、Lag 分布、下游耗时定位 |
| 解释扩容反弹 | 能说出最短板、分区队列、下游瓶颈、重试风暴 |
| 处理死信 | 能说清告警、修复、重放、补偿 |
| 做商业设计 | 能设计订单、ES 同步、短信、采集任务方案 |
| 面试表达 | 标准回答简洁,追问能回到原理页 |
本章小结
真正学会 MQ,不是会背“异步、解耦、削峰”,而是能把发送可靠性、Broker 存储、消费者确认、幂等、顺序、事务、重试、死信、堆积、Lag、背压、补偿、选型串成一条完整链路。
如果某个环节讲不清,就回到对应知识点页重新看流程图、Demo 和商业场景。MQ 最怕半懂不懂,因为它不会立刻暴露问题,往往是在活动高峰、下游故障、重试风暴、堆积爆发时才把设计缺陷一次性放大。
