Skip to content

消息队列从零到生产级掌握

消息队列不能只学成“异步、削峰、解耦”三个词。真正到商业项目里,你必须能解释:消息从生产者到 Broker 再到消费者的每一步怎么工作,哪里可能丢,哪里可能重复,为什么会堆积,扩容为什么可能没用,顺序消息为什么牺牲吞吐,事务消息解决什么又不解决什么,RabbitMQ、Kafka、RocketMQ 到底怎么选。

一句话建立主线:

MQ 是分布式系统里的异步通信和缓冲层。它把同步调用变成事件或任务投递,用 Broker 存储消息,再由消费者按能力处理,从而实现解耦、削峰、最终一致和可回放。

学习目标

学完这一页,你要能做到:

  1. 解释 MQ 为什么存在,不用 MQ 会怎样。
  2. 解释生产者、Broker、Topic、Queue、Partition、ConsumerGroup、Offset、ACK 的关系。
  3. 解释一条消息从发送到消费成功的全过程。
  4. 解释消息丢失发生在哪三个阶段,以及怎么降低风险。
  5. 解释为什么消息会重复消费,消费端为什么必须幂等。
  6. 解释顺序消息为什么通常只能局部有序。
  7. 解释事务消息、本地消息表、Outbox、CDC 的关系。
  8. 解释延迟消息、重试队列、死信队列适合什么场景。
  9. 解释消息堆积、Lag、Lag 分布、Ready、Unacked、背压。
  10. 解释扩容消费者后短暂有效又堆积的根因。
  11. 对比 RabbitMQ、Kafka、RocketMQ 的模型、存储、ACK、事务、顺序、集群和选型。
  12. 能按订单、支付、库存、搜索同步、通知、采集入库等商业场景设计 MQ 链路。

如果你已经读完主线,但还不知道怎么把 MQ 原理落到订单、ES 同步、堆积排查和面试里,继续做:MQ 商业场景训练营。它把本地消息表、消费者幂等、ACK、重试死信、Lag 分布、扩容反弹和三大 MQ 选型串成可验证训练。

学习路线

mermaid
flowchart TD
    A["同步调用问题<br/>慢、耦合、扛不住高峰"] --> B["MQ 基础模型<br/>Producer、Broker、Consumer"]
    B --> C["发送流程<br/>确认、重试、持久化"]
    C --> D["消费流程<br/>拉取/推送、ACK、Offset"]
    D --> E["可靠性<br/>不丢、重复、幂等"]
    E --> F["顺序和事务<br/>局部顺序、最终一致"]
    F --> G["异常处理<br/>重试、死信、延迟"]
    G --> H["堆积与背压<br/>Lag、扩容、下游瓶颈"]
    H --> I["选型<br/>RabbitMQ、Kafka、RocketMQ"]

第一步:没有 MQ 会怎样

订单创建后要做很多事:

  1. 扣库存。
  2. 发积分。
  3. 发短信。
  4. 同步 ES。
  5. 写操作日志。
  6. 通知第三方。

如果全部同步调用:

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

问题:

问题后果
链路长用户等待时间变长
任一下游慢主接口被拖慢
任一下游失败主流程可能失败
高峰流量全部压力打到数据库和下游
新增下游要改订单主流程

使用 MQ 后:

mermaid
flowchart TD
    A["创建订单接口"] --> B["订单落库"]
    B --> C["发送订单创建事件"]
    C --> D["快速返回用户"]
    C --> E["库存消费者"]
    C --> F["积分消费者"]
    C --> G["短信消费者"]
    C --> H["ES 同步消费者"]

MQ 的价值:

价值解释
异步主流程不等所有下游完成
解耦新增消费者不影响生产者核心代码
削峰Broker 暂存高峰消息,消费者慢慢处理
最终一致主库先成功,下游通过消息追上
可回放Kafka 这类日志模型可以按 Offset 重新消费

第二步:MQ 通用模型

mermaid
flowchart TD
    A["Producer<br/>生产者"] --> B["Broker<br/>消息服务器"]
    B --> C["Topic / Exchange<br/>消息分类或路由入口"]
    C --> D["Queue / Partition<br/>实际存储和并行单元"]
    D --> E["ConsumerGroup<br/>消费组"]
    E --> F["Consumer<br/>消费者实例"]

