首页 > 数据库 >Redis+MQ高并发秒杀技术方案与实现

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

来源:互联网 2026-07-24 09:10:26

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

前言

电商秒杀场景下,瞬间涌入的海量请求往往达到数万甚至数十万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的最终一致性,避免数据断层;
  • 流水对账作为最后一道防线,及时发现并修复异常。

侠游戏发布此文仅为了传递信息,不代表侠游戏网站认同其观点或证实其描述

热游推荐

更多
湘ICP备2026025700号-3 湘公网安备 43070302000280号
All Rights Reserved
本站为非盈利网站,不接受任何广告。本站所有软件,都由网友
上传,如有侵犯你的版权,请发邮件给xiayx666@163.com
抵制不良色情、反动、暴力游戏。注意自我保护,谨防受骗上当。
适度游戏益脑,沉迷游戏伤身。合理安排时间,享受健康生活。