AIHub

消息队列:异步、削峰与可靠投递

进阶21 分钟读完2026-08-06#后端#消息队列#Kafka#RabbitMQ
消息队列:异步、削峰与可靠投递

开篇:下单成功,用户却等了 5 秒

一个典型的下单接口要做什么?校验库存、扣减库存、创建订单、扣优惠券、发短信通知、加积分、通知仓库拣货……如果这些全部同步串行在一个 HTTP 请求里完成,每步 100ms,用户就要盯着转圈圈的页面等上好几秒;更糟的是,短信服务商一抖动,整个下单接口跟着超时——明明订单都创建成功了。

问题出在哪?出在"所有事情必须立刻、亲自、按顺序做完"的假设上。现实中,发短信、加积分、通知仓库这些事不需要在下单的那一刻完成,晚几秒没有任何问题。消息队列(Message Queue,MQ)就是为此而生:把"必须立刻做"和"可以稍后做"分开,后者打包成一条消息扔进队列,由别的服务慢慢消费。

MQ 的三大价值:异步、解耦、削峰

生产者、队列与消费者的关系示意

异步化:下单接口扣完库存、建好订单,往队列里丢一条"订单已创建"的消息就立刻返回,RT 从 5 秒降到 200ms。短信、积分由消费者后台慢慢处理。用户体验的改善是直接的。

解耦:没有 MQ 时,订单系统要亲自调用短信服务、积分服务、仓储服务——每加一个下游,订单代码就要改一处;任何一个下游挂掉,下单链路就多一个故障点。有了 MQ,订单系统只负责发消息,根本不知道也不关心谁在消费。新增一个"数据分析服务"想订阅订单事件?直接消费同一个 topic,订单代码一行不动。系统之间从"蜘蛛网状的直接调用"变成"各自只跟队列说话"。

削峰填谷:秒杀开始的瞬间,10 万请求同时涌入,数据库每秒只能扛 5000 笔写入。没有 MQ 就是直接打挂。有了 MQ:请求先快速写入队列返回"排队中",消费者按数据库能承受的速率(比如每秒 4000 条)匀速消费——洪峰被队列这个"水库"拦下,摊平到之后的几十秒里消化掉。队列把瞬时的写压力变成了延时的匀速压力,这是任何缓存都做不到的事。

Kafka 与 RabbitMQ:不是谁更好,是场景不同

两者都是 MQ,但设计目标不同,选型前先想清楚你要什么。

RabbitMQ 是传统的"消息代理"(broker):核心是灵活的路由(direct/topic/fanout 等交换机类型)、丰富的投递语义、单条消息级别的确认。它像一个尽职的邮局,适合业务消息:订单事件、任务派发、RPC 异步化——吞吐量要求万级/秒、但对路由和可靠性语义要求高的场景。

Kafka 本质是分布式提交日志:消息按 topic 分区(partition)顺序追加,消费者自己记录读到哪了(offset),可以反复重放历史消息。它像一卷永不覆盖(保留期内)的磁带,适合海量数据管道:日志收集、行为埋点、流式计算,单机十万级/秒起步的吞吐是它的主场。代价是路由能力弱、单条消息管理粗。

一句话:业务事件用 RabbitMQ,数据洪流用 Kafka。不要为用 Kafka 而用 Kafka——运维一套 Kafka 集群的成本,小团队往往吃不消。

可靠投递与幂等消费:MQ 的真正难点

把消息发出去只是开始。工程上真正的问题是两个:消息会不会丢?会不会重复?

先说结论:消息队列无法保证"恰好一次"(exactly-once)语义,工业标准做法是"至少一次投递(at-least-once)+ 消费端幂等"。

消息从生产者到消费者要经过三段:生产者→队列、队列自身存储、队列→消费者。防丢要三段都做:生产者开启发送确认(RabbitMQ 的 confirm 机制、Kafka 的 acks=all);队列开启持久化(消息和队列都持久化到磁盘,防止 broker 重启丢消息);消费者处理完业务再手动 ack,而不是收到就 ack——否则消费者处理到一半宕机,消息就没了。

那重复从哪来?消费者处理成功、但 ack 在网络中丢了,队列认为没消费,重新投递——消息就被处理了两次。所以消费端必须幂等:同一条消息处理 N 次和处理一次,效果完全相同。常用手段:

// 幂等消费示例:用消息唯一 ID 去重
async function handleOrderCreated(msg) {
  const ok = await redis.set(`mq:dedup:${msg.id}`, '1', 'NX', 'EX', 86400);
  if (!ok) return; // 已处理过,直接丢弃
  await createOrderFollowUp(msg.orderId);
}