不同 MQ 叫法不同:

通用概念RabbitMQKafkaRocketMQ
消息分类Exchange + RoutingKeyTopicTopic + Tag
存储单元QueuePartitionMessageQueue
消费进度Queue ACK 状态OffsetConsumer Offset
消费组Consumer Group 概念较弱但可按队列/消费者组织Consumer GroupConsumerGroup
BrokerRabbitMQ NodeKafka BrokerRocketMQ Broker

第三步:一条消息的完整链路

mermaid
flowchart TD
    A["业务代码构建消息"] --> B["生产者发送"]
    B --> C{"Broker 是否收到"}
    C -- "否" --> D["发送失败/重试/记录异常"]
    C -- "是" --> E["Broker 写入内存或磁盘"]
    E --> F["返回发送确认"]
    F --> G["消息等待投递或拉取"]
    G --> H["消费者获取消息"]
    H --> I["执行业务逻辑"]
    I --> J{"业务成功吗"}
    J -- "成功" --> K["ACK 或提交 Offset"]
    J -- "失败" --> L["重试、延迟、死信"]

一条消息真正可靠,不是看某一个点,而是看三个阶段:

阶段风险防线
生产阶段发送失败、超时、重复发送发送确认、重试、业务消息表
Broker 阶段未持久化宕机、主从未同步持久化、复制、副本、刷盘策略
消费阶段业务失败但已 ACK、ACK 丢失成功后 ACK、幂等、重试、死信

第四步:消息为什么会丢

生产者阶段丢

场景:

  1. 业务写库成功。
  2. 发送 MQ 时网络异常。
  3. 代码没有重试或记录。
  4. 下游永远收不到事件。

解决:

方案说明
发送确认Broker 确认收到才算成功
失败重试短暂网络问题自动重试
本地消息表业务数据和消息记录同事务提交
事务消息RocketMQ 半消息机制保证本地事务和消息发送一致
Outbox/CDC通过业务库变更可靠驱动消息

Broker 阶段丢

场景:

  1. Broker 收到消息。
  2. 还没刷盘或复制。
  3. Broker 宕机。
  4. 消息丢失。

解决:

  1. 开启持久化。
  2. 配置副本或主从复制。
  3. 选择合适刷盘策略。
  4. 对核心消息等待确认。

消费者阶段丢

最常见危险是提前 ACK:

mermaid
flowchart TD
    A["消费者收到消息"] --> B["提前 ACK"]
    B --> C["Broker 认为处理成功"]
    C --> D["消费者执行业务"]
    D --> E{"业务失败或进程宕机"}
    E --> F["消息不会再投递<br/>业务结果丢失"]

正确思路:

mermaid
flowchart TD
    A["消费者收到消息"] --> B["执行业务"]
    B --> C{"业务成功吗"}
    C -- "成功" --> D["ACK / 提交 Offset"]
    C -- "失败" --> E["不 ACK / 重试 / 死信"]

第五步:为什么会重复消费

多数 MQ 更容易保证“至少一次投递”,而不是“只投递一次”。

重复消费来源:

来源例子
ACK 丢失业务成功了,但 ACK 网络失败
消费者宕机处理完还没提交 offset
Broker 重试认为消费者没成功
RebalanceKafka 分区重新分配
生产者重试发送超时后重复发送
补偿任务反复扫描未完成消息

因此消费端必须幂等。

消费幂等 Demo

消费日志表:

sql
create table mq_consume_log (
  id bigint primary key auto_increment,
  consumer_group varchar(64) not null,
  message_key varchar(128) not null,
  status varchar(32) not null,
  created_at datetime not null,
  unique key uk_group_msg (consumer_group, message_key)
);

消费者伪代码:

java
public class OrderEventConsumer {
    public void onMessage(OrderCreatedEvent event) {
        boolean first = consumeLogRepository.tryInsert("stock-consumer", event.messageKey());
        if (!first) {
            return;
        }

        stockService.lockStock(event.orderNo(), event.skuId(), event.count());
        consumeLogRepository.markSuccess("stock-consumer", event.messageKey());
    }
}

