游乐游手机版
首页/数据库/文章详情

Redis+MQ高并发秒杀技术方案与实现详解

时间:2026-07-29 06:21
基于Redis缓存与RocketMQ消息队列,实现高并发秒杀系统。前端使用令牌防止重复请求,Lua脚本原子性预扣减库存并记录流水,半消息经本地事务校验后提交,消费者消费消息后执行数据库扣减操作,确保最终一致性。

前言

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

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 tokenMap = redisTemplate.opsForHash().entries(uid);        for (Map.Entry entry : tokenMap.entrySet()) {            String token = (String) entry.getKey();            String flowJson = (String) entry.getValue();            SeckillFlow flow = JSON.parseObject(flowJson, SeckillFlow.class);             // 2. 检查MySQL是否有对应订单            Integer count = jdbcTemplate.queryForObject(                "SELECT COUNT(1) FROM orders WHERE product_id = ? AND uid = ? AND token = ?",                Integer.class,                flow.getProduct(), flow.getUid(), token            );             if (count == 0) {                // 3. 未找到订单 → 人工介入或自动回滚Redis库存                log.warn("发现不一致:Redis有流水但MySQL无订单,product={}, uid={}", flow.getProduct(), uid);                // redisTemplate.opsForValue().increment(flow.getProduct(), Integer.parseInt(flow.getChange()));            }        }    }}

这个定时任务对比Redis流水和订单表的数据,如果发现Redis有流水但MySQL没有对应订单,说明订单生成失败,就需要人工介入补单或回滚Redis库存,避免少卖;反过来,如果订单表有记录但MySQL库存未扣减,则要触发库存补扣,防止多卖。

总而言之,整个方案通过预扣减 + 事务消息 + 对账三重机制,为秒杀系统上了三道保险。Redis扛住高并发,事务消息保证一致性,对账再兜底,这才是应对高并发秒杀的成熟打法。

总结

Redis+MQ方案通过预扣减 + 事务消息 + 对账三重机制,完美解决了高并发秒杀的核心痛点:

  • Redis承担高并发读写,通过Lua脚本确保原子性,防止超卖;
  • MQ事务消息保障Redis与MySQL的最终一致性,避免数据断层;
  • 流水对账作为最后一道防线,及时发现并修复异常。
来源:https://www.jb51.net/database/360213mig.htm
上一篇Redis BGSAVE内存不足异常解决方案 下一篇Redis缓存雪崩:原理、防御策略与工程实践指南
本站内容用于信息整理与展示,如有侵权或内容问题请及时联系处理。

相关推荐

补充同频道和同主题内容,方便继续浏览更多相关内容。

同类最新

继续查看同栏目最近更新的文章。

更多
Redis是什么:核心特性、架构与应用场景解析
数据库 · 2026-09-01

Redis是什么:核心特性、架构与应用场景解析

Redis是一款基于内存的键值型NoSQL数据库,以超高读写速度和丰富的数据结构著称。本文系统梳理Redis的核心特性、架构组成、性能优势及典型应用场景,并通过与Memcached、MySQL、MongoDB的对比,帮助开发者快速判断Redis是否适合当前业务需求。

Windows 安装 MongoDB 完整图文教程
数据库 · 2026-09-01

Windows 安装 MongoDB 完整图文教程

本文详细介绍在 Windows 系统上安装 MongoDB 的完整流程。从官网下载 MSI 安装包开始,逐步演示自定义安装路径、配置 Windows 服务、跳过 MongoDB Compass 等关键选项,并提供通过系统服务列表验证安装是否成功的方法,帮助开发者快速搭建本地 MongoDB 环境。

Linux 安装 MongoDB 完整指南:依赖配置、环境变量与服务启动
数据库 · 2026-09-01

Linux 安装 MongoDB 完整指南:依赖配置、环境变量与服务启动

本文详解在 Linux 系统下安装 MongoDB 的完整流程,涵盖依赖包安装、二进制包下载解压、环境变量配置、数据与日志目录创建及服务启动验证。通过标准化命令与路径说明,帮助开发者快速完成部署并确认服务状态。

MacOS安装MongoDB完整教程
数据库 · 2026-09-01

MacOS安装MongoDB完整教程

本文介绍在MacOS系统下安装MongoDB的完整流程,涵盖下载、解压、目录配置、环境变量设置及服务启动。通过明确的命令与参数说明,帮助开发者快速完成环境搭建并验证安装结果。

Ubuntu系统安装与配置Redis完整指南
数据库 · 2026-09-01

Ubuntu系统安装与配置Redis完整指南

本文详解在Ubuntu系统中安装Redis的两种主流方式:apt在线安装与源码编译安装。涵盖版本选择逻辑、服务启停与状态检查、连接验证方法,以及在线练习工具与桌面GUI客户端的对比与使用建议,帮助开发者快速搭建并验证Redis运行环境。