前言电商秒杀场景下,瞬间涌入的海量请求往往达到数万甚至数十万QPS,对系统构成严峻考验。传统数据库单表架构难以支撑,而Redis与消息队列(MQ)的组合,凭借高性能与可靠性,成为应对高并发秒杀的黄金搭档。本文将深入拆解这套方案,从整体架构到具体实现,揭示其如何抵御秒杀冲击。方案总览整体流程概括为:用
电商秒杀场景下,瞬间涌入的海量请求往往达到数万甚至数十万QPS,对系统构成严峻考验。传统数据库单表架构难以支撑,而Redis与消息队列(MQ)的组合,凭借高性能与可靠性,成为应对高并发秒杀的黄金搭档。本文将深入拆解这套方案,从整体架构到具体实现,揭示其如何抵御秒杀冲击。

长期稳定更新的攒劲资源: >>>点此立即查看<<<
整体流程概括为:用户请求 → 前端生成Token → Redis执行Lua脚本(预扣减+防重+流水)→ 发送RocketMQ事务消息 → [本地事务校验Redis结果] → MQ消息确认(COMMIT/ROLLBACK)→ 消费者消费消息 → MySQL扣减库存+记录订单。
秒杀系统的核心诉求是抗并发、防超卖、保一致。Redis+MQ方案通过“前端拦截 - 中间缓冲 - 后端落地”三层架构实现目标:
预扣减流程是整个方案的起点,具体步骤如下:
开始
│
├─ 生成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,脚本执行后返回扣减后的库存或错误信息。
接下来是关键步骤:如何保证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方案通过预扣减 + 事务消息 + 对账三重机制,完美解决高并发秒杀的核心痛点:
侠游戏发布此文仅为了传递信息,不代表侠游戏网站认同其观点或证实其描述