关键点:

  1. 不要先查再处理,因为并发下可能两个消费者都查不到。
  2. 用唯一索引争抢处理权。
  3. 业务操作本身最好也有状态机或唯一约束。

第六步:顺序消息为什么难

全局顺序很贵。因为所有消息都必须进入一个队列,由一个消费者串行处理。

mermaid
flowchart TD
    A["所有订单消息"] --> B["同一个队列"]
    B --> C["一个消费者串行处理"]
    C --> D["吞吐很低"]

商业项目通常要的是局部顺序:

同一个订单的状态变化有序,不同订单之间可以并行。

mermaid
flowchart TD
    A["orderId=1001"] --> B["队列 1"]
    C["orderId=1002"] --> D["队列 2"]
    E["orderId=1003"] --> F["队列 3"]
    B --> G["消费者 A 顺序处理"]
    D --> H["消费者 B 顺序处理"]
    F --> I["消费者 C 顺序处理"]

实现方式:

MQ局部顺序方式
Kafka相同 key 路由到同一个 Partition
RocketMQ相同 sharding key 或消息组进入同一队列
RabbitMQ单 Queue 内有序,多个消费者并发会破坏严格处理顺序

顺序消息的风险:

  1. 某条消息失败会阻塞后续消息。
  2. 热点 key 会压垮单个分区或队列。
  3. 扩容消费者不一定提升同 key 吞吐。
  4. 业务异步处理后提前 ACK 会破坏顺序语义。

第七步:事务消息和本地消息表

分布式场景常见问题:

  1. 订单库提交成功。
  2. 发送消息失败。
  3. 库存、积分、ES 都不知道订单创建。

本地消息表:

mermaid
flowchart TD
    A["本地事务开始"] --> B["写订单表"]
    B --> C["写本地消息表"]
    C --> D["本地事务提交"]
    D --> E["投递任务扫描消息表"]
    E --> F["发送 MQ"]
    F --> G{"发送成功吗"}
    G -- "成功" --> H["标记已发送"]
    G -- "失败" --> I["下次继续重试"]

RocketMQ 事务消息:

mermaid
sequenceDiagram
    participant P as Producer
    participant B as Broker
    participant DB as 本地数据库
    P->>B: 发送半事务消息
    B-->>P: 半消息写入成功
    P->>DB: 执行本地事务
    DB-->>P: 本地事务成功
    P->>B: Commit 消息
    B-->>P: 消息可投递

区别:

方案解决点注意
本地消息表业务数据和消息记录同事务需要扫描任务和投递状态
RocketMQ 事务消息本地事务和消息发送一致消费端仍然要幂等
Kafka 事务Kafka 内部生产和消费位点事务外部数据库一致性仍常用 Outbox/CDC
CDC订阅数据库 binlog/WAL 生成消息链路复杂,需要处理顺序和补偿

事务消息只解决“本地事务和消息发送一致”,不保证消费者一定成功。消费端仍然要重试、死信、补偿和幂等。

第八步:重试、延迟、死信

消费者失败后不要无限快速重试,否则会形成重试风暴。

mermaid
flowchart TD
    A["消费失败"] --> B{"是否短暂故障"}
    B -- "是" --> C["延迟重试"]
    C --> D{"重试次数超限吗"}
    D -- "否" --> E["再次消费"]
    D -- "是" --> F["进入死信队列"]
    B -- "否" --> F
    F --> G["告警、人工处理、补偿任务"]
机制适合
立即重试短暂轻微抖动,但次数要少
延迟重试下游短暂不可用、限流
死信队列数据格式错、业务规则不满足、长期失败
补偿任务从事实源重新生成或修复结果

失败分类很重要:

错误是否适合重试
网络超时适合延迟重试
数据库临时锁等待适合重试
参数格式错误不适合无限重试
外键不存在先确认是否上游延迟,否则死信
业务状态不允许不适合快速重试,应进入异常表

第九步:消息堆积与 Lag

堆积本质:

text
生产速度 > 消费速度

如果长期成立,堆积只会越来越多。

核心公式:

