Skip to content

MQ 商业场景训练营

MQ 不能只背“异步、削峰、解耦”。真正到项目和面试里,你要能解释一条消息从业务事务、生产者、Broker、消费者、ACK、重试、死信、补偿到最终一致的全过程;还要能解释为什么会丢、为什么会重复、为什么会堆积、为什么扩容消费者短暂有效后又开始堆积。

训练目标:用订单、支付、库存、ES 同步、短信通知、采集任务入库这些商业场景,把可靠性、幂等、事务消息、堆积、Lag 分布、死信和补偿跑通。

训练总流程

mermaid
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 异步处理。

为什么不能全部同步调用

mermaid
flowchart TD
    A["创建订单"] --> B["写订单库"]
    B --> C["扣库存"]
    C --> D["发积分"]
    D --> E["发短信"]
    E --> F["同步 ES"]
    F --> G["返回用户"]

这样的问题是链路长、下游慢会拖慢主接口、某个下游失败会影响订单创建。

使用 MQ 后

mermaid
flowchart TD
    A["创建订单"] --> B["订单落库"]
    B --> C["发送 OrderCreated 事件"]
    C --> D["MQ Broker"]
    D --> E["库存消费者"]
    D --> F["积分消费者"]
    D --> G["短信消费者"]
    D --> H["ES 同步消费者"]

消息体 Demo

json
{
  "eventId": "evt-20260706-0001",
  "eventType": "OrderCreated",
  "orderNo": "O202607060001",
  "userId": 1001,
  "occurredAt": "2026-07-06T10:00:00",
  "version": 1
}

必须有业务唯一标识,例如 eventIdorderNo + eventType + version,否则消费者不好做幂等。

训练二:本地消息表保证可靠发送

问题

订单写库成功,但 MQ 发送失败怎么办?如果直接在事务外发送,一旦进程宕机,消息可能丢。

本地消息表

sql
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)
);

订单事务内同时写订单和本地消息:

java
@Transactional
public void createOrder(CreateOrderCommand command) {
    Order order = orderRepository.save(command.toOrder());

    OutboxEvent event = OutboxEvent.orderCreated(order);
    outboxEventRepository.save(event);
}

后台任务发送:

java
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());
        }
    }
}

原理图

mermaid
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 重试,都可能让同一条消息再次投递。

幂等表

sql
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:

java
@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());
}

原理

mermaid
flowchart TD
    A["收到消息"] --> B["插入 consume_log 唯一键"]
    B --> C{"插入成功"}
    C -- "否" --> D["重复消息,直接返回"]
    C -- "是" --> E["执行业务"]
    E --> F["标记消费成功"]

幂等不是为了让重复消息不存在,而是为了重复消息来了也不重复产生副作用。

训练四:业务成功后再 ACK

错误做法

消费者刚收到消息就 ACK,然后再写数据库或 ES。如果后续业务失败,Broker 已经认为消息成功,消息不会再投递。

mermaid
flowchart TD
    A["收到消息"] --> B["提前 ACK"]
    B --> C["写 ES"]
    C --> D{"写 ES 失败"}
    D --> E["消息已丢失,需要人工补偿"]

正确顺序

mermaid
flowchart TD
    A["收到消息"] --> B["幂等检查"]
    B --> C["执行业务"]
    C --> D{"业务成功"}
    D -- "是" --> E["ACK / 提交 Offset"]
    D -- "否" --> F["不 ACK / 抛异常 / 进入重试"]

标准回答:

text
消费者不能提前 ACK。ACK 或 offset 提交表示 Broker 可以认为消息已经处理完成,如果业务还没成功就 ACK,后续失败会导致消息丢失。正确做法是业务成功、幂等状态落库之后再 ACK;失败则重试、死信或补偿。

训练五:ES 同步失败怎么办

场景

MySQL 是订单事实源,ES 是搜索视图。订单状态更新后发送 MQ,同步 ES。如果 ES 更新失败,不能静默丢。

处理链路

mermaid
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,失败率或重试量上升。

核心公式

text
净消化速度 = 成功消费 TPS - 生产 TPS
追平时间 = 当前堆积量 / 净消化速度

如果成功消费 TPS 小于生产 TPS,系统永远追不上,必须先限流或提高真实消费能力。

排查流程

mermaid
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 又下降。

mermaid
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 分布才能说明哪里卡住。

text
partition-0 lag: 10
partition-1 lag: 12
partition-2 lag: 180000
partition-3 lag: 9

这种不是整体慢,而是 partition-2 有热点、慢消息或消费者异常。

mermaid
flowchart TD
    A["总 Lag 很高"] --> B["拆分到分区/队列"]
    B --> C{"Lag 是否均匀"}
    C -- "均匀" --> D["整体消费能力不足"]
    C -- "不均匀" --> E["局部热点或阻塞"]
    E --> F["看消息 Key / 错误日志 / 单条耗时"]

训练九:RabbitMQ Ready 和 Unacked

指标含义常见原因
Ready 高消息在队列里,还没投递给消费者消费者不足、消费者掉线、prefetch 太小
Unacked 高已投递给消费者,但还没 ACK消费者处理慢、业务卡住、prefetch 太大

治理思路:

  1. Ready 高先看消费者是否在线、消费速率、队列数。
  2. Unacked 高先看业务耗时、线程池、ACK 时机、prefetch。
  3. 不能盲目把 prefetch 调很大,否则消息堆在消费者内存里,Broker 看起来少了,业务并没有处理完。

训练十:不同 MQ 选型

产品更适合典型场景
RabbitMQ复杂路由、任务队列、传统业务消息通知、审批、普通异步任务
Kafka高吞吐、事件流、日志、CDC、可回放埋点、日志、binlog、行为流
RocketMQ业务消息、事务消息、顺序、重试订单、支付、库存、交易链路

选型不能只看性能。要看团队经验、部署运维、可靠性、重试模型、事务消息、顺序消息、监控工具和业务生态。

最终验收清单

做完这页后,你应该能回答:

  1. 为什么 MQ 不能只说异步、削峰、解耦?
  2. 订单落库成功但 MQ 发送失败怎么办?
  3. 本地消息表和事务消息分别解决什么问题?
  4. 消费端为什么必须幂等?
  5. 为什么不能提前 ACK?
  6. ES 同步失败怎么补偿?
  7. 消息堆积先看哪些指标?
  8. 净消化速度和追平时间怎么算?
  9. 扩容消费者后为什么可能再次堆积?
  10. Lag 分布是什么,为什么比总 Lag 更重要?
  11. RabbitMQ Ready 和 Unacked 高分别说明什么?
  12. Kafka、RabbitMQ、RocketMQ 怎么选?

关联知识点

知识点入口
MQ 主线从零到生产级掌握
消息堆积消息堆积与背压
扩容反弹扩容后再次堆积
Lag 分布Lag 分布与定位
RabbitMQRabbitMQ 核心模型
KafkaKafka 总览
RocketMQRocketMQ 总览
MQ 面试消息队列面试题