秒杀下单的核心矛盾是:请求量很大,但数据库扣库存和创建订单的处理能力有限。如果每个请求都同步完成查询、校验、扣库存和写订单,数据库很快就会成为瓶颈。
这次优化把秒杀流程拆成两段:Redis 负责快速判断库存和一人一单,数据库相关的订单处理放到后台线程异步执行。实现先从 JVM 阻塞队列开始,再解决阻塞队列的内存和可靠性问题,最后把订单消息迁移到 Redis Stream 消费组。
一、同步下单为什么需要优化
原有下单流程通常包含以下步骤:
- 查询优惠券和秒杀信息;
- 判断秒杀是否开始、是否结束以及库存是否充足;
- 查询当前用户是否已经下单;
- 扣减数据库库存;
- 创建订单。
其中,优惠券查询、重复下单查询、库存更新和订单保存都会访问数据库。高并发到来时,大量请求会同时占用数据库连接和执行 SQL,真正需要写数据库的步骤会互相争抢资源。
优化后的处理顺序是:请求先把库存和一人一单判断放到 Redis 中完成。判断成功后,只把订单所需的用户 ID、优惠券 ID 和订单 ID 交给异步线程,数据库扣库存和创建订单由后台慢慢处理。这样,请求线程不需要等待完整的下单事务结束,可以更快返回订单 ID。
二、用 Redis Lua 完成秒杀资格判断
2.1 把秒杀库存预热到 Redis
创建秒杀优惠券时,数据库保存优惠券和秒杀信息,同时把库存写入 Redis。库存 key 按优惠券 ID 区分,例如:
@Override@TransactionalpublicvoidaddSeckillVoucher(Vouchervoucher){// 保存优惠券save(voucher);// 保存秒杀信息SeckillVoucherseckillVoucher=newSeckillVoucher();seckillVoucher.setVoucherId(voucher.getId());seckillVoucher.setStock(voucher.getStock());seckillVoucher.setBeginTime(voucher.getBeginTime());seckillVoucher.setEndTime(voucher.getEndTime());seckillVoucherService.save(seckillVoucher);// 保存秒杀库存到Redis中stringRedisTemplate.opsForValue().set(RedisConstants.SECKILL_STOCK_KEY+voucher.getId(),voucher.getStock().toString());}请求到达时直接读取 Redis 中的库存,不需要先查询数据库。Redis 中还需要为每张优惠券维护一个 Set,用来记录已经抢购成功的用户 ID。
2.2 为什么库存判断和重复下单判断要放进 Lua
库存判断、重复下单判断、扣减库存和记录用户必须作为一个整体执行。如果先读取库存,再判断用户是否下单,最后分别执行扣库存和写入用户,多个请求之间可能在这些步骤中交错执行,造成超卖或重复下单。
Redis 执行 Lua 脚本时会把脚本中的命令作为一个原子操作处理,因此可以把这几步放在同一个脚本中:
-- 1.参数列表-- 1.1.优惠券idlocalvoucherId=ARGV[1]-- 1.2.用户idlocaluserId=ARGV[2]-- 2.数据key-- 2.1.库存keylocalstockKey='seckill:stock:'..voucherId-- 2.2.订单keylocalorderKey='seckill:order:'..voucherId-- 3.脚本业务-- 3.1.判断库存是否充足 get stockKeyif(tonumber(redis.call('get',stockKey))<=0)then-- 3.2.库存不足,返回1return1end-- 3.2.判断用户是否下单 SISMEMBER orderKey userIdif(redis.call('sismember',orderKey,userId)==1)then-- 3.3.存在,说明是重复下单,返回2return2end-- 3.4.扣库存 incrby stockKey -1redis.call('incrby',stockKey,-1)-- 3.5.下单(保存用户)sadd orderKey userIdredis.call('sadd',orderKey,userId)return0脚本返回 1 表示库存不足,返回 2 表示用户已经购买过,返回 0 表示抢购资格校验通过。这里的 Set 只记录“已经通过 Redis 抢购资格判断”的用户,数据库订单仍然由后续异步流程创建。
2.3 Java 调用 Lua 脚本
Spring Data Redis 可以把 resources 目录下的 Lua 文件加载成 DefaultRedisScript:
privatestaticfinalDefaultRedisScript<Long>SECKILL_SCRIPT;static{SECKILL_SCRIPT=newDefaultRedisScript<>();SECKILL_SCRIPT.setLocation(newClassPathResource("seckill.lua"));SECKILL_SCRIPT.setResultType(Long.class);}请求方法只负责获取用户 ID、生成订单 ID、执行脚本并根据返回值决定是否继续:
@OverridepublicResultseckillVoucher(LongvoucherId){//获取用户LonguserId=UserHolder.getUser().getId();//获取订单idlongorderId=redisIdWork.nextId("order");//1.调用lua脚本Longresult=stringRedisTemplate.execute(SECKILL_SCRIPT,Collections.emptyList(),voucherId.toString(),userId.toString(),String.valueOf(orderId));intr=result.intValue();//2.判断结果是否为0if(r!=0){//2.1不为0,无购买资格returnResult.fail(r==1?"库存不足":"不能重复下单");}//2.2 为0,有购买资格,订单消息已经写入消息队列returnResult.ok(orderId);}当前 Java 调用传入了三个参数:优惠券 ID、用户 ID 和订单 ID。Lua 脚本如果要把订单消息写入 Stream,还必须定义local orderId = ARGV[3];否则后面的 XADD 命令拿不到订单 ID。
三、第一版实现:阻塞队列加异步线程池
3.1 先把数据库下单放到后台线程
Lua 脚本返回 0 后,说明请求已经完成了库存预扣和一人一单判断。此时可以创建 VoucherOrder 对象,放入 JVM 的 BlockingQueue,然后立即返回订单 ID。
privatestaticfinalExecutorServiceSECKILL_ORDER_EXECUTOR=Executors.newSingleThreadExecutor();privatefinalBlockingQueue<VoucherOrder>orderTasks=newArrayBlockingQueue<>(1024*1024);@PostConstructprivatevoidinit(){SECKILL_ORDER_EXECUTOR.submit(newVoucherOrderHandler());}privateclassVoucherOrderHandlerimplementsRunnable{@Overridepublicvoidrun(){while(true){try{//1.获取订单中的队列消息VoucherOrdervoucherOrder=orderTasks.take();//2.创建订单handleVoucherOrder(voucherOrder);}catch(Exceptione){log.error("处理订单异常:",e);}}}}ArrayBlockingQueue是有界阻塞队列。生产者把订单放入队列,消费者线程通过 take 阻塞等待消息;队列为空时,消费者不会空转,队列有消息时才继续处理。单线程执行器保证当前实例内按顺序处理订单。
3.2 消费线程中的一人一单和数据库操作
订单处理线程从队列拿到消息后,再获取用户维度的分布式锁,查询数据库订单、扣减库存并保存订单:
privatevoidhandleVoucherOrder(VoucherOrdervoucherOrder){// 1.获取用户LonguserId=voucherOrder.getUserId();// 2.创建锁对象RLockredisLock=redissonClient.getLock("lock:order:"+userId);// 3.尝试获取锁booleanisLock=redisLock.tryLock();// 4.判断是否获得锁成功if(!isLock){// 获取锁失败,直接返回失败或者重试log.error("不允许重复下单!");return;}try{proxy.createVoucherOrder(voucherOrder);}finally{// 释放锁redisLock.unlock();}}@TransactionalpublicvoidcreateVoucherOrder(VoucherOrdervoucherOrder){LonguserId=voucherOrder.getUserId();// 5.1.查询订单intcount=query().eq("user_id",userId).eq("voucher_id",voucherOrder.getVoucherId()).count();// 5.2.判断是否存在if(count>0){// 用户已经购买过了log.error("用户已经购买过一次!");return;}// 6.扣减库存booleansuccess=seckillVoucherService.update().setSql("stock = stock - 1")// set stock = stock - 1.eq("voucher_id",voucherOrder.getVoucherId()).gt("stock",0)// where id = ? and stock > 0.update();if(!success){log.error("库存不足!");return;}// 7.创建订单save(voucherOrder);}异步线程不会自动继承请求线程中的 Spring 事务上下文,事务方法需要通过 Spring 代理调用。直接使用 this 调用会绕过事务代理,导致 @Transactional 不生效,因此阻塞队列版本需要保存订单服务代理对象,再由消费者线程调用代理方法。
3.3 为什么异步后响应更快
同步方案要等数据库查询、扣库存和保存订单都完成后才能返回。异步方案把 Redis 中可以快速完成的资格判断放在请求线程,数据库写操作交给后台线程。请求线程只负责生成订单 ID、投递订单消息和返回结果,等待时间明显缩短。
这里的“异步”只改变调用方是否等待,不代表数据库操作消失了。订单最终仍然要经过库存更新和订单保存,后台消费者的处理能力决定了消息堆积速度和订单落库速度。
四、阻塞队列实现存在的问题
4.1 JVM 内存限制
BlockingQueue 的数据保存在当前 Java 进程的堆内存中。队列容量虽然可以设置得很大,但仍然受到 JVM 内存上限约束。秒杀流量持续升高时,订单消息可能不断堆积,最终导致内存压力甚至触发频繁 GC 或内存溢出。
此外,单个实例中的队列只能被这个实例自己的消费者线程读取。应用扩容后,每个实例都有一份独立队列,订单消息不会自动在多个实例之间共享,实例之间也无法统一协调积压情况。
4.2 数据安全问题
队列中的消息还没有落到 Redis 或数据库时,数据只存在 JVM 内存中。应用重启、服务器宕机或进程异常退出,队列里的订单消息都会丢失。
消费者通过 take 取出消息后,如果数据库处理过程中发生异常,消息已经从队列中移除,BlockingQueue 没有消息确认和待处理列表,程序也无法根据消息 ID 自动恢复。因此,阻塞队列适合演示异步流程或低风险的临时任务,不适合承载需要可靠投递的订单消息。
要解决这两个问题,订单消息需要放到独立于应用 JVM 的持久化存储中,并且要能记录“消息已经被哪个消费者取走但还没有处理完成”。Redis Stream 的消费者组提供了这两种能力。
五、使用 Redis Stream 改造消息队列
5.1 Stream 的消息模型
Redis Stream 是 Redis 提供的日志型数据结构。生产者通过 XADD 把一条带字段和值的消息追加到 Stream,Redis 为消息生成递增的消息 ID。消费者可以按 ID 读取历史消息,也可以使用阻塞读取等待新消息。
异步秒杀中需要一个名为 stream.orders 的 Stream,以及一个消费者组 g1。消费者组会记录组的消费位置,并为每个消费者维护 Pending Entries List,简称 pending-list:
- 消费者读取到消息后,消息先进入当前消费者的 pending-list;
- 业务处理成功后,通过 XACK 确认消息,消息才会从 pending-list 移除;
- 消费者处理过程中宕机时,消息仍然保留在 pending-list,可以重新读取处理。
创建消息队列和消费者组的命令如下:
XGROUP CREATE stream.orders g1 0 MKSTREAM其中,0 表示从 Stream 的第一条消息开始消费,MKSTREAM 表示 Stream 不存在时自动创建。生产环境使用时,消费者组创建操作应保证只执行一次;重复创建同名消费者组会返回已存在错误,需要按项目启动流程处理。
5.2 让 Lua 在资格通过后直接发送订单消息
Redis 资格判断和消息投递需要保持一致。Lua 脚本只有在库存充足、用户未下单、库存扣减和用户记录都成功后,才向 stream.orders 写入订单消息:
-- 1.参数列表-- 1.1.优惠券idlocalvoucherId=ARGV[1]-- 1.2.用户idlocaluserId=ARGV[2]-- 1.3.订单idlocalorderId=ARGV[3]-- 2.数据key-- 2.1.库存keylocalstockKey='seckill:stock:'..voucherId-- 2.2.订单keylocalorderKey='seckill:order:'..voucherId-- 3.脚本业务-- 3.1.判断库存是否充足 get stockKeyif(tonumber(redis.call('get',stockKey))<=0)then-- 3.2.库存不足,返回1return1end-- 3.2.判断用户是否下单 SISMEMBER orderKey userIdif(redis.call('sismember',orderKey,userId)==1)then-- 3.3.存在,说明是重复下单,返回2return2end-- 3.4.扣库存 incrby stockKey -1redis.call('incrby',stockKey,-1)-- 3.5.下单(保存用户)sadd orderKey userIdredis.call('sadd',orderKey,userId)-- 3.6.发送消息到队列中, XADD stream.orders * k1 v1 k2 v2 ...redis.call('xadd','stream.orders','*','userId',userId,'voucherId',voucherId,'id',orderId)return0XADD 与前面的库存扣减、用户记录处于同一个 Lua 脚本中。脚本返回 0 后,订单消息已经进入 Redis Stream,Java 请求线程只需要返回订单 ID,不再把订单对象放进 JVM 队列。
5.3 消费者组读取新消息
项目启动后创建单线程执行器,并提交订单消费者任务。消费者使用 g1 消费组和 c1 消费者,从最后消费位置读取新消息;没有消息时最多阻塞两秒:
privatestaticfinalExecutorServiceSECKILL_ORDER_EXECUTOR=Executors.newSingleThreadExecutor();@PostConstructprivatevoidinit(){SECKILL_ORDER_EXECUTOR.submit(newVoucherOrderHandler());}privateclassVoucherOrderHandlerimplementsRunnable{privatefinalStringqueueName="stream.orders";@Overridepublicvoidrun(){while(true){try{// 1.获取消息队列中的订单信息List<MapRecord<String,Object,Object>>list=stringRedisTemplate.opsForStream().read(Consumer.from("g1","c1"),StreamReadOptions.empty().count(1).block(Duration.ofSeconds(2)),StreamOffset.create(queueName,ReadOffset.lastConsumed()));//2.判断消息获取是否成功if(list==null||list.isEmpty()){// 如果为null,说明没有消息,继续下一次循环continue;}// 3.解析消息中的订单信息MapRecord<String,Object,Object>record=list.get(0);Map<Object,Object>values=record.getValue();VoucherOrdervoucherOrder=BeanUtil.fillBeanWithMap(values,newVoucherOrder(),true);// 4.如果获取成功,可以下单createVoucherOrder(voucherOrder);// 5.确认消息 XACK stream.orders g1 idstringRedisTemplate.opsForStream().acknowledge(queueName,"g1",record.getId());}catch(Exceptione){log.error("处理订单异常:",e);//处理异常消息handlePendingList();}}}}StreamReadOptions 的 count(1) 限制每次最多读取一条消息,block(Duration.ofSeconds(2)) 让没有消息时的读取进入阻塞等待。ReadOffset.lastConsumed() 对应消费者组的未消费位置,消息被读取后会进入当前消费者的 pending-list。
订单处理成功后才执行 acknowledge。XACK 的作用是从消费者组的 pending-list 中移除指定消息,表示这条订单消息已经完成处理。如果 createVoucherOrder 抛出异常,代码不会执行 XACK,而是进入 pending-list 恢复逻辑。
5.4 处理 pending-list 中未确认的消息
消费者进程可能在“读取消息”和“确认消息”之间异常退出。重新启动后,不能只读取新消息,还要先处理之前留在 pending-list 中的消息。项目通过相同的消费者组和消费者读取起始 ID 0,循环处理未确认消息:
privatevoidhandlePendingList(){while(true){try{// 1.获取pending-list中的订单信息List<MapRecord<String,Object,Object>>list=stringRedisTemplate.opsForStream().read(Consumer.from("g1","c1"),StreamReadOptions.empty().count(1),StreamOffset.create(queueName,ReadOffset.from("0")));// 2.判断消息获取是否成功if(list==null||list.isEmpty()){// 如果获取失败,说明pending-list没有异常消息,结束循环break;}// 3.解析消息中的订单信息MapRecord<String,Object,Object>record=list.get(0);Map<Object,Object>values=record.getValue();VoucherOrdervoucherOrder=BeanUtil.fillBeanWithMap(values,newVoucherOrder(),true);// 4.如果获取成功,可以下单createVoucherOrder(voucherOrder);// 5.确认消息 XACK stream.orders g1 idstringRedisTemplate.opsForStream().acknowledge(queueName,"g1",record.getId());}catch(Exceptione){log.error("处理pending订单异常",e);try{Thread.sleep(20);}catch(InterruptedExceptionex){Thread.currentThread().interrupt();return;}}}}新消息读取使用最后消费位置,pending-list 恢复使用 ID 0,这两个读取入口不能混淆。前者负责接收尚未分配给消费者的新消息,后者负责重新处理已经分配但没有确认的消息。
5.5 Stream 版本的完整调用链
改造后的订单链路可以按下面的顺序理解:
- 请求线程生成订单 ID,并调用 Lua 脚本;
- Lua 检查 Redis 库存和用户下单集合;
- 校验通过后,Lua 原子扣减库存、记录用户并向 stream.orders 写入订单消息;
- 请求线程立即返回订单 ID;
- 后台消费者从 g1 消费组读取订单消息;
- 消费者把消息字段转换成 VoucherOrder,获取用户维度的 Redisson 锁;
- 查询数据库订单、扣减数据库库存、保存订单;
- 数据库处理成功后执行 XACK;
- 处理异常时保留 pending 状态,后续从 pending-list 继续消费。
这里 Redis 的库存预扣和数据库的最终扣库存各自承担不同职责:Redis 负责在高并发入口快速筛掉库存不足和重复请求,数据库负责最终落库。数据库更新仍然带有 stock > 0 条件,避免异步消费过程中把库存扣成负数。
总结
秒杀优化的第一步是把库存和一人一单判断放到 Redis Lua 中,以原子方式快速完成资格校验;第二步是把数据库下单从请求线程中移走,先用 BlockingQueue 验证异步处理流程;第三步是使用 Redis Stream 消费组替代 JVM 阻塞队列,让订单消息脱离单个应用实例的内存,并通过 pending-list 和 XACK 支持异常恢复。
阻塞队列适合说明异步模型,但受 JVM 内存限制,应用重启还会丢消息。Redis Stream 把消息存储、消费位置、待确认消息和恢复流程放在 Redis 中,当前项目的 stream.orders、g1 和 c1 正好对应这条异步下单链路。