ARTICLE DETAIL

建站实战干货

来自一线的建站与推广经验沉淀,每一条都经过真实交付验证。

秒杀系统架构设计06:MQ 异步架构

2026/8/10 14:01:10 拓冰建站 浏览量
秒杀系统架构设计06:MQ 异步架构 秒杀的 MQ 异步架构与消息可靠性本文是「10Wqps 秒杀架构」系列的第六篇聚焦 RocketMQ 在秒杀架构中的角色——不只是削峰解耦更重要的是保障消息零丢失和幂等消费。一、开篇在 异步下单架构篇 中我们讨论了为什么秒杀需要异步以及三阶段异步如何将响应时间降到 30ms。随后的 数据一致性架构篇 里Redis 预扣成功后下一步是通过 MQ 发送下单消息。就这一步如果消息丢失了就会出现库存扣了但订单没生成——也就是掉单。掉单对业务的伤害比超卖更深超卖可以通过补偿发货、退款解决掉单则是用户明明抢到了却没有订单投诉升级、信任崩塌。本文以 RocketMQ 为主线从消息的发送、传输、消费三个环节逐一剖析可靠性保障并延伸到幂等性、垃圾消息和消息积压的应对方案。二、问题拆解2.1 消息丢失的三个环节一条消息从生产者到消费者会经历三个阶段每个阶段都可能丢失生产者 ──(1)发送──→ Broker ──(2)存储──→ 磁盘 ──(3)消费──→ 消费者环节 1 — 发送丢失网络抖动、Broker 宕机导致生产者发送失败环节 2 — 存储丢失Broker 刷盘前宕机内存中的消息全部丢失异步刷盘环节 3 — 消费丢失消费者收到消息后在业务处理完成前崩溃消息未 ACK2.2 只靠 MQ 自带机制够吗RocketMQ 提供了事务消息、同步刷盘、重试队列等机制来应对上述问题。但在真实生产环境中百密一疏——任何环节都可能出现预期外的异常。比如事务消息的回查机制依赖网络网络分区时状态不确定同步刷盘有性能代价高并发下可能选择异步刷盘消费者代码的 Bug 可能导致消费逻辑成功但未 ACK因此MQ 自带机制是技术防线业务层面的兜底方案本地消息表 定时扫描是终极保障。三、核心方案3.1 消息零丢失的三条防线防线一RocketMQ 事务消息保障环节 1事务消息确保本地操作和消息发送的原子性——要么都成功要么都不执行。RocketMQ Broker生产者秒杀服务RocketMQ Broker生产者秒杀服务alt[本地事务成功][本地事务失败]如果第 4 步因网络问题未到达 Broker1. 发送 half 消息半消息消费者不可见2. half 消息发送成功3. 执行本地事务Redis 库存预扣4. COMMIT消息对消费者可见4. ROLLBACK消息被丢弃5. 回查Callback Check本地事务状态6. 返回 COMMIT 或 ROLLBACK关键点Half 消息发送成功后消费者不会立即消费——只有 COMMIT 后消息才可见如果生产者未发送 COMMIT/ROLLBACK网络超时或进程 crashBroker 会定期回查回查接口checkLocalTransaction由业务方实现返回真实的事务状态// 发送事务消息TransactionMQProducerproducernewTransactionMQProducer(seckill-group);producer.setTransactionListener(newTransactionListener(){OverridepublicLocalTransactionStateexecuteLocalTransaction(Messagemsg,Objectarg){try{// 执行 Redis 库存预扣Lua 原子操作deductStock((SeckillOrder)arg);returnLocalTransactionState.COMMIT_MESSAGE;}catch(Exceptione){returnLocalTransactionState.ROLLBACK_MESSAGE;}}OverridepublicLocalTransactionStatecheckLocalTransaction(MessageExtmsg){// 回查根据消息中的订单 ID 查询 Redis 预扣记录是否存在StringorderIdmsg.getKeys();if(redis.exists(pre:order:orderId)){returnLocalTransactionState.COMMIT_MESSAGE;}returnLocalTransactionState.ROLLBACK_MESSAGE;}});防线二同步刷盘 主从复制保障环节 2RocketMQ 支持两种刷盘模式模式行为可靠性性能异步刷盘消息写入 PageCache 即返回成功Broker 宕机会丢失未刷盘消息高同步刷盘消息写入磁盘后才返回成功极高低约降 50%秒杀场景对性能要求极高但消息丢失又是不可接受的。折中方案同步刷盘 主从同步复制消息写入 Master 磁盘并同步到 Slave 后才返回成功性能代价可接受因为 MQ 处理的是下单消息QPS 远低于秒杀请求同步刷盘的性能代价被削峰稀释了# Broker 配置 flushDiskTypeSYNC_FLUSH # 同步刷盘 brokerRoleSYNC_MASTER # 同步主从复制防线三手动 ACK 消费重试保障环节 3消费者在完成所有业务操作订单入库、库存确认、积分更新等后才发送 ACK。中途失败不 ACKBroker 会自动重投。RocketMQMessageListener(topicSECKILL_ORDER_TOPIC,consumerGrouporder-create-group)publicclassOrderCreateConsumerimplementsRocketMQListenerMessageExt{OverridepublicvoidonMessage(MessageExtmessage){try{SeckillOrderorderparseMessage(message);// 1. 幂等校验if(orderService.isDuplicate(order.getOrderId())){return;// 重复消费直接返回会 ACK}// 2. 订单入库 库存确认同一事务orderService.createOrderInTransaction(order);// 3. 成功 → 手动 ACK// 注意RocketMQListener 正常返回即视为 ACK}catch(Exceptione){// 抛异常 → 不 ACK → Broker 重试最多 3 次// 3 次仍失败 → 进入死信队列thrownewRuntimeException(消费失败触发重试,e);}}}重试配置# 消费重试策略 maxReconsumeTimes3 # 最大重试 3 次 # 重试间隔: 10s, 30s, 1min, 2min... (递增)3 次重试失败后消息进入死信队列%DLQ%order-create-group需人工介入处理。3.2 终极兜底本地消息表 定时扫描上述三条防线覆盖了绝大多数场景但在极端情况下如 Broker 集群整体故障、磁盘损坏消息仍可能丢失。这个概率极低但后果极严重的风险需要业务层面兜底。本地消息表方案┌────────────────────────────────────────────────────────────┐ │ 生产者服务 │ │ │ │ 1. 开启本地事务 │ │ 2. INSERT INTO local_msg (msg_id, payload, status待发送)│ │ 3. 发送 MQ 消息 │ │ 4. 如果 MQ 发送成功 → 提交事务 │ │ 如果 MQ 发送失败 → 不提交事务等待定时任务重试 │ └────────────────────────────────────────────────────────────┘ ┌────────────────────────────────────────────────────────────┐ │ 定时任务 │ │ │ │ 每 5 分钟执行: │ │ 1. 查询 local_msg 表中 status待发送 且创建时间5分钟的记录 │ │ 2. 重新发送 MQ 消息 │ │ 3. 更新 status 为 已发送 │ │ 4. 重试次数 最大值的记录 → 人工介入 │ └────────────────────────────────────────────────────────────┘ ┌────────────────────────────────────────────────────────────┐ │ 消费者服务 │ │ │ │ 1. 消费 MQ 消息 │ │ 2. 执行业务订单入库 │ │ 3. 更新 local_msg 表中对应消息的 status 为 已消费 │ │ 4. 失败 → 不更新状态定时任务补偿 │ └────────────────────────────────────────────────────────────┘这个方案的本质是“用数据库的持久化能力弥补 MQ 的不可靠窗口”。本地消息表存储在 MySQL 中只要 MySQL 事务提交成功消息就不会丢失——定时任务会不断重试直到成功。3.3 幂等性一锁二判三更新消息重试机制保证了不丢失但也带来了新问题重复消费。同一个下单消息可能被消费多次生产者重试 消费者重试 网络超时导致 ACK 丢失如果下游不处理幂等就会产生重复订单。幂等性方案的核心框架——一锁二判三更新一锁: 获取分布式锁以消息唯一标识为 Key 二判: 查消息处理表/业务表判断是否已处理 三更新: 执行业务 写入处理记录同一事务具体实现TransactionalpublicvoidcreateOrder(SeckillOrderorder){// 一锁以订单号为 Key 加分布式锁StringlockKeylock:order:order.getOrderId();RLocklockredissonClient.getLock(lockKey);try{lock.lock(5,TimeUnit.SECONDS);// 二判查业务表是否已存在if(orderMapper.existsByOrderId(order.getOrderId())){log.warn(重复消费已跳过: orderId{},order.getOrderId());return;}// 三更新入库 写处理记录天然在同一事务中orderMapper.insert(order);// 订单表中 order_id 有唯一索引重复插入会触发 DuplicateKeyException// 这是数据库层面的第二道幂等保障}finally{lock.unlock();}}双重幂等保障保障层机制作用应用层Redis 分布式锁 查表判断拦截 99% 的重复消费数据库层订单号唯一索引uk_order_id兜底即使锁和判断都失败DB 也会拒绝重复插入订单号生成使用雪花算法Snowflake在秒杀服务本地生成不需要 RPC 调用且天然全局唯一// 雪花算法: 1bit 符号 41bit 时间戳 10bit 工作机器 12bit 序列号longorderIdsnowflakeIdGenerator.nextId();// 示例: 71234567890123456783.4 垃圾消息重试次数限制当消费者代码存在 Bug 导致某条消息永远无法成功消费时如果没有重试上限这条消息会在重试队列中无限循环——这就是垃圾消息。它不仅浪费系统资源还可能阻碍正常消息的消费。解决消费重试最大 3 次后转入死信队列同时本地消息表的重试次数也设上限// 定时任务扫描本地消息表时ListLocalMessagemsgslocalMsgMapper.selectRetryable(now);for(LocalMessagemsg:msgs){if(msg.getRetryCount()MAX_RETRY){// 如 10 次localMsgMapper.updateStatus(msg.getId(),FAILED);alertService.send(消息重试达到上限: msgIdmsg.getId());continue;}// 重试发送sendMQMessage(msg);localMsgMapper.incrementRetryCount(msg.getId());}3.5 消息积压监控与快速恢复积压发现通过 RocketMQ Web Console 或 Prometheus 监控消息堆积量 生产总量 - 消费总量。设置告警阈值积压 10000 条告警通知积压 50000 条紧急告警 自动扩容消费者积压处理紧急扩容消费者临时增加消费者实例数通过 K8s HPA 或手动扩容提升消费并发度调整consumeThreadMin和consumeThreadMax跳过非核心逻辑临时降级消费者只做订单入库跳过积分更新、通知发送等非核心操作事后补数积压消除后通过定时任务补做降级阶段跳过的操作# 消费者线程配置 consumeThreadMin20 consumeThreadMax64四、边界与异常4.1 本地消息表的性能开销每次发送消息前都要写一次本地消息表会增加约 5-10ms 的延迟。对于秒杀下单消息QPS 5000这个开销可以接受。但对于更高频的操作如更新缓存本地消息表方案过重应仅用于发送失败后果严重的消息场景。4.2 死信队列的消息如何处理死信队列DLQ中的消息需要人工或脚本处理查看失败原因日志修复根本问题如补上缺失的数据通过 RocketMQ Console 重投死信消息或编写批量重投脚本4.3 幂等锁的过期时间Redis 分布式锁的过期时间要大于业务执行时间。如果业务因 DB 慢查询等原因执行了 30 秒而锁在 10 秒后过期则重复消费可能绕过幂等判断。设置保守的锁过期时间如 60 秒依靠finally { unlock() }主动释放。4.4 事务消息的回查间隔RocketMQ 默认回查间隔为 60 秒这意味着如果生产者 COMMIT 时网络超时消费者最多要等 60 秒才能消费到这条消息。对于秒杀场景可以调整回查间隔为 10-15 秒producer.setCheckImmunityTimeInSeconds(10);五、总结消息 0 丢失需要三层保障事务消息发送环节 同步刷盘存储环节 手动 ACK消费环节缺一不可。本地消息表 定时扫描是终极兜底用 DB 的持久化能力弥补 MQ 的理论极限适用于发送失败后果严重的核心链路。幂等性的核心是一锁二判三更新分布式锁 查表判断 业务写入同一事务数据库唯一索引作为最后防线。垃圾消息靠重试上限控制设置最大重试次数消费 3 次 本地表 10 次超限转入死信队列并人工处理。消息积压需要监控 自动扩容积压 10000 告警 50000 紧急扩容消费者支持临时降级非核心消费逻辑。