或者利用数据库唯一约束(比如"以订单号建积分流水,订单号唯一索引"),重复插入直接失败即视为已处理。幂等设计是消费端代码的标配,不是可选项。

死信队列与延时消息

死信队列(DLQ):一条消息被消费失败重试多次(比如 3 次)后仍失败,怎么办?一直重试会堵住队列,直接丢弃会丢数据。标准做法是把它移入一个专门的"死信队列",让正常消息继续流动,运维或定时任务再人工/自动排查死信——是数据错了还是下游挂了。DLQ 是消息系统的"重症监护室",任何严肃使用 MQ 的系统都必须配置,否则毒消息(poison message)迟早堵住你的主队列。

延时消息:订单 30 分钟未支付自动关闭,就是典型的延时场景。RabbitMQ 可通过 TTL + 死信队列的组合实现(消息先待在有过期时间的队列里,过期后转入真实消费队列),新版本也有延时插件;RocketMQ 原生支持延时等级;Redis ZSet 也能手搓一个简易延时队列。别用轮询数据库(每分钟扫一遍未支付订单)——那是上个时代的做法,又慢又伤库。

消息的顺序与积压:两个容易忽略的问题

顺序性:订单的"创建 → 支付 → 完成"三条消息如果被并发消费,"支付"可能先于"创建"被处理。真相是:全局有序代价极高,通常只需要局部有序。做法是让同一个业务实体的消息路由到同一分区/同一队列(比如按订单 ID 取模),由一个消费者顺序处理,不同订单之间仍然并行。Kafka 的分区机制天然支持这一点;RabbitMQ 则靠单队列单消费者实现。接受"跨实体无序"这个现实,能省去大量无谓的架构纠结。

积压处理:消费者挂了半天,队列里堆了百万条消息,恢复后怎么办?第一,先监控住——**消费积压量(lag)**是 MQ 最重要的告警指标;第二,临时扩容消费者实例(MQ 的消费模型允许水平扩展,Kafka 注意分区数是并发上限);第三,评估能不能"跳过"——对某些有时效的消息(如实时位置上报),陈旧消息直接丢弃比慢慢追更合理。削峰填谷和积压,本质上是同一枚硬币的两面:队列帮你存住了洪峰,你迟早要把它消费完。

实战:订单系统的 MQ 改造

把开头的同步下单改造成事件驱动,完整链路如下:

  1. 下单接口:校验 → 扣库存(数据库原子操作)→ 创建订单 → 发送 order.created 消息 → 立即返回。RT 从 5 秒降到 200ms。
  2. 短信服务:消费 order.created,发通知短信。短信服务商抖动?消息在队列里等着,恢复了再发,下单完全不受影响
  3. 积分服务:消费同一消息,用订单号唯一约束保证幂等,重复消息直接跳过。
  4. 未支付关单:下单时同时发一条 30 分钟延时消息,消费者到点后检查订单状态,未支付则关闭并回补库存。
  5. 异常兜底:所有消费者失败重试 3 次后进入死信队列,配告警通知值班同学。

改造前,五个系统强耦合、一损俱损;改造后,每个环节独立伸缩、独立容错。这就是 MQ 在真实架构中的样子。

常见误区

  • 误区一:引入 MQ 解决一切。 MQ 用复杂度换来了异步和解耦:链路变长、排查变难、要处理重复和乱序。两个服务之间简单的调用,别硬上 MQ。
  • 误区二:忘记幂等。 以为"消息不会重复"是新手的标志性错误,线上迟早教你做人。
  • 误区三:消息体塞大对象。 消息里放 ID 和关键字段即可,完整数据让消费者自己去查;几 MB 的消息会拖垮整个集群的吞吐。
  • 误区四:不监控队列积压。 消费速度跟不上生产速度时,积压(lag)会持续上涨,这是 MQ 系统最重要的告警指标,没有之一。

小结

消息队列的本质是用存储换时间、用缓冲换稳定:异步化降 RT,解耦降维护成本,削峰保系统性命。选型记住"业务事件 RabbitMQ、数据洪流 Kafka";可靠性记住"至少一次投递 + 消费端幂等";死信队列和延时消息是绕不开的两个工程配件。至此,单机到集群的核心组件我们都过了一遍,下一篇进入部署环节,看 Docker 如何把这些服务标准化地跑起来。


系列导航:上一篇《缓存体系:Redis 数据结构到缓存三大坑》(/content/backend-expert-08-redis-cache) · 下一篇《Docker 与现代化部署:从容器到 CI/CD》(/content/backend-expert-10-docker-deploy)

相关教程

消息队列:异步、削峰与可靠投递 | AIHub