MQ 商业场景训练营
MQ 不能只背“异步、削峰、解耦”。真正到项目和面试里,你要能解释一条消息从业务事务、生产者、Broker、消费者、ACK、重试、死信、补偿到最终一致的全过程;还要能解释为什么会丢、为什么会重复、为什么会堆积、为什么扩容消费者短暂有效后又开始堆积。
训练目标:用订单、支付、库存、ES 同步、短信通知、采集任务入库这些商业场景,把可靠性、幂等、事务消息、堆积、Lag 分布、死信和补偿跑通。
训练总流程
flowchart TD
A["业务事务"] --> B["生成业务事件"]
B --> C["可靠发送到 MQ"]
C --> D["Broker 持久化"]
D --> E["消费者拉取或推送"]
E --> F["业务处理"]
F --> G{"处理成功"}
G -- "是" --> H["ACK / 提交 Offset"]
G -- "否" --> I["重试 / 死信 / 补偿"]
I --> J["排查堆积和失败"]训练一:订单创建事件
场景
用户创建订单后,主流程只负责订单落库。积分、短信、ES 同步、操作日志这些下游通过 MQ 异步处理。
为什么不能全部同步调用
flowchart TD
A["创建订单"] --> B["写订单库"]
B --> C["扣库存"]
C --> D["发积分"]
D --> E["发短信"]
E --> F["同步 ES"]
F --> G["返回用户"]这样的问题是链路长、下游慢会拖慢主接口、某个下游失败会影响订单创建。
使用 MQ 后
flowchart TD
A["创建订单"] --> B["订单落库"]
B --> C["发送 OrderCreated 事件"]
C --> D["MQ Broker"]
D --> E["库存消费者"]
D --> F["积分消费者"]
D --> G["短信消费者"]
D --> H["ES 同步消费者"]消息体 Demo
{
"eventId": "evt-20260706-0001",
"eventType": "OrderCreated",
"orderNo": "O202607060001",
"userId": 1001,
"occurredAt": "2026-07-06T10:00:00",
"version": 1
}必须有业务唯一标识,例如 eventId 或 orderNo + eventType + version,否则消费者不好做幂等。
训练二:本地消息表保证可靠发送
问题
订单写库成功,但 MQ 发送失败怎么办?如果直接在事务外发送,一旦进程宕机,消息可能丢。
本地消息表
create table outbox_event (
id bigint primary key auto_increment,
event_id varchar(64) not null,
event_type varchar(64) not null,
aggregate_id varchar(64) not null,
payload text not null,
status varchar(20) not null,
retry_count int not null default 0,
next_retry_time datetime not null,
created_at datetime not null,
updated_at datetime not null,
unique key uk_event_id (event_id),
key idx_status_retry_time (status, next_retry_time)
);订单事务内同时写订单和本地消息:
@Transactional
public void createOrder(CreateOrderCommand command) {
Order order = orderRepository.save(command.toOrder());
OutboxEvent event = OutboxEvent.orderCreated(order);
outboxEventRepository.save(event);
}后台任务发送:
public void publishPendingEvents() {
List<OutboxEvent> events = outboxEventRepository.findPending(100);
for (OutboxEvent event : events) {
try {
mqTemplate.send("order.created", event.getPayload());
outboxEventRepository.markSent(event.getEventId());
} catch (Exception ex) {
outboxEventRepository.markRetry(event.getEventId(), ex.getMessage());
}
}
}原理图
flowchart TD
A["业务事务开始"] --> B["写订单表"]
B --> C["写 outbox_event"]
C --> D["事务提交"]
D --> E["后台扫描待发送事件"]
E --> F["发送 MQ"]
F --> G{"发送成功"}
G -- "是" --> H["标记 SENT"]
G -- "否" --> I["记录失败并重试"]这样保证“业务数据”和“待发送消息”在同一个数据库事务里,要么都成功,要么都失败。
训练三:消费者幂等
为什么一定要幂等
多数 MQ 更容易保证至少投递一次。消费者处理成功但 ACK 失败、网络超时、重平衡、Broker 重试,都可能让同一条消息再次投递。
幂等表
create table consume_log (
id bigint primary key auto_increment,
consumer_group varchar(128) not null,
event_id varchar(64) not null,
status varchar(20) not null,
created_at datetime not null,
unique key uk_consumer_event (consumer_group, event_id)
);消费者 Demo:
@Transactional
public void onMessage(OrderCreatedEvent event) {
boolean first = consumeLogRepository.tryInsert(
"es-sync-consumer", event.getEventId()
);
if (!first) {
return;
}
esSyncService.syncOrder(event.getOrderNo());
consumeLogRepository.markSuccess("es-sync-consumer", event.getEventId());
}原理
flowchart TD
A["收到消息"] --> B["插入 consume_log 唯一键"]
B --> C{"插入成功"}
C -- "否" --> D["重复消息,直接返回"]
C -- "是" --> E["执行业务"]
E --> F["标记消费成功"]幂等不是为了让重复消息不存在,而是为了重复消息来了也不重复产生副作用。
训练四:业务成功后再 ACK
错误做法
消费者刚收到消息就 ACK,然后再写数据库或 ES。如果后续业务失败,Broker 已经认为消息成功,消息不会再投递。
flowchart TD
A["收到消息"] --> B["提前 ACK"]
B --> C["写 ES"]
C --> D{"写 ES 失败"}
D --> E["消息已丢失,需要人工补偿"]正确顺序
flowchart TD
A["收到消息"] --> B["幂等检查"]
B --> C["执行业务"]
C --> D{"业务成功"}
D -- "是" --> E["ACK / 提交 Offset"]
D -- "否" --> F["不 ACK / 抛异常 / 进入重试"]标准回答:
消费者不能提前 ACK。ACK 或 offset 提交表示 Broker 可以认为消息已经处理完成,如果业务还没成功就 ACK,后续失败会导致消息丢失。正确做法是业务成功、幂等状态落库之后再 ACK;失败则重试、死信或补偿。训练五:ES 同步失败怎么办
场景
MySQL 是订单事实源,ES 是搜索视图。订单状态更新后发送 MQ,同步 ES。如果 ES 更新失败,不能静默丢。
处理链路
flowchart TD
A["订单状态变更"] --> B["MySQL 提交"]
B --> C["发送 MQ"]
C --> D["ES 同步消费者"]
D --> E{"写 ES 成功"}
E -- "是" --> F["ACK"]
E -- "否" --> G["延迟重试"]
G --> H{"超过重试上限"}
H -- "否" --> D
H -- "是" --> I["死信队列 / 异常表"]
I --> J["补偿任务从 MySQL 重建 ES"]为什么可以补偿
ES 不是事实源,MySQL 才是。只要 MySQL 有完整状态,就可以通过补偿任务重新同步 ES。不能做的是:ES 失败后 ACK 掉消息且没有异常记录。
训练六:消息堆积排查
先判断是不是异常
计划内削峰:活动高峰消息上涨,下游按能力慢慢消费,最老消息延迟在 SLA 内。
异常堆积:Lag 持续上涨,最老消息等待时间超过 SLA,失败率或重试量上升。
核心公式
净消化速度 = 成功消费 TPS - 生产 TPS
追平时间 = 当前堆积量 / 净消化速度如果成功消费 TPS 小于生产 TPS,系统永远追不上,必须先限流或提高真实消费能力。
排查流程
flowchart TD
A["发现堆积"] --> B["看生产 TPS"]
B --> C["看成功消费 TPS"]
C --> D["看失败率和重试量"]
D --> E["看 Lag 分布"]
E --> F{"是否局部队列/分区高"}
F -- "是" --> G["热点 Key / 慢消息 / 顺序阻塞"]
F -- "否" --> H["整体消费能力不足或下游慢"]
H --> I["查 DB / ES / Redis / 远程接口耗时"]
G --> J["定位具体消息和队列"]训练七:为什么扩容后短暂有效又堆积
原理
扩容消费者后,刚开始消费 TPS 上升,Lag 下降。但如果真正瓶颈是数据库、ES、Redis、远程接口、连接池、线程池或 Broker,更多消费者会把下游打得更慢,单条消息耗时升高,最终消费 TPS 又下降。
flowchart TD
A["消费者扩容"] --> B["并发上涨"]
B --> C["短期消费 TPS 上升"]
C --> D["更多压力打到下游"]
D --> E["DB/ES/接口变慢"]
E --> F["单条消息耗时升高"]
F --> G["成功消费 TPS 回落"]
G --> H["Lag 再次上涨"]常见根因
| 根因 | 表现 | 处理 |
|---|---|---|
| 分区/队列数不足 | 消费者多于分区也没用 | 增加分区或队列,重新分配 |
| 下游 DB 慢 | DB CPU/慢 SQL 上升 | 优化 SQL、批量、限流 |
| ES 写入慢 | bulk 队列、refresh 压力 | 批量写、调 refresh、限流 |
| 重试风暴 | 失败消息反复重试 | 延迟重试、死信隔离 |
| 热点 Key | 少数分区 Lag 特别高 | 调整 key、拆热点 |
| 顺序消息阻塞 | 一个队列被坏消息卡住 | 修复坏消息、人工补偿 |
训练八:Lag 分布定位
总 Lag 只能说明积压规模,Lag 分布才能说明哪里卡住。
partition-0 lag: 10
partition-1 lag: 12
partition-2 lag: 180000
partition-3 lag: 9这种不是整体慢,而是 partition-2 有热点、慢消息或消费者异常。
flowchart TD
A["总 Lag 很高"] --> B["拆分到分区/队列"]
B --> C{"Lag 是否均匀"}
C -- "均匀" --> D["整体消费能力不足"]
C -- "不均匀" --> E["局部热点或阻塞"]
E --> F["看消息 Key / 错误日志 / 单条耗时"]训练九:RabbitMQ Ready 和 Unacked
| 指标 | 含义 | 常见原因 |
|---|---|---|
| Ready 高 | 消息在队列里,还没投递给消费者 | 消费者不足、消费者掉线、prefetch 太小 |
| Unacked 高 | 已投递给消费者,但还没 ACK | 消费者处理慢、业务卡住、prefetch 太大 |
治理思路:
- Ready 高先看消费者是否在线、消费速率、队列数。
- Unacked 高先看业务耗时、线程池、ACK 时机、prefetch。
- 不能盲目把 prefetch 调很大,否则消息堆在消费者内存里,Broker 看起来少了,业务并没有处理完。
训练十:不同 MQ 选型
| 产品 | 更适合 | 典型场景 |
|---|---|---|
| RabbitMQ | 复杂路由、任务队列、传统业务消息 | 通知、审批、普通异步任务 |
| Kafka | 高吞吐、事件流、日志、CDC、可回放 | 埋点、日志、binlog、行为流 |
| RocketMQ | 业务消息、事务消息、顺序、重试 | 订单、支付、库存、交易链路 |
选型不能只看性能。要看团队经验、部署运维、可靠性、重试模型、事务消息、顺序消息、监控工具和业务生态。
最终验收清单
做完这页后,你应该能回答:
- 为什么 MQ 不能只说异步、削峰、解耦?
- 订单落库成功但 MQ 发送失败怎么办?
- 本地消息表和事务消息分别解决什么问题?
- 消费端为什么必须幂等?
- 为什么不能提前 ACK?
- ES 同步失败怎么补偿?
- 消息堆积先看哪些指标?
- 净消化速度和追平时间怎么算?
- 扩容消费者后为什么可能再次堆积?
- Lag 分布是什么,为什么比总 Lag 更重要?
- RabbitMQ Ready 和 Unacked 高分别说明什么?
- Kafka、RabbitMQ、RocketMQ 怎么选?
关联知识点
| 知识点 | 入口 |
|---|---|
| MQ 主线 | 从零到生产级掌握 |
| 消息堆积 | 消息堆积与背压 |
| 扩容反弹 | 扩容后再次堆积 |
| Lag 分布 | Lag 分布与定位 |
| RabbitMQ | RabbitMQ 核心模型 |
| Kafka | Kafka 总览 |
| RocketMQ | RocketMQ 总览 |
| MQ 面试 | 消息队列面试题 |
