KEEL · 龙骨 · A CURRICULUM FOR THE AI ERA
05 · 消息队列:削峰、解耦与可靠性 — keel 龙骨
这一章回答:什么时候该引入消息队列、削峰链路怎么搭、以及消息到底怎么保证不丢不重。
这一章回答:什么时候该引入消息队列、削峰链路怎么搭、以及消息到底怎么保证不丢不重。
消息队列解决的其实是三件不同的事,混为一谈就容易过度设计:异步化(不该让用户等的事)、削峰(瞬时流量超过系统容量)、解耦(下游增减不该影响上游)。
一、削峰链路:以秒杀为例
请求 → 校验与资格判断 → Redis 预扣库存(原子操作)
↓ 成功
投递 MQ(立即返回给用户"排队中")
↓
消费者按自身能力匀速落库 → 写订单、扣真实库存
要点:
- 真实扣减后置:Redis 预扣是为了快速拒绝超额请求,真正的库存扣减在消费者里用第 04 章的乐观锁完成;
- 消费者速率可控:这是削峰的本质——下游按自己的承受力消费,而不是被流量冲垮;
- 结果通知:用轮询或 SSE/WS 告知用户最终结果,不要让请求同步等待(长连接怎么选、怎么撑住,见 《实时推送与 SSE》)。
二、可靠性四要素
任何一处缺失都会导致丢消息,四者要成套配置:
| 环节 | 措施 | 具体配置 |
|---|---|---|
| Broker 持久化 | 交换机、队列、消息都持久化 | durable=True、DeliveryMode.PERSISTENT |
| 生产者确认 | 投递结果可感知 | confirm_delivery(),失败则记录并重投 |
| 消费者确认 | 处理成功才 ACK | 手动 ACK;失败 nack(requeue=True) 或进死信 |
| 消费幂等 | 重复投递不产生副作用 | 业务键去重(见下) |
# 生产者侧的要点:连接池 + 确认 + 持久化
connection = await aio_pika.connect_robust(url, heartbeat=60)
channel = await connection.channel()
await channel.set_qos(prefetch_count=10) # 限制未确认消息数量
await channel.confirm_delivery() # 开启发布确认
await channel.default_exchange.publish(
aio_pika.Message(body=payload, delivery_mode=aio_pika.DeliveryMode.PERSISTENT),
routing_key="order.create",
)
# 消费者侧:手动 ACK / 失败重入或死信
async with message.process(requeue=True): # 异常自动处理
await handle(message.body)
两个容易配错的参数:
prefetch_count:限制单个消费者未确认的消息数。不设的话会把大量消息推给一个慢消费者,造成"看起来都在处理、实际大量超时"。希望严格均摊时设 1,追求吞吐时设 10 左右并按压测调。heartbeat:长时间空闲的连接会被中间设备静默断开。设 60 秒左右,配合connect_robust自动重连。
三、幂等:重复投递一定会发生
MQ 的语义通常是"至少一次",意味着重复投递是正常现象。幂等不能靠"应该不会重复"来回避。
做法是为每条消息定义稳定的业务幂等键:
幂等键 = business_id + message_type # 例如 order_12345 + stock_deduct
处理前先查这个键是否处理过(Redis SET NX 或数据库唯一索引),处理完记录。唯一索引是最好的幂等实现——它让数据库替你把关,不依赖应用逻辑正确。
四、积压治理
队列堆积时按这个顺序处理,不要一上来就加消费者:
- 先看消费速率是否有异常:慢查询、外部依赖超时、异常重试循环;
- 再扩容消费者:注意分区/队列数量限制(Kafka 消费者数不能超过分区数);
- 限流上游:必要时在入口削峰,保护下游;
- 设置兜底:队列最大长度(超限走降级或拒绝)、消息 TTL(过期进死信,避免无限增长);
- 保留可重放能力:死信要能查、能重放,否则故障恢复后数据缺口无法补齐。
五、RabbitMQ 与 Kafka 的取舍
| 维度 | RabbitMQ | Kafka |
|---|---|---|
| 主打 | 可靠投递、灵活路由 | 高吞吐、可回放 |
| 消费语义 | 消费后删除 | 按 offset 保留,可重放 |
| 适合 | 业务事件、任务分发 | 日志流、行为数据、事件溯源 |
| 并行上限 | 队列维度 | 分区维度(分区数决定消费者上限) |
动手:可观察结果
| 产出 | 判断标准 |
|---|---|
| 一条削峰链路 | 10 倍峰值流量下,下游写库速率保持平稳,库存不超卖 |
| 可靠性测试 | kill 掉消费者进程,重启后消息仍在且不重复发货 |
| 幂等验证 | 手工重复投递同一条消息,业务只生效一次 |
| 积压处理记录 | 人为制造堆积,按五步处理并记录每步效果 |
完成标志:能在管理界面看到秒杀期间队列堆积后平稳消化,且最终库存与订单数一致。
故障注入
| 注入方式 | 观察 |
|---|---|
不设 prefetch_count |
消息是否大量堆积在某个慢消费者、处理超时率是否上升 |
| 关闭持久化并重启 Broker | 未消费的消息是否丢失 |
| 关闭发布确认 | 投递失败时是否无人知晓 |
| 消费成功前就 ACK | 处理失败的消息是否永久丢失 |
| 去掉幂等键 | 重复投递是否导致重复扣款/重复发放 |
| 不给队列设长度与 TTL | 长时间堆积是否撑爆内存或磁盘 |
自测题
- 为什么"至少一次"投递下幂等是必须的?幂等键怎么设计才稳定?
prefetch_count设得太大和太小分别会发生什么?- 手动 ACK 应该在业务的哪个位置?为什么不能在收到消息时就 ACK?
- 队列积压时的处理顺序是什么?为什么先看消费速率而不是先扩容?
- 什么情况下应该用 Kafka 而不是 RabbitMQ?举一条判断依据。