消息队列最容易被误解为“把数据放进去,消费者再取出来”。这个说法没有错,但遗漏了工程上的关键:消息队列是在上游和下游之间加入一个可暂存、可确认、可重试、可观测的异步边界。
它解决的不是某一个产品的配置问题,而是分布式系统中反复出现的矛盾:处理时间不同、处理能力不同、可用时间不同,以及一个业务事实需要被多个系统独立处理。
本文先建立 MQ 的通用判断框架,再比较 RabbitMQ、Kafka、RocketMQ 的核心模型;所有可靠性和故障场景则统一以 RabbitMQ 为例。这样,ACK、死信队列、prefetch、Quorum Queue 等术语都有明确的业务语境。
什么是消息队列
可以用一句话概括 MQ 的取舍:以存储换时间、以异步换解耦、以最终一致性换吞吐量。
消息队列可以先理解为生产者和消费者之间的中间层。生产者只负责说明“发生了什么”或“需要做什么”;消费者在合适的时间取得消息并执行处理。
基本工作流程
生产者 ──发布──→ Broker / MQ ──投递──→ 消费者
│ │
│ ├─ 业务事务
├─ 路由、持久化 ├─ ACK / 重试
├─ 副本、死信 └─ 幂等、补偿
└─ 队列或日志
以订单创建为例,订单服务不再同步等待库存、积分、通知、风控全部完成,而是发布 OrderCreated 这一已发生的事实。下游系统按自己的速度和可用性消费它。
订单服务 → OrderCreated → MQ
├─ 库存服务
├─ 积分服务
├─ 通知服务
└─ 风控服务
这不是消除复杂性,而是将复杂性从一次同步请求转移到一条可恢复的异步链路。
为什么它不只是普通 Queue
语言中的 Queue 常常只是内存里的 FIFO 容器:enqueue → dequeue。分布式 MQ 还必须处理跨网络传输、磁盘持久化、多副本、消费进度、消费者故障转移、确认与重试、权限、路由和监控。因此 Broker 更像专门保存和分发事件的分布式系统。
不要混淆三个层次:
| 层次 | 典型对象 | 通信范围 | 主要能力 |
|---|---|---|---|
| 线程间 | BlockingQueue | 一个进程 | 任务交接与阻塞协调 |
| 进程间 | POSIX/System V 消息队列 | 一台机器 | 内核管理的进程通信 |
| 服务间 | RabbitMQ、Kafka、RocketMQ | 跨机器、跨服务 | 存储、复制、路由、恢复与观测 |
操作系统队列和分布式 MQ 都通过中间缓冲区让发送方与接收方不必同时执行,但前者受单机内核容量限制,后者需要面对网络故障、集群存储和业务恢复。
消息、命令与事件
消息是传输载体;命令和事件是其业务语义。
命令:请扣减库存、请发送短信
事件:订单已支付、用户已注册
命令隐含一个希望被执行的动作;事件描述已发生的事实。跨系统发布时,事实事件通常更容易扩展:以后增加审计、画像或推荐消费者时,无须回头修改订单服务。消息队列降低了服务地址和调用顺序的耦合,但双方仍耦合于事件名、字段语义、版本、顺序和交付语义,因此消息 Schema 也需要版本演进。
MQ 解决哪些问题
异步处理:缩短关键路径,不减少工作量
假设创建订单耗时 50 ms,扣库存 100 ms、积分 80 ms、短信 200 ms。同步链路约需 430 ms;若订单服务可靠地发布事件后立即返回,用户响应可接近 60 ms,其余工作由消费者完成。
异步减少的是请求线程的等待和关键路径长度,不是总计算量。序列化、网络、Broker 落盘、重试、状态跟踪都会增加新的成本。短信、报表、索引更新、图片转码适合异步;支付密码校验、库存是否足够等会决定当前事务结果的步骤,不能只为更快响应而盲目异步。它类似操作系统把重活放到工作队列:前台快速记录任务,后台稍后处理。
削峰与限流:蓄水池,不是扩容器
令生产速率为 P、消费速率为 C。当 P > C 时:
积压增长速率 = P - C
若秒杀瞬间写入 10,000 条/秒,下游仅能处理 4,000 条/秒,每秒会新增 6,000 条积压。MQ 能吸收短期洪峰,避免数据库和下游被瞬间压垮;但如果生产速率长期大于消费速率,它只是延迟系统失效。
削峰能成立至少要满足:高峰是暂时的、低峰能追回积压、消息允许等待、Broker 容量足够。它与网络路由器的缓冲区相似:出口带宽长期不足时,缓冲区最终会满,随后只能拒绝或丢弃流量。
解耦、故障隔离与多订阅
没有 MQ 时,PLM 可能直接调用数据中心、ERP、ISCP;增加一个下游就要改上游。改为发布“物料已创建”后,上游只关心发生了什么。
MQ 提供三类解耦:
- 空间解耦:生产者不必知道消费者地址。
- 时间解耦:双方不必同时在线。
- 数量解耦:生产者不必知道有多少订阅者。
当 ERP 暂时不可用,消息可先留在 Broker,恢复后继续消费,避免 ERP 故障直接导致 PLM 创建物料失败。发布订阅还能让一份事件可靠分发给多个下游,每个下游维护自己的处理状态。它可以看作跨进程、可持久化、带重试的观察者模式,而不是进程内同步回调。
其他方案能否替代
| 方案 | 能覆盖什么 | 主要边界 |
|---|---|---|
| 同步 RPC | 立即获取下游结果 | 下游变慢或故障会传递给上游 |
| 数据库中间表 | 暂存、扫描、补发 | 多订阅、并发、进度、死信和积压治理需要自建 |
| 定时任务 | 周期处理、补偿 | 实时性较低,频繁扫描会压数据库 |
| Redis List/Stream | 简单队列、轻量事件流 | 持久化、事务、故障恢复需谨慎设计 |
| 专业 MQ | 异步、削峰、路由、可靠分发 | 引入运维、契约和一致性复杂度 |
所以问题不是“只有 MQ 能做到什么”,而是当需求同时包含跨系统、高并发、多消费者、可靠恢复和积压管理时,是否值得让每个业务系统重复实现一套通用机制。
MQ 还能承载哪些业务语义
分布式事务与最终一致性
最危险的窗口是本地数据库已提交但消息发送失败,或消息已发送但本地事务最终回滚。目标通常不是让所有服务进入一个大 ACID 事务,而是保证“本地事实成立”与“对应事件最终可见”不长期矛盾。
支付服务本地事务:更新支付单 + 写入 PaymentSucceeded Outbox
↓
投递器发布事件
↓
订单、会计、营销服务分别处理
常见手段包括本地消息表/Transactional Outbox、CDC 读取数据库变更、Broker 事务消息、Saga 补偿。它们共同遵循“先可靠记录意图,再异步完成后续动作”。这和数据库先写 WAL 再修改数据页的思想相通。最终仍要有对账与补偿,因为消息系统无法替业务判断“钱是否真的到账”。
顺序:通常只需要按业务实体有序
订单创建、支付、发货、完成不能颠倒,但没有必要让所有订单共享一个串行通道。全局有序像一把全局锁,正确但吞吐最低;更实用的是局部有序:
同一 orderId:创建 → 支付 → 发货
不同 orderId:允许并行
Kafka 通常用 Key 将同一实体放入同一 Partition;RocketMQ 可用 FIFO 消息组。RabbitMQ 也应围绕同一个业务键设计路由、单消费者边界或分片策略。可按 orderId、accountId、userId、deviceId 分组;若所有消息使用同一个 Key,系统会退化成单通道,若 Key 过粗则会产生热点。
延时、即时通信与数据流
延时消息适合订单超时关闭、自动确认收货、会员到期提醒、延时重试:创建订单后投递“30 分钟后检查”的消息。延时时间到达只表示消息可以投递,Broker 调度粒度、积压、消费者负载和网络都会影响实际处理时刻,所以它不是精确定时器。
即时通信中,WebSocket、TCP、MQTT 解决客户端怎么连接;MQ 解决服务端如何分发、保存和处理事件。典型 IoT 路径是“设备 → MQTT/接入网关 → MQ → 实时计算、告警、存储”。用户连接映射、离线消息、多端同步、弱网重连仍属于连接和业务层的职责。
数据流场景关注的不只是“任务是否完成”,还包括“发生过什么、读到哪里、能否重算”。日志、埋点、CDC 可进入 Kafka 或 RabbitMQ Stream,再被 Flink、Spark、数仓和实时指标消费。它与 binlog、WAL、事件溯源共享一个观念:状态是结果,日志保存变化过程;保留变化过程,才能回放并重建状态。
RabbitMQ 基础知识:先认识最小模型
RabbitMQ 是一个消息 Broker。应用程序既可以作为 Producer 向它发布消息,也可以作为 Consumer 从它接收消息;Broker 则负责维护拓扑、路由、存储和投递状态。下面这组概念是阅读后文的前提。
Virtual Host、用户与权限
一个 RabbitMQ 集群可以划分多个 Virtual Host(vhost)。vhost 类似消息资源的逻辑租户或命名空间:不同 vhost 内可以存在同名的 Exchange、Queue 和 Binding,用户权限也按 vhost 授予。开发、测试、生产环境,或不同业务租户,通常应使用不同 vhost,而不是将所有资源放在默认 vhost 中。
RabbitMQ Broker
├─ /dev → 开发 Exchange、Queue、用户权限
├─ /test → 测试 Exchange、Queue、用户权限
└─ /prod → 生产 Exchange、Queue、用户权限
六个核心角色
Producer
│ publish(message, exchange, routingKey)
▼
Exchange ──Binding rule──→ Queue ──deliver──→ Consumer
│ │
└─ 保存消息 └─ ack / nack / reject
- Producer:构造并发布消息;通常指定 Exchange 和 Routing Key。
- Exchange:不保存待消费消息,负责按规则将消息路由到一个或多个 Queue。
- Binding:Exchange 到 Queue 的连接及其匹配规则。
- Routing Key:生产者带上的路由地址;是否参与匹配取决于 Exchange 类型。
- Queue:保存待投递的消息;多个消费者消费同一 Queue 时属于竞争消费。
- Consumer:接收并处理消息,再用确认结果告诉 Broker 如何处置本次投递。
一个基础误区是认为“Producer 直接把消息放进 Queue”。在 AMQP 0-9-1 中,常见路径是 Producer 先发布给 Exchange;即使使用默认 Exchange,本质上也仍经过了一次按队列名路由。若消息不能匹配任何 Queue,是否返回生产者、丢弃或交给备用拓扑,应在发布策略中显式设计。
Exchange 的四种类型
| 类型 | 匹配方式 | 常见用途 |
|---|---|---|
| Direct | Routing Key 与 Binding Key 完全相等 | 按事件类别、日志级别分发 |
| Fanout | 忽略 Routing Key,广播给所有绑定 Queue | 用户注册、配置变更、在线通知 |
| Topic | * 匹配一个词,# 匹配零或多个词 | 地区、业务、事件类型组合订阅 |
| Headers | 按消息 Header 的键值匹配 | 路由维度不适合放进字符串时 |
例如事件键 order.cn.created 可以由 order.*.created 匹配,也可由 order.# 匹配。Topic 的“词”由 . 分隔;* 不会匹配两段,# 则可匹配零段或多段。
Queue 的声明属性
Queue 不只是名称,还带有生命周期和可靠性属性:
- durable:Broker 重启后保留队列定义;要同时配合持久消息才能降低消息丢失风险。
- exclusive:仅创建它的连接可用,连接断开后删除;适合临时回复队列,不适合关键业务队列。
- auto-delete:最后一个消费者取消订阅后自动删除;适合临时订阅。
- TTL:可为消息或队列设定存活时间;过期消息可被丢弃或进入死信路径。
- 最大长度/字节数:限制 Queue 可承受的积压,并配合溢出策略保护 Broker。
这些属性属于拓扑契约。同名 Queue 若以不同属性重复声明,RabbitMQ 会认为定义冲突并关闭当前 Channel;因此应将拓扑声明集中管理,避免不同服务各自猜测配置。
消息属性、持久化与确认
消息体之外还可带 Content Type、Headers、Correlation ID、Message ID、Timestamp、Expiration、Priority、Delivery Mode 等属性。Correlation ID 常用于请求/响应关联,Message ID 适合追踪但不能替代业务幂等键,Priority 只有在队列声明支持优先级时才生效。
持久队列只表示队列元数据可恢复;消息是否可恢复还取决于消息的持久化投递模式,以及 Broker 是否已经确认接管。消费者收到消息后有三类常见处置:
basic.ack → 成功,Broker 移除该消息
basic.nack/reject + requeue → 失败,重新入队或进入死信路径
不确认且消费者连接断开 → 未确认消息通常会被重新投递
自动确认模式会在投递后立即视为成功,适合可丢失的实时通知,不适合需要写数据库、调用支付或发券的关键业务。手动确认让消费者能在业务成功后再 ACK,但也要求处理重复投递。
一次最小发布与消费流程
1. 客户端连接到指定 vhost,建立 Connection 与 Channel
2. 声明 Exchange、Queue,并建立 Binding
3. Producer 发布 message + routing key
4. Exchange 按类型和 Binding 找到目标 Queue
5. Queue 按可用消费者和 prefetch 投递消息
6. Consumer 执行业务本地事务
7. 成功 ack;临时失败重试;不可恢复失败进入死信或人工流程
后文的可靠投递、积压、优先级与队列类型,都是对这条最小路径中某一个环节的细化。
三种 MQ 的核心抽象
以同一条订单创建事件为例:
{
"event": "order.created",
"orderId": "1001",
"region": "cn",
"amount": 299,
"vip": true
}
| 产品 | 核心路径 | 核心问题 |
|---|---|---|
| RabbitMQ | Exchange → Queue → ACK | 这条消息应该投递给哪些队列? |
| Kafka | Topic → Partition → Offset | 事件写到哪里、顺序如何、读到哪里? |
| RocketMQ | Topic → Message Queue → Offset | 如何并行处理业务消息,并提供顺序、延时、事务语义? |
RabbitMQ:路由优先
生产者通常不直接面向 Queue,而是将消息发布给 Exchange。Exchange 根据 Binding 与 Routing Key 决定向哪些 Queue 复制消息。
Producer ──order.created.cn──→ Topic Exchange
├─ inventory.queue
├─ points.queue
└─ cn.analytics.queue
direct 做精确匹配,fanout 广播,topic 以 . 分词,* 匹配一个词、# 匹配零个或多个词;headers 可按 Header 匹配。RabbitMQ 的重点是“满足条件的消息应送给谁”。多个消费者消费同一个 Queue 是竞争关系,每条消息只给其中一个;一个 Exchange 绑定多个 Queue 才是发布订阅,每个 Queue 各有一份消息。
Kafka:追加日志优先
Kafka 将事件追加到 Topic 的某个 Partition,消费者以 Offset 记录位置:
order-events / Partition 1
offset 0 → OrderCreated(1001)
offset 1 → OrderPaid(1001)
offset 2 → OrderShipped(1001)
同一 orderId 作为 Key 通常进入同一分区,从而获得局部顺序。消息不会因某个消费者读取而立刻删除;不同 Consumer Group 各自维护 Offset,可以独立消费或回放。Kafka 的 Broker 不擅长复杂业务路由,库存系统若只关心创建事件,通常自行过滤、拆分 Topic,或通过 Kafka Streams 产生派生 Topic。
RocketMQ:业务消息语义与二级过滤
RocketMQ 也以 Topic 下多个 Message Queue 实现并行存储和 Offset 消费,并提供 Tag 与属性过滤:
Topic = OrderCreated
Tag = VIP
region = cn
库存服务可订阅全部,VIP 服务订阅 Tag,中国区统计用属性条件。它比 Kafka 原生消费更方便表达业务筛选,但复杂路由拓扑通常不如 RabbitMQ 灵活。
同一需求下,RabbitMQ 可用 order.created.cn.vip 和三个 Topic Binding 自动复制消息;Kafka 可由三个消费者组自行过滤,或借助流处理拆 Topic;RocketMQ 则将筛选表达为 Tag 与属性。路由键、分区键与 Offset 不是可互换概念:Routing Key 影响 RabbitMQ 投递目标,Kafka Key 主要影响分区、顺序与并行度,Offset 表示消费进度。
RabbitMQ 的协议、连接与拓扑
AMQP 0-9-1 与 AMQP 1.0
AMQP 0-9-1 明确定义了 Exchange、Queue、Binding、Routing Key、发布与消费操作:
Producer → Exchange → Binding → Queue → Consumer
AMQP 1.0 不是其简单升级版。它采用 Connection → Session → Link,以 Sender Link 与 Receiver Link 在 Source 和 Target 之间传输消息,Target 是队列、主题还是日志资源由 Broker 解释。因此 AMQP 1.0 更像可跨产品互操作的通用消息传输模型。
AMQP 0-9-1 常用 prefetch 限制未确认消息数量;AMQP 1.0 由 Receiver 给 Sender 发放 credit,细粒度控制一条 Link 能接收的交付数量。AMQP 1.0 还定义 accepted、rejected、released、modified 等交付状态及 settlement:立即结算接近 at-most-once,接收方确认后结算接近 at-least-once。它们都不能自动让任意数据库和第三方副作用变成 exactly-once。
AMQP 1.0 的消息可含 Header、Delivery Annotations、Message Annotations、Properties、Application Properties、Body、Footer,Body 可为二进制、AMQP Value 或 AMQP Sequence;0-9-1 更偏属性加二进制消息体。
五种常见模式
| 模式 | Exchange/拓扑 | 消费关系 | 用途 |
|---|---|---|---|
| 简单模式 | 默认 Exchange → 一个 Queue | 一个消费者 | 简单异步任务 |
| Work Queue | 一个 Queue → 多个 Worker | 竞争消费 | 横向扩展耗时任务 |
| Publish/Subscribe | Fanout → 多个 Queue | 每个 Queue 都收到 | 广播事件 |
| Routing | Direct → 多个 Queue | 精确匹配 | 日志级别、类别分发 |
| Topic | Topic → 多个 Queue | 通配匹配 | 多维事件订阅 |
简单模式实际仍经过默认 Exchange:Queue 声明后会用自己的名称绑定默认 Exchange,生产者以队列名为 Routing Key 发布。Work Queue 应结合手动 ACK、持久队列、持久消息和合理的 prefetch;它的目的不是广播,而是把每条任务交给一个可用 Worker。Fanout 忽略 Routing Key,向每个绑定 Queue 复制消息。Topic 中 order.*.created 匹配三段式创建事件,order.# 匹配任意层级订单事件。
Connection 与 Channel
TCP/TLS 连接需要文件描述符、内核缓冲、心跳、网络状态和 TLS 握手,不应每次发布都建立。AMQP 0-9-1 在一个长期 Connection 上复用多个 Channel:
TCP Connection
├─ Channel 1:发布订单
├─ Channel 2:消费日志
└─ Channel 3:确认消息
帧包含 Channel ID,因此同一连接上可以交错传输 [ch1][ch2][ch1][ch3]。每个 Channel 保存自己的消费者、unacked 消息、prefetch、Publisher Confirm、事务和异常状态。声明同名但属性不同的队列、绑定不存在的 Exchange 等协议错误通常只会关闭当前 Channel,而不是整个 Connection。
Channel 一般不应由多个线程并发共享:帧顺序、确认与状态容易混乱。合理方式是长期 Connection 加每线程或每发布者独立的长期 Channel,或使用安全的 Channel Pool;不要每发一条消息就创建和关闭 Channel。
RabbitMQ 可靠投递:从发布到业务完成
可靠性不能只看“消息是否持久化”,而要检查完整链路:
Producer → 网络 → RabbitMQ → 网络 → Consumer → 业务数据库
发布端:Publisher Confirm 与重复发送
生产者调用发布 API 后,网络可能在 Broker 收到前后中断。Publisher Confirm 用于确认 Broker 已承担这条消息,但确认响应也可能丢失:Broker 已收,生产者却以为失败并重发。因此可靠发布天然会产生重复的可能,需要发送重试与消费者幂等协同设计。
本地事务与发布之间的窗口可通过 Outbox 处理:业务数据和待发布事件在同一个本地事务写入,投递器再带重试地发送。对关键队列,还要明确消息持久化、刷盘、副本以及何时向生产者确认。
Broker:队列、复制与存储责任
持久化队列和持久消息能降低进程重启导致的数据丢失,但不等于自动高可用。Classic Queue 默认没有复制;Quorum Queue 使用 Raft,让 Leader 将日志复制给 Follower,在多数副本确认后提交。更强的可靠性意味着更多磁盘、网络和 CPU 成本,也会提高延迟。
消费端:为什么必须在业务成功后 ACK
错误顺序是:
收到消息 → ACK → 写数据库 → 进程崩溃
此时 RabbitMQ 已删除消息,业务却没有完成。常见顺序是:
收到消息 → 执行业务本地事务 → 成功后 basic.ack
但反向窗口仍然存在:数据库已提交、ACK 尚未送达时消费者崩溃,RabbitMQ 会重新投递。因此 Broker 的 ACK 表示“本次投递已经处理”,不等于“世界上绝不会再出现同一业务事件”。这正是 RabbitMQ 场景中重复消费的来源。
重复消费与幂等
RabbitMQ 的至少一次投递语义优先保证“不轻易丢”,代价是“可能重复”。应使用业务唯一键,而不是仅依赖 Broker 的投递 ID。
| 方法 | 做法 | 适用与限制 |
|---|---|---|
| 唯一索引 | 为 orderId、支付流水号等建唯一约束 | 去重判断与写入可在数据库原子完成 |
| 状态机 | UNPAID → PAID → SHIPPED | 防止重复事件让状态倒退或重复推进 |
| 消费记录表 | 同一事务写消费记录和业务数据 | 适合需要审计消费状态的场景 |
| Redis 去重 | SET key value NX EX | 适合时间窗口;要考虑 TTL、故障和非原子边界 |
幂等不是“代码绝不运行第二次”,而是同一操作运行一次或多次,最终业务结果相同。将订单状态设为 PAID 比“支付次数加一”容易幂等。发券、扣款、第三方调用等外部副作用,还应把业务幂等键继续传递给下游;它相当于异步系统中的 HTTP Idempotency-Key。
失败、重试与死信
业务临时失败时,消费者可以 nack 或 reject 并选择是否 requeue。但无条件 requeue=true 容易让一条毒消息立刻反复投递,耗尽消费者。更常见的策略是有限次数重试、延时重试队列,以及超过阈值后路由到死信队列。
死信队列不是垃圾桶,而是自动恢复失败后等待诊断、修复和补偿的隔离区。完整的业务兜底还包括对账任务、补偿任务、告警和人工处理入口。RabbitMQ 能帮助保存与重投递消息,但只有业务终点能判定库存是否扣减、钱是否到账、短信是否真的送达。
RabbitMQ 消息积压:从 Ready、Unacked 到恢复
消息积压不是根因,而是生产速度大于有效消费速度的症状。在 RabbitMQ 中,排查不能只看 Queue 总消息数,还要区分:
- Ready:仍在队列内,尚未投递给消费者的消息。
- Unacked:已经投递给消费者、尚未确认的消息。
- 发布速率与确认速率:持续失衡时积压会增长。
- 最老消息年龄:比单纯数量更贴近业务时效。
- 失败率、重投递率、死信量:判断是否存在毒消息或重试风暴。
若消费者稳定处理 2,000 条/秒、生产者仍写入 1,500 条/秒,净清理速率只有 500 条/秒。恢复时间要按下面计算,而不是只看消费吞吐:
净清理速率 = 消费速率 - 生产速率
清空时间 ≈ 当前积压 ÷ 净清理速率
先找根因
常见原因包括:
- 生产突然增加:大促、全量重跑、上游死循环、重试风暴、爬虫或攻击。
- 消费者变慢:慢 SQL、数据库连接池耗尽、第三方超时、锁竞争、GC 停顿、单条消息过大。
- 有效并行度不足:Worker 数太少,或队列、路由、顺序约束使任务无法并行。
- 消息分布不均:热点业务键让少数队列或消费者承担大部分流量。
- Broker 受限:磁盘、网络、内存水位、节点或副本异常拖慢投递。
prefetch 是 RabbitMQ 的重要变量。它限定一个消费者可以持有多少未确认消息:过大时,慢消费者会把大量消息预取到本地,其他消费者无事可做;过小时,网络往返会降低吞吐。更关键的是,已经进入 Unacked 的普通消息无法被后来到达的高优先级消息插队。
止血、恢复与复盘
处理顺序通常是:
控制新增积压
→ 保护数据库和第三方下游
→ 修复慢消费根因
→ 提高有效并行度
→ 清理历史积压
→ 补齐监控、容量和回放预案
临时增加消费者前,要确认下游数据库能承受、业务允许并行、顺序不会被破坏。若瓶颈本身是数据库,盲目扩消费者只会使故障更严重。批量消费和批量写入可减少网络往返、SQL 解析与事务提交,但批次越大,失败重试成本、等待时间和内存占用也越高。
严重积压时可临时转储到文件或新 Stream/Topic,先释放在线队列压力,再按受控速率重放。转储必须保留业务唯一键、原始时间、消息版本、原队列/路由信息和重放批次,否则重放可能重复更新,甚至让旧事件覆盖新状态。
不同业务的处理策略不同:通知可合并、降级或丢弃过期内容;行为日志可批量压缩与延迟处理;订单、支付、库存不能随意丢弃,必须保留幂等、顺序、对账和人工补偿;实时推荐中的旧事件可能价值快速降低,更适合直接计算最新状态。
RabbitMQ 队列类型与优先级
Classic、Quorum 与 Stream
| 类型 | 消费后消息 | 可靠性和用途 |
|---|---|---|
| Classic Queue | ACK 后删除 | 灵活、成本较低,默认无复制 |
| Quorum Queue | ACK 后删除 | Raft 多副本,适合关键任务 |
| Stream | 按保留策略删除 | 追加日志、Offset、重复读取与回放 |
Classic Queue 支持非持久、独占、自动删除、TTL、优先级等,适合临时队列、reply queue、开发环境或允许较低可靠性的任务。节点故障时,默认无复制的未保护数据可能丢失。
Quorum Queue 有 Leader 和 Follower,Leader 写入后复制到多数节点,再确认提交;Leader 故障时可选举新 Leader。它适合支付、订单、库存等关键命令,代价是更高的网络、磁盘与 CPU 开销;不适合非持久或 exclusive 队列需求。
Stream 是持久化追加日志:消息被某个消费者读取后不会立刻删除,多个消费者分别维护 Offset,可用于实时消费、故障恢复、报表重算和审计回放。它适合日志、埋点、行为事件、CQRS、事件溯源和高吞吐流,不适合大量短生命周期临时队列。
三个例子可以帮助选择:支付处理使用 Quorum Queue;在线用户状态用 Fanout Exchange 分发到多个队列;用户行为分析使用 Stream,让实时推荐、报表、风控和离线回放分别读取。
优先级队列的边界
Classic Queue 可用 x-max-priority 声明最大优先级,内部可近似理解为维护多个优先级子队列:高优先级先被取走,同级消息仍大致 FIFO。优先级范围越大,需要的调度、CPU 和内存越多,业务通常用紧急/高/普通/低四级或 1 到 10 级已经足够。
优先级只影响仍留在 Queue 中、尚未投递的消息。若 prefetch=1000,普通消息已大量处于消费者的 Unacked 缓冲区,后来的高优先级消息不能把它们抢回来。因此优先级队列应配合较小 prefetch。Quorum Queue 的严格优先级会更坚决地优先投递高优先级消息,但持续高优先级流量也可能让低优先级任务饥饿。
最终决策:先问问题,再选机制
引入 MQ 或设计 RabbitMQ 拓扑前,依次回答:
- 任务能否稍后完成?当前请求是否必须等待下游结果?
- 高峰是短暂的还是长期能力不足?消息允许等待多久?
- 事件需要广播、复杂路由,还是保存与回放?
- 真正需要的是全局顺序,还是按订单、账户、用户的局部顺序?
- 发布确认、持久化、副本和消费者 ACK 分别如何设置?
- 重复投递时,业务唯一键、状态机和消费记录如何设计?
- 重试上限、延时重试、死信、对账和人工补偿如何闭环?
- Ready、Unacked、消息年龄、消费速率、重投递率和死信量是否可观测?
看完本文后,应当能够回答:为什么 MQ 能削峰却不能修复长期容量不足;RabbitMQ、Kafka、RocketMQ 分别在抽象什么;为什么可靠投递需要幂等;以及 RabbitMQ 出现重复消费、积压或优先级失效时应该从哪里开始排查。