text
消化时间 = 当前堆积量 / (消费 TPS - 生产 TPS)

如果消费 TPS 小于生产 TPS,永远追不上。

Lag 是什么

Kafka 中:

text
Lag = LOG-END-OFFSET - CURRENT-OFFSET

RocketMQ 中常看 MessageQueue 的 Diff。

RabbitMQ 没有完全一样的 offset Lag,但要看:

指标含义
Ready已入队,未投递
Unacked已投递,未确认
Consumers消费者数量
Ack rate确认速率

Lag 分布

不能只看总 Lag。

mermaid
flowchart TD
    A["总 Lag 12000"] --> B["Partition 0: 100"]
    A --> C["Partition 1: 11600"]
    A --> D["Partition 2: 150"]
    A --> E["Partition 3: 150"]

这种不是整体慢,而是局部通道卡住。可能原因:

  1. 热点 key。
  2. 单条慢消息。
  3. 顺序消息阻塞。
  4. 某消费者异常。
  5. 重试风暴集中。

第十步:扩容消费者后又堆积

这个面试题很高频。

现象:

  1. 扩容后消费 TPS 短暂升高。
  2. Lag 开始下降。
  3. 过一会儿 Lag 又上升。

根因通常是:扩容打出了更深层瓶颈。

mermaid
flowchart TD
    A["消费者扩容"] --> B["并发提高"]
    B --> C["更多请求打到下游"]
    C --> D{"下游能否承受"}
    D -- "能" --> E["堆积持续下降"]
    D -- "不能" --> F["DB/ES/接口变慢"]
    F --> G["单条消费耗时升高"]
    G --> H["消费 TPS 再次下降"]
    H --> I["再次堆积"]

排查方向:

现象可能原因
所有分区都高整体消费能力不足或下游慢
少数分区高热点 key、慢消息、单消费者异常
Unacked 高消费者处理慢或 ACK 问题
重试队列高失败消息反复重试
CPU 不高但 Lag 高卡在 DB、锁、接口、连接池
消费者多但无效分区/队列数不足

第十一步:背压怎么设计

背压是下游处理不过来时,让上游或拉取端慢下来。

常见措施:

背压手段
生产者限流、降级、暂停非核心消息
Broker队列长度告警、磁盘水位保护
消费者Kafka pause/resume、RabbitMQ prefetch
线程池有界队列、拒绝策略
下游数据库限流、连接池保护
业务延迟重试、死信隔离、补偿任务

没有背压会怎样?

  1. 消费者本地内存被拉爆。
  2. 线程池队列无限增长。
  3. 数据库连接池耗尽。
  4. Broker 磁盘持续上涨。
  5. 故障从下游扩散到上游。

第十二步:RabbitMQ 从原理上怎么理解

RabbitMQ 核心是 Exchange、Queue、Binding。

mermaid
flowchart TD
    A["Producer"] --> B["Exchange"]
    B --> C["Binding 规则"]
    C --> D["Queue 1"]
    C --> E["Queue 2"]
    D --> F["Consumer A"]
    E --> G["Consumer B"]

适合:

  1. 复杂路由。
  2. 传统任务队列。
  3. 延迟队列、死信队列。
  4. 企业应用集成。

重点:

说明
Exchange决定消息路由
Queue存储消息
ACK消费成功确认
prefetch控制未确认消息数量
DLX死信交换机

RabbitMQ 排查堆积时重点看 Ready 和 Unacked。

第十三步:Kafka 从原理上怎么理解

Kafka 核心是 Topic、Partition、Log、Offset、Consumer Group。

mermaid
flowchart TD
    A["Producer"] --> B["Topic"]
    B --> C["Partition 0<br/>追加日志"]
    B --> D["Partition 1<br/>追加日志"]
    B --> E["Partition 2<br/>追加日志"]
    C --> F["Consumer Group<br/>按分区消费"]

适合:

  1. 高吞吐事件流。
  2. 日志采集。
  3. 用户行为埋点。
  4. Binlog/CDC 数据管道。
  5. 实时计算。
  6. 可回放事件。

重点:

