ARTICLE DETAIL

资讯详情

深耕郑州网站建设与运营推广的一线实战洞察。

【黑马点评 | 第九篇】秒杀优化-异步实现

【黑马点评 | 第九篇】秒杀优化-异步实现 秒杀下单的核心矛盾是请求量很大但数据库扣库存和创建订单的处理能力有限。如果每个请求都同步完成查询、校验、扣库存和写订单数据库很快就会成为瓶颈。这次优化把秒杀流程拆成两段Redis 负责快速判断库存和一人一单数据库相关的订单处理放到后台线程异步执行。实现先从 JVM 阻塞队列开始再解决阻塞队列的内存和可靠性问题最后把订单消息迁移到 Redis Stream 消费组。一、同步下单为什么需要优化原有下单流程通常包含以下步骤查询优惠券和秒杀信息判断秒杀是否开始、是否结束以及库存是否充足查询当前用户是否已经下单扣减数据库库存创建订单。其中优惠券查询、重复下单查询、库存更新和订单保存都会访问数据库。高并发到来时大量请求会同时占用数据库连接和执行 SQL真正需要写数据库的步骤会互相争抢资源。优化后的处理顺序是请求先把库存和一人一单判断放到 Redis 中完成。判断成功后只把订单所需的用户 ID、优惠券 ID 和订单 ID 交给异步线程数据库扣库存和创建订单由后台慢慢处理。这样请求线程不需要等待完整的下单事务结束可以更快返回订单 ID。二、用 Redis Lua 完成秒杀资格判断2.1 把秒杀库存预热到 Redis创建秒杀优惠券时数据库保存优惠券和秒杀信息同时把库存写入 Redis。库存 key 按优惠券 ID 区分例如OverrideTransactionalpublicvoidaddSeckillVoucher(Vouchervoucher){// 保存优惠券save(voucher);// 保存秒杀信息SeckillVoucherseckillVouchernewSeckillVoucher();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_KEYvoucher.getId(),voucher.getStock().toString());}请求到达时直接读取 Redis 中的库存不需要先查询数据库。Redis 中还需要为每张优惠券维护一个 Set用来记录已经抢购成功的用户 ID。2.2 为什么库存判断和重复下单判断要放进 Lua库存判断、重复下单判断、扣减库存和记录用户必须作为一个整体执行。如果先读取库存再判断用户是否下单最后分别执行扣库存和写入用户多个请求之间可能在这些步骤中交错执行造成超卖或重复下单。Redis 执行 Lua 脚本时会把脚本中的命令作为一个原子操作处理因此可以把这几步放在同一个脚本中-- 1.参数列表-- 1.1.优惠券idlocalvoucherIdARGV[1]-- 1.2.用户idlocaluserIdARGV[2]-- 2.数据key-- 2.1.库存keylocalstockKeyseckill:stock:..voucherId-- 2.2.订单keylocalorderKeyseckill: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 文件加载成 DefaultRedisScriptprivatestaticfinalDefaultRedisScriptLongSECKILL_SCRIPT;static{SECKILL_SCRIPTnewDefaultRedisScript();SECKILL_SCRIPT.setLocation(newClassPathResource(seckill.lua));SECKILL_SCRIPT.setResultType(Long.class);}请求方法只负责获取用户 ID、生成订单 ID、执行脚本并根据返回值决定是否继续OverridepublicResultseckillVoucher(LongvoucherId){//获取用户LonguserIdUserHolder.getUser().getId();//获取订单idlongorderIdredisIdWork.nextId(order);//1.调用lua脚本LongresultstringRedisTemplate.execute(SECKILL_SCRIPT,Collections.emptyList(),voucherId.toString(),userId.toString(),String.valueOf(orderId));intrresult.intValue();//2.判断结果是否为0if(r!0){//2.1不为0,无购买资格returnResult.fail(r1?库存不足:不能重复下单);}//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_EXECUTORExecutors.newSingleThreadExecutor();privatefinalBlockingQueueVoucherOrderorderTasksnewArrayBlockingQueue(1024*1024);PostConstructprivatevoidinit(){SECKILL_ORDER_EXECUTOR.submit(newVoucherOrderHandler());}privateclassVoucherOrderHandlerimplementsRunnable{Overridepublicvoidrun(){while(true){try{//1.获取订单中的队列消息VoucherOrdervoucherOrderorderTasks.take();//2.创建订单handleVoucherOrder(voucherOrder);}catch(Exceptione){log.error(处理订单异常:,e);}}}}ArrayBlockingQueue是有界阻塞队列。生产者把订单放入队列消费者线程通过 take 阻塞等待消息队列为空时消费者不会空转队列有消息时才继续处理。单线程执行器保证当前实例内按顺序处理订单。3.2 消费线程中的一人一单和数据库操作订单处理线程从队列拿到消息后再获取用户维度的分布式锁查询数据库订单、扣减库存并保存订单privatevoidhandleVoucherOrder(VoucherOrdervoucherOrder){// 1.获取用户LonguserIdvoucherOrder.getUserId();// 2.创建锁对象RLockredisLockredissonClient.getLock(lock:order:userId);// 3.尝试获取锁booleanisLockredisLock.tryLock();// 4.判断是否获得锁成功if(!isLock){// 获取锁失败直接返回失败或者重试log.error(不允许重复下单);return;}try{proxy.createVoucherOrder(voucherOrder);}finally{// 释放锁redisLock.unlock();}}TransactionalpublicvoidcreateVoucherOrder(VoucherOrdervoucherOrder){LonguserIdvoucherOrder.getUserId();// 5.1.查询订单intcountquery().eq(user_id,userId).eq(voucher_id,voucherOrder.getVoucherId()).count();// 5.2.判断是否存在if(count0){// 用户已经购买过了log.error(用户已经购买过一次);return;}// 6.扣减库存booleansuccessseckillVoucherService.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 把一条带字段和值的消息追加到 StreamRedis 为消息生成递增的消息 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.优惠券idlocalvoucherIdARGV[1]-- 1.2.用户idlocaluserIdARGV[2]-- 1.3.订单idlocalorderIdARGV[3]-- 2.数据key-- 2.1.库存keylocalstockKeyseckill:stock:..voucherId-- 2.2.订单keylocalorderKeyseckill: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 StreamJava 请求线程只需要返回订单 ID不再把订单对象放进 JVM 队列。5.3 消费者组读取新消息项目启动后创建单线程执行器并提交订单消费者任务。消费者使用 g1 消费组和 c1 消费者从最后消费位置读取新消息没有消息时最多阻塞两秒privatestaticfinalExecutorServiceSECKILL_ORDER_EXECUTORExecutors.newSingleThreadExecutor();PostConstructprivatevoidinit(){SECKILL_ORDER_EXECUTOR.submit(newVoucherOrderHandler());}privateclassVoucherOrderHandlerimplementsRunnable{privatefinalStringqueueNamestream.orders;Overridepublicvoidrun(){while(true){try{// 1.获取消息队列中的订单信息ListMapRecordString,Object,ObjectliststringRedisTemplate.opsForStream().read(Consumer.from(g1,c1),StreamReadOptions.empty().count(1).block(Duration.ofSeconds(2)),StreamOffset.create(queueName,ReadOffset.lastConsumed()));//2.判断消息获取是否成功if(listnull||list.isEmpty()){// 如果为null说明没有消息继续下一次循环continue;}// 3.解析消息中的订单信息MapRecordString,Object,Objectrecordlist.get(0);MapObject,Objectvaluesrecord.getValue();VoucherOrdervoucherOrderBeanUtil.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中的订单信息ListMapRecordString,Object,ObjectliststringRedisTemplate.opsForStream().read(Consumer.from(g1,c1),StreamReadOptions.empty().count(1),StreamOffset.create(queueName,ReadOffset.from(0)));// 2.判断消息获取是否成功if(listnull||list.isEmpty()){// 如果获取失败说明pending-list没有异常消息结束循环break;}// 3.解析消息中的订单信息MapRecordString,Object,Objectrecordlist.get(0);MapObject,Objectvaluesrecord.getValue();VoucherOrdervoucherOrderBeanUtil.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 正好对应这条异步下单链路。
返回列表