前言
你可能遇到过电商秒杀场景中瞬间涌入的海量请求,数万甚至数十万QPS对系统而言堪称严峻考验。传统数据库单表架构根本无法支撑,而Redis与消息队列(MQ)的组合,凭借其卓越性能与高可靠性,已成为应对这类高并发场景的黄金搭档。本文将从整体架构到具体实现,详细拆解这套方案,看看它是如何扛住秒杀冲击的。

方案总览
整体流程可概括为:用户请求 → 前端生成Token → Redis执行Lua脚本(预扣减+防重+流水)→ 发送RocketMQ事务消息 → [本地事务校验Redis结果] → MQ消息确认(COMMIT/ROLLBACK)→ 消费者消费消息 → MySQL扣减库存+记录订单。
秒杀系统的核心目标在于:抗高并发、防止超卖、保证数据一致性。而Redis+MQ方案通过“前端拦截 - 中间缓冲 - 后端落地”的三层架构完美实现了这一目标:
- 前端拦截:Redis借助Lua脚本原子性地处理库存预扣减,过滤掉无效请求;
- 中间缓冲:MQ(如RocketMQ)通过事务消息削峰填谷,确保流量平稳进入数据库;
- 后端落地:MySQL最终存储库存与订单数据,通过事务消息保障与Redis的一致性。
流程拆解(示例代码)
Redis 库存预扣减
我们先来看预扣减的流程,这是整个方案的起点。
开始
│
├─ 生成Token(前端)
│
├─ 前端携带Token请求秒杀│
├─ 执行Lua脚本│ │
│ ├─ 检查Token是否存在(Hash结构)│ │ ├─ 存在 → 返回“重复提交”
│ │ └─ 不存在 → 继续│ │
│ ├─ 获取Redis库存(String结构)│ │ ├─ 库存不足 → 返回“库存不足”
│ │ └─ 库存充足 → 继续│ │
│ ├─ 扣减Redis库存并更新│ │
│ └─ 记录流水到Hash结构│
├─ 返回扣减结果(成功/失败)│
结束
关键的Lua脚本是整个逻辑的核心。这个脚本究竟完成了哪些操作?
-- 启用Redis命令复制,确保脚本在集群环境中正确同步redis.replicate_commands() -- 1. 防重提交校验:通过用户ID+Token判断是否重复提交-- KEYS[2]为用户ID(uid),ARGV[2]为本次请求的Tokenif redis.call('hexists', KEYS[2], ARGV[2]) == 1 then return redis.error_reply('repeat submit') -- 重复提交,返回错误end -- 2. 库存充足性校验local product_id = KEYS[1] -- 商品IDlocal stock = redis.call('get', KEYS[1]) -- 获取当前库存if not stock then -- 库存不存在(如商品未上架) return redis.error_reply('product not found')endif tonumber(stock) < tonumber(ARGV[1]) then -- 库存不足 return redis.error_reply('stock is not enough')end -- 3. 执行库存扣减local remaining_stock = tonumber(stock) - tonumber(ARGV[1])redis.call('set', KEYS[1], tostring(remaining_stock)) -- 更新库存 -- 4. 记录交易流水(用于后续一致性校验)local time = redis.call('time') -- 获取当前时间(秒+微秒)local currentTimeMillis = (time[1] * 1000) + math.floor(time[2] / 1000) -- 转换为毫秒时间戳-- 存储流水到Hash结构:用户ID → Token → 流水详情redis.call('hset', KEYS[2], ARGV[2], cjson.encode({ action = '扣减库存', product = product_id, from = stock, -- 扣减前库存 to = remaining_stock, -- 扣减后库存 change = ARGV[1], -- 扣减数量 token = ARGV[2], timestamp = currentTimeMillis })) return remaining_stock -- 返回扣减后库存脚本逻辑清晰明了:先检查重复提交,再判断库存是否充足,接着执行扣减并更新库存,最后将交易流水记录到Hash结构中。整个操作具备原子性,不会出现并发问题。
然后在Java中调用这个脚本,代码非常直接:
@Servicepublic class SeckillService { @Autowired private StringRedisTemplate redisTemplate; // 加载Lua脚本 private DefaultRedisScript stockScript; @PostConstruct public void init() { stockScript = new DefaultRedisScript<>(); stockScript.setScriptSource(new ResourceScriptSource(new ClassPathResource("seckill.lua"))); stockScript.setResultType(Long.class); } /** * 执行Redis库存预扣减 * @param productId 商品ID * @param uid 用户ID * @param quantity 购买数量 * @param token 防重Token * @return 扣减后库存(-1表示失败) */ public Long preDeductStock(String productId, String uid, Integer quantity, String token) { try { // 执行Lua脚本:KEYS = [商品ID, 用户ID],ARGV = [数量, Token] return redisTemplate.execute( stockScript, Arrays.asList(productId, uid), quantity.toString(), token ); } catch (Exception e) { log.error("Redis预扣减失败", e); return -1L; } }} 这里传入的参数依次为商品ID、用户ID、购买数量和防重Token,脚本执行后返回扣减后的库存或错误信息。
MySQL 库存扣减
接下来是关键步骤:如何保证Redis预扣减和MySQL最终扣减的一致性?这里使用了RocketMQ的事务消息。
扣减流程如下:
开始
│
├─ 发送半消息到RocketMQ│
├─ 执行本地事务│ │
│ ├─ 检查Redis流水是否存在│ │ ├─ 存在 → 提交消息(COMMIT)
│ │ └─ 不存在 → 回滚消息(ROLLBACK)│ │
│ └─ 未知状态 → 等待回查│
├─ RocketMQ回查机制│ ├─ 有流水 → 提交消息
│ └─ 无流水 → 回滚消息│
├─ 消息被消费│ │
│ ├─ 查询数据库当前版本号(乐观锁)│ │
│ ├─ 执行库存扣减(WHERE version = 当前版本)│ │ ├─ 扣减成功 → 记录数据库流水
│ │ └─ 扣减失败 → 抛出异常(触发重试)│
结束
系统首先向RocketMQ发送一条半消息,此时消息处于不可消费状态,需要等待确认。
// 发送半消息public void sendHalfMessage(String productId, String uid, String token, Integer quantity) { // 构建消息 Message message = new Message( "seckill_topic", // 主题 "stock_deduct", // 标签 JSON.toJSONString(new SeckillMessage(productId, uid, token, quantity)).getBytes() ); // 发送事务消息 TransactionSendResult result = rocketMQTemplate.sendMessageInTransaction( "seckill_producer_group", // 生产者组 message, null // 本地事务参数(可传递上下文) ); log.info("半消息发送结果:{}", result.getSendStatus());}半消息发送后,系统会执行本地事务来校验Redis预扣减是否成功。如果Redis中的Lua脚本执行成功(库存预扣减完成且流水已记录),系统就会向RocketMQ返回提交指令,消息变为可消费状态;如果失败(如库存不足或重复提交),则返回回滚指令,消息被丢弃。如果RocketMQ长时间未收到结果,就会触发回查机制,系统会再次检查Redis中是否存在对应流水来决定是提交还是回滚。
@Componentpublic class SeckillTransactionListener implements TransactionListener { @Autowired private StringRedisTemplate redisTemplate; // 执行本地事务 @Override public LocalTransactionState executeLocalTransaction(Message msg, Object arg) { try { SeckillMessage message = JSON.parseObject(new String(msg.getBody()), SeckillMessage.class); // 检查Redis中是否存在对应流水(验证预扣减成功) Boolean flag = redisTemplate.opsForHash().hasKey( message.getUid(), // Hash key:用户ID message.getToken() // Hash field:Token ); return flag ? RocketMQLocalTransactionState.COMMIT : RocketMQLocalTransactionState.ROLLBACK; } catch (Exception e) { return RocketMQLocalTransactionState.UNKNOWN; // 未知状态,触发回查 } } // 消息回查(解决超时未确认问题) @Override public LocalTransactionState checkLocalTransaction(MessageExt msg) { SeckillMessage message = JSON.parseObject(new String(msg.getBody()), SeckillMessage.class); // 回查逻辑:再次检查流水是否存在 Boolean flag = redisTemplate.opsForHash().hasKey(message.getUid(), message.getToken()); return flag ? RocketMQLocalTransactionState.COMMIT : RocketMQLocalTransactionState.ROLLBACK; }}消息被确认后,消费者开始处理。消费者的任务是从消息中获取信息,然后执行MySQL库存扣减操作。这里必须保证幂等性:如果消费失败,MQ会自动重试,直到成功或达到最大重试次数(此时需要人工介入)。
@Component@RocketMQMessageListener( topic = "seckill_topic", consumerGroup = "seckill_consumer_group", messageModel = MessageModel.CLUSTERING)public class SeckillConsumer implements RocketMQListener{ @Autowired private JdbcTemplate jdbcTemplate; @Override public void onMessage(MessageExt message) { SeckillMessage msg = JSON.parseObject(new String(message.getBody()), SeckillMessage.class); String productId = msg.getProductId(); int quantity = msg.getQuantity(); // 数据库扣减(使用乐观锁防超卖) String sql = "UPDATE product_stock " + "SET stock = stock - ?, version = version + 1 " + "WHERE product_id = ? AND stock >= ? AND version = ?"; // 1. 查询当前版本号 Integer version = jdbcTemplate.queryForObject( "SELECT version FROM product_stock WHERE product_id = ?", Integer.class, productId ); // 2. 执行扣减(乐观锁保证原子性) int rows = jdbcTemplate.update(sql, quantity, productId, quantity, version); if (rows > 0) { // 扣减成功:记录数据库流水 jdbcTemplate.update( "INSERT INTO stock_flow (product_id, quantity, op_type, create_time) " + "VALUES (?, ?, 'SECKILL', NOW())", productId, quantity ); // 确认消费成功(返回ACK) } else { // 扣减失败:触发重试(MQ默认重试机制) throw new RuntimeException("数据库扣减失败,触发重试"); } }}
这里采用了乐观锁,通过版本号保证并发操作的正确性。如果更新受影响的行数为0,说明库存已被其他请求扣减或版本号不匹配,需要抛出异常来触发重试。
最后一道防线是一致性保障。为防止Redis与MySQL数据不一致,系统通过定时任务定期对账:
@Scheduled(cron = "0 0 */1 * * ?") // 每小时执行一次public void reconcileStock() { // 1. 扫描Redis中未同步到MySQL的流水 Set uids = redisTemplate.keys("uid:*"); // 假设用户ID前缀为uid: for (String uid : uids) { Map 这个定时任务对比Redis流水和订单表的数据,如果发现Redis有流水但MySQL没有对应订单,说明订单生成失败,就需要人工介入补单或回滚Redis库存,避免少卖;反过来,如果订单表有记录但MySQL库存未扣减,则要触发库存补扣,防止多卖。
总而言之,整个方案通过预扣减 + 事务消息 + 对账三重机制,为秒杀系统上了三道保险。Redis扛住高并发,事务消息保证一致性,对账再兜底,这才是应对高并发秒杀的成熟打法。
总结
Redis+MQ方案通过预扣减 + 事务消息 + 对账三重机制,完美解决了高并发秒杀的核心痛点:
- Redis承担高并发读写,通过Lua脚本确保原子性,防止超卖;
- MQ事务消息保障Redis与MySQL的最终一致性,避免数据断层;
- 流水对账作为最后一道防线,及时发现并修复异常。