说明
Partition并行和顺序基本单位
Offset消费进度
Consumer Group同组分摊,不同组广播
Retention消息按时间或大小保留
Rebalance消费者变化后重新分配分区

Kafka 的强项是日志模型和高吞吐,不是复杂路由。

第十四步:RocketMQ 从原理上怎么理解

RocketMQ 核心是 Topic、MessageQueue、NameServer、Broker、ConsumerGroup。

mermaid
flowchart TD
    A["Producer"] --> B["NameServer 获取路由"]
    B --> C["Broker"]
    C --> D["CommitLog"]
    C --> E["ConsumeQueue"]
    E --> F["ConsumerGroup 消费"]

适合:

  1. 电商交易消息。
  2. 顺序消息。
  3. 事务消息。
  4. 延迟消息。
  5. 消费重试和死信。
  6. 金融级异步链路。

重点:

说明
NameServer路由发现,不存消息
Broker存储消息
CommitLog消息主存储
ConsumeQueue消费索引
ConsumerGroup消费组
事务消息半消息 + 本地事务 + 回查

RocketMQ 在业务消息、事务消息和顺序消息上更贴近 Java 后端常见交易场景。

第十五步:三大 MQ 选型

对比项RabbitMQKafkaRocketMQ
核心模型Exchange + QueueTopic + Partition + LogTopic + MessageQueue
强项复杂路由、ACK、死信高吞吐、回放、流处理业务消息、事务、顺序、重试
顺序Queue 内有序Partition 内有序Queue 内有序
事务常配合业务表/确认机制Kafka 内部事务,外部用 Outbox/CDC原生事务消息更常用
消息保留偏队列消费按时间/大小保留支持堆积和重试
典型场景企业任务、延迟死信、路由日志、埋点、CDC、实时计算订单、支付、库存、交易事件
排查重点Ready、Unacked、prefetchLag、Partition、RebalanceDiff、MessageQueue、重试、死信

选型建议:

  1. 需要复杂路由和传统任务队列:RabbitMQ。
  2. 需要高吞吐、可回放、数据管道:Kafka。
  3. 需要交易消息、事务消息、顺序消息:RocketMQ。
  4. 团队已有稳定运维经验,优先沿用现有体系。
  5. 不要为了“高级”引入不熟悉的 MQ。

商业场景一:订单创建后异步处理

mermaid
flowchart TD
    A["订单服务本地事务"] --> B["订单落库"]
    B --> C["发送 order.created 消息"]
    C --> D["库存消费者锁库存"]
    C --> E["积分消费者发积分"]
    C --> F["短信消费者发通知"]
    C --> G["ES 消费者同步索引"]

落地要点:

  1. 订单号唯一。
  2. 消息 key 使用订单号。
  3. 生产者发送失败要重试或本地消息表。
  4. 消费者用订单号幂等。
  5. ES 同步失败进入死信和补偿。
  6. 监控 Lag、失败率、死信数量。

商业场景二:医疗数据采集异步入库

mermaid
flowchart TD
    A["采集任务"] --> B["读取接口/文件"]
    B --> C["生成采集批次消息"]
    C --> D["清洗消费者"]
    D --> E["入库消费者"]
    E --> F["资产变更事件"]
    F --> G["ES 同步消费者"]
    F --> H["审计日志消费者"]

落地要点:

  1. taskId + sourceId + batchNo 做幂等键。
  2. 原始数据先落库或对象存储,避免处理失败丢数据。
  3. 清洗失败进入异常表,不无限重试。
  4. 入库用唯一约束防重复。
  5. ES 同步失败可以从数据库事实源补偿。

线上排查总流程

mermaid
flowchart TD
    A["MQ 问题"] --> B{"表现是什么"}
    B -- "消息丢失" --> C["查生产确认、Broker持久化、消费ACK"]
    B -- "重复消费" --> D["查ACK/Offset/重试/Rebalance"]
    B -- "消息堆积" --> E["查生产TPS、消费TPS、Lag分布"]
    B -- "顺序错乱" --> F["查key路由、并发、异步ACK"]
    B -- "消费失败" --> G["查异常、重试、死信、业务数据"]
    E --> H["判断整体慢还是局部热点"]
    H --> I["扩容、限流、优化下游、隔离失败消息"]

排查证据:

证据看什么
生产日志消息是否发送成功,有无重试
Broker 指标队列深度、磁盘、水位、副本
消费日志消费耗时、异常、ACK 时机
Lag 分布是否局部热点或整体慢
重试队列是否失败消息反复重试
死信队列是否有不可恢复消息
下游 DB/ES慢 SQL、写入耗时、连接池
业务表幂等记录、状态机、补偿结果

常见误区

误区为什么错正确做法
MQ 一定不丢消息三个阶段都可能失败发送确认、持久化、ACK、补偿
MQ 能保证不重复多数是至少一次消费端幂等
堆积就加消费者可能瓶颈在分区、下游、热点先看 Lag 分布和耗时
可以直接清队列可能丢业务结果先确认能否从事实源补偿
全局顺序很常见全局顺序吞吐极低设计业务局部顺序
事务消息保证全链路一致只保证本地事务和消息发送一致消费仍需幂等和补偿
死信就是垃圾消息死信是异常证据要告警、分析、修复、补偿

面试标准回答

MQ 有什么作用

text
MQ 主要用于异步解耦、削峰填谷、最终一致和事件驱动。生产者把消息发送到 Broker,Broker 负责存储和投递,消费者按自身能力处理。这样主链路可以更短,下游慢或短暂不可用时不会直接拖垮上游,高峰流量也可以先堆在 Broker 里再慢慢消费。

MQ 怎么保证消息不丢

text
要分三段看。生产阶段要有发送确认、失败重试或本地消息表;Broker 阶段要持久化、复制和合适刷盘策略;消费阶段要在业务真正成功后再 ACK 或提交 Offset,失败进入重试或死信。即使这样也要有补偿和对账,因为 MQ 可靠性是链路设计,不是单个开关。

MQ 为什么会重复消费

text
多数 MQ 更容易保证至少一次投递。消费者处理成功但 ACK 或 Offset 提交失败、网络超时、Broker 重试、消费者宕机、Rebalance、生产者重试都可能导致重复消费。所以消费端必须幂等,常用业务唯一键、消费日志表、唯一索引、状态机或乐观锁保证重复消息不会产生重复副作用。

消息堆积怎么排查

text
先确认 Topic、Queue、ConsumerGroup,再看生产 TPS、消费 TPS、总 Lag、最老消息等待时间和 Lag 分布。所有分区都高通常是整体消费能力不足或下游慢;少数分区高通常是热点 key、慢消息、消费者异常、顺序阻塞或重试风暴。处理时先限流止血,再优化消费者、扩容、批处理、隔离失败消息、处理死信和补偿,不能直接清队列或随便 reset offset。

RabbitMQ、Kafka、RocketMQ 怎么选

text
RabbitMQ 适合复杂路由、传统任务队列、ACK、死信和延迟场景;Kafka 适合高吞吐事件流、日志、埋点、CDC、实时计算和可回放;RocketMQ 适合 Java 业务消息、事务消息、顺序消息、延迟消息、消费重试和交易链路。选型还要看团队运维经验、现有生态、可靠性要求和业务模型。

关联知识点

知识点继续学习
消息队列总览MQ 总览
消息堆积与背压堆积总览
扩容后再次堆积扩容反弹
Lag 分布Lag 分布与定位
RabbitMQRabbitMQ 核心模型
RabbitMQ ACK 死信ACK 重试和死信
KafkaKafka 总览
Kafka 分区 Offset分区顺序与 Offset
Kafka 可靠性可靠性幂等与事务
RocketMQRocketMQ 总览
RocketMQ 顺序消息顺序消息
RocketMQ 事务消息事务消息
分布式事务可靠消息方案
MQ 面试消息队列面试题

本章小结

MQ 从零到生产级掌握,不是背“异步、解耦、削峰”,而是把生产、存储、消费、ACK、Offset、幂等、顺序、事务、重试、死信、堆积、背压、补偿、监控和选型串起来。真正成熟的 MQ 方案一定有可靠发送、消费幂等、失败隔离、死信处理、Lag 监控、补偿对账和容量评估。