大营销平台 —— 活动SKU库存扣减业务及其一致性处理

发布时间:2026/7/22 10:41:44
大营销平台 —— 活动SKU库存扣减业务及其一致性处理 一、前言前面我们搭建了活动订单业务的整体骨架设计了整体的活动订单流程我们在责任链中设计了两个节点一个用于校验一个用于扣减库存但是在上一节我们是没有写逻辑的所以这一节的第一件事就是去补齐逻辑。除此之外我们需要引入MQ来处理SKU的库存一致性这里的一致性不是指的前面抽奖那种延迟队列异步扣减库存而是当缓存中的库存为0时就不要再浪费性能去执行定时任务调度了直接用MQ传一个清空库存的消息直接清空延迟队列并且将数据库中的库存直接清0。二、SKU校验和SKU库存扣减这两个业务我们放在了上一节的两个空节点中去实现Slf4j Component(activity_base_action) public class ActivityBaseActionChain extends AbstractActionChain { Override public boolean action(ActivitySkuEntity activitySkuEntity, ActivityEntity activityEntity, ActivityCountEntity activityCountEntity) { log.info(活动责任链-基础信息【有效期、状态、库存sku】校验开始:sku:{} activityId:{}, activitySkuEntity.getSku(), activityEntity.getActivityId()); //校验活动状态 if (!ActivityStateVO.open.equals(activityEntity.getState())) { throw new AppException(ResponseCode.ACTIVITY_STATE_ERROR.getCode(), ResponseCode.ACTIVITY_STATE_ERROR.getInfo()); } //校验活动是否在有效期内 Date currentDate new Date(); if (activityEntity.getBeginDateTime().after(currentDate) || activityEntity.getEndDateTime().before(currentDate)) { throw new AppException(ResponseCode.ACTIVITY_DATE_ERROR.getCode(), ResponseCode.ACTIVITY_DATE_ERROR.getInfo()); } //校验活动sku库存 if (activitySkuEntity.getStockCountSurplus() 0) { throw new AppException(ResponseCode.ACTIVITY_SKU_STOCK_ERROR.getCode(), ResponseCode.ACTIVITY_SKU_STOCK_ERROR.getInfo()); } return next().action(activitySkuEntity, activityEntity, activityCountEntity); } }Slf4j Component(activity_sku_stock_action) public class ActivitySkuStockActionChain extends AbstractActionChain { Resource private IActivityDispatch activityDispatch; Resource private IActivityRepository activityRepository; Override public boolean action(ActivitySkuEntity activitySkuEntity, ActivityEntity activityEntity, ActivityCountEntity activityCountEntity) { log.info(活动责任链-商品库存处理【扣减】开始:sku:{} activityId:{}, activitySkuEntity.getSku(), activityEntity.getActivityId()); // 扣减缓存中的库存 boolean status activityDispatch.subtractionActivitySkuStock(activitySkuEntity.getSku(), activityEntity.getEndDateTime()); // true库存扣减成功 if (status) { log.info(活动责任链-商品库存处理【有效期、状态、库存(sku)】成功。sku:{} activityId:{}, activitySkuEntity.getSku(), activityEntity.getActivityId()); // 写入延迟队列延迟消费更新库存记录 activityRepository.activitySkuStockConsumeSendQueue(ActivitySkuStockKeyVO.builder() .sku(activitySkuEntity.getSku()) .activityId(activityEntity.getActivityId()) .build()); return true; } throw new AppException(ResponseCode.ACTIVITY_SKU_STOCK_ERROR.getCode(), ResponseCode.ACTIVITY_SKU_STOCK_ERROR.getInfo()); } }其中扣减数据库中SKU库存的动作放到延迟队列中去实现Override public void activitySkuStockConsumeSendQueue(ActivitySkuStockKeyVO activitySkuStockKeyVO) { String cacheKey Constants.RedisKey.ACTIVITY_SKU_COUNT_QUERY_KEY; RBlockingQueueActivitySkuStockKeyVO blockingQueue redisService.getBlockingQueue(cacheKey); RDelayedQueueActivitySkuStockKeyVO delayedQueue redisService.getDelayedQueue(blockingQueue); delayedQueue.offer(activitySkuStockKeyVO, 3, TimeUnit.SECONDS); }而上面扣减redis中库存的方法我们留到后面讲这里先只看整体流程。三、SKU库存预热这里的预热就是将数据库的库存数据同步到redis中去和抽奖策略领域中的奖品库存一样我们将每个扣减的数据都上锁避免运营时的失误详情请看第一阶段的解决超卖问题和库存回补后的重复扣减问题章节。这里还是通过实现两个接口来划分功能边界/** * author 印东升 * description 活动装配预热 * create 2026-07-20 17:49 */ public interface IActivityArmory { boolean assembleActivitySku(Long sku); }/** * author 印东升 * description 活动调度【扣减库存】 * create 2026-07-20 18:00 */ public interface IActivityDispatch { /** * * param sku 互动SKU * param endDateTime 活动结束时间根据结束时间设置加锁的key为结束时间 * return */ Boolean subtractionActivitySkuStock(Long sku, Date endDateTime); }/** * author 印东升 * description 活动装配预热 * create 2026-07-20 17:50 */ Slf4j Service public class ActivityArmory implements IActivityArmory,IActivityDispatch { Resource private IActivityRepository activityRepository; Override public boolean assembleActivitySku(Long sku) { ActivitySkuEntity activitySkuEntity activityRepository.queryActivitySku(sku); cacheActivitySkuStockCount(sku, activitySkuEntity.getStockCount()); //预热活动【查询时预热到缓存】 activityRepository.queryRaffleActivityByActivityId(activitySkuEntity.getActivityId()); //预热活动次数【查询时预热到缓存】 activityRepository.queryRaffleActivityCountByActivityCountId(activitySkuEntity.getActivityCountId()); return false; } private void cacheActivitySkuStockCount(Long sku, Integer stockCount) { String cacheKey Constants.RedisKey.ACTIVITY_SKU_STOCK_COUNT_KEY sku; activityRepository.cacheActivitySkuStockCount(cacheKey, stockCount); } Override public Boolean subtractionActivitySkuStock(Long sku, Date endDateTime) { String cacheKey Constants.RedisKey.ACTIVITY_SKU_STOCK_COUNT_KEY sku; return activityRepository.subtractionActivitySkuStock(sku,cacheKey,endDateTime); } }下面是仓储中的实现其实很简单就两行至于redisService的API详见第一阶段。Override public void cacheActivitySkuStockCount(String cacheKey, Integer stockCount) { if (redisService.isExists(cacheKey)) return; redisService.setAtomicLong(cacheKey, stockCount); }四、SKU库存扣减一致性处理1.细节部分这里我们注意到在责任链的第二个SKU库存扣减节点中有这么一个方法来用于redis缓存中的SKU库存扣减// 扣减库存 boolean status activityDispatch.subtractionActivitySkuStock(activitySkuEntity.getSku(), activityEntity.getEndDateTime());这里我将完整链路展示出来ActivitySkuStockActionChain - ActivityArmory - ActivityRepositoryActivityArmory:Override public Boolean subtractionActivitySkuStock(Long sku, Date endDateTime) { String cacheKey Constants.RedisKey.ACTIVITY_SKU_STOCK_COUNT_KEY sku; return activityRepository.subtractionActivitySkuStock(sku,cacheKey,endDateTime); }ActivityRepository:Override public Boolean subtractionActivitySkuStock(Long sku, String cacheKey, Date endDateTime) { long surplus redisService.decr(cacheKey); if (surplus 0){ //库存消耗完毕。发送MQ消息直接清空库存 eventPublisher.publish(activitySkuStockZeroMessageEvent.topic(),activitySkuStockZeroMessageEvent.buildEventMessage(sku)); return false; } else if (surplus0) { //库存小于0恢复为0 redisService.setAtomicLong(cacheKey,0); return false; } //1.按照cacheKey decr 后的值如99 98 97 String lockKey cacheKey Constants.UNDERLINE surplus; long expireMillis endDateTime.getTime() - System.currentTimeMillis() TimeUnit.DAYS.toMillis(1); Boolean lock redisService.setNx(lockKey, expireMillis, TimeUnit.MILLISECONDS); if (!lock){ log.info(活动sku库存加锁失败 {},lockKey); } return lock; }而异步扣减的定时任务如下每五秒去扣减一个在延迟队列中的库存/** * author 印东升 * description 更新活动sku库存任务 * create 2026-07-21 11:37 */ Slf4j public class UpdateActivitySkuStockJob { Resource private ISkuStock skuStock; Scheduled(cron 0/5 * * * * ?) public void exec(){ try{ log.info(定时任务更新活动sku库存【延迟队列获取降低对数据库的更新频次不要产生竞争】); ActivitySkuStockKeyVO activitySkuStockKeyVO skuStock.takeQueueValue(); if (nullactivitySkuStockKeyVO) return; log.info(定时任务更新活动sku库存 sku:{} activityId:{},activitySkuStockKeyVO.getSku(),activitySkuStockKeyVO.getActivityId()); skuStock.updateActivitySkuStock(activitySkuStockKeyVO.getSku()); } catch (Exception e) { log.error(定时任务更新活动sku库存失败,e); } } }2.MQ部分注意对于库存的异步扣减依旧是采用的Redisson的延迟队列一个一个扣减MQ的作用仅仅是在库存为0后停止任务调度节约定时任务的开销。而MQ部分也就是下面的代码就是将库存为0的消息传给监听器监听器去直接清空数据库中的库存并停止定时任务。if (surplus 0){ //库存消耗完毕。发送MQ消息直接清空库存 eventPublisher.publish(activitySkuStockZeroMessageEvent.topic(),activitySkuStockZeroMessageEvent.buildEventMessage(sku)); return false; }这个“库存为0的消息”我们用一个类来表示并且指定了交换机和传的信息清空指定sku的库存。/** * author 印东升 * description 活动sku库存清空消息 * create 2026-07-21 11:02 */ Component public class ActivitySkuStockZeroMessageEvent extends BaseEventLong { Value(${spring.rabbitmq.topic.activity_sku_stock_zero}) private String topic; Override public EventMessageLong buildEventMessage(Long sku) { return EventMessage.Longbuilder() .id(RandomStringUtils.randomNumeric(11)) .timestamp(new Date()) .data(sku) .build(); } Override public String topic() { return topic; } }监听器用于清空队列并且清空库存/** * author 印东升 * description 活动sku库存耗尽 * create 2026-07-21 11:42 */ Slf4j Component public class ActivitySkuStockZeroCustomer { Value(activity_sku_stock_zero) private String topic; Resource private ISkuStock skuStock; RabbitListener(queuesToDeclare Queue(value activity_sku_stock_zero)) public void listener(String message){ try{ log.info(监听活动sku库存消耗为0消息 topic:{} message:{},topic,message); //转换对象 BaseEvent.EventMessageLong eventMessage JSON.parseObject(message,new TypeReferenceBaseEvent.EventMessageLong(){ }.getType()); Long sku eventMessage.getData(); //更新库存 skuStock.clearActivitySkuStock(sku); //清空队列(此时就不需要延迟更新数据库记录了) skuStock.clearQueueValue(); }catch (Exception e){ log.error(监听活动sku库存消耗为0消息消费失败 topic:{} message:{},topic,message); throw e; } } }而具体的实现清空队列清空库存放在了Service中实现然后再在仓储中实现其实就是个透传仓储透到服务服务再用到监听器。Override public ActivitySkuStockKeyVO takeQueueValue() throws InterruptedException { return activityRepository.takeQueueValue(); } Override public void clearQueueValue() { activityRepository.clearQueueValue(); } Override public void updateActivitySkuStock(Long sku) { activityRepository.updateActivitySkuStock(sku); } Override public void clearActivitySkuStock(Long sku) { activityRepository.clearActivitySkuStock(sku); }Override public ActivitySkuStockKeyVO takeQueueValue() { String cacheKey Constants.RedisKey.ACTIVITY_SKU_COUNT_QUERY_KEY; RBlockingQueueActivitySkuStockKeyVO destinationQueue redisService.getBlockingQueue(cacheKey); return destinationQueue.poll(); } Override public void clearActivitySkuStock(Long sku) { raffleActivitySkuDao.clearActivitySkuStock(sku); } Override public void clearQueueValue() { String cacheKey Constants.RedisKey.ACTIVITY_SKU_COUNT_QUERY_KEY; RBlockingQueueActivitySkuStockKeyVO destinationQueue redisService.getBlockingQueue(cacheKey); destinationQueue.clear(); } Override public void updateActivitySkuStock(Long sku) { raffleActivitySkuDao.updateActivitySkuStock(sku); }五、测试/** * 测试库存消耗和最终一致更新 * 1. raffle_activity_sku 库表库存可以设置20个 * 2. 清空 redis 缓存 flushall * 3. for 循环20次消耗完库存最终数据库剩余库存为0 */ Test public void test_createSkuRechargeOrder() throws InterruptedException { for (int i 0; i 20; i) { try { SkuRechargeEntity skuRechargeEntity new SkuRechargeEntity(); skuRechargeEntity.setUserId(xiaofuge); skuRechargeEntity.setSku(9011L); // outBusinessNo 作为幂等仿重使用同一个业务单号2次使用会抛出索引冲突 Duplicate entry 700091009111 for key uq_out_business_no 确保唯一性。 skuRechargeEntity.setOutBusinessNo(RandomStringUtils.randomNumeric(12)); String orderId raffleOrder.createSkuRechargeOrder(skuRechargeEntity); log.info(测试结果{}, orderId); } catch (AppException e) { log.warn(e.getInfo()); } } new CountDownLatch(1).await(); }当库存扣减完了MQ发送消息到监听器而如果只是扣减部分就会通过定时任务一个一个地扣减六、Bug修改和细节优化1.Bug优化测试看上去是没有问题的但是实际上这里有一个bug就是当我们有多个sku时这些sku共用一个延迟队列当其中一个sku的库存为0时整个延迟队列的其他sku也会被全部删除。因此这里是有问题的注释掉就好只是以前残留的sku库存为0的消息还会在队列中//清空队列(此时就不需要延迟更新数据库记录了) 有点bug //skuStock.clearQueueValue();优化节省 处理sku库存为0的消息的 性能//优化每次判断出队列的SKU是否库存已经为0为0就跳过节省去update的性能 Long currentStock skuStock.getSkuStockSurplus(activity_sku_stock_count_key activitySkuStockKeyVO.getSku()); if (currentStock 0) { log.info(定时任务活动sku库存已经为0跳过该消息 sku:{} activityId:{}, activitySkuStockKeyVO.getSku(), activitySkuStockKeyVO.getActivityId()); return; }2.延迟优化延迟优化首先做异步更新的目的是减少直接访问数据库的次数因此基于这一点其他的部分可以优化比如下面是我给出的优化通过循环将连续的、sku的库存为已经0的全部一次性跳过这样可以大幅降低redis和mysql数据库的数据更新的延迟假如有20条这样的信息我通过一次性跳过就可以将延迟缩短100秒并且是不影响性能的本质是在内存中跑不会直接访问数据库。Slf4j Component public class UpdateActivitySkuStockJob { Resource private ISkuStock skuStock; Scheduled(cron 0/5 * * * * ?) public void exec() { try { while(true) {//优化让那些sku库存为0的消息一次性跳过如果连续的话 log.info(定时任务更新活动sku库存【延迟队列获取降低对数据库的更新频次不要产生竞争】); ActivitySkuStockKeyVO activitySkuStockKeyVO skuStock.takeQueueValue(); if (null activitySkuStockKeyVO) return; //优化每次判断出队列的SKU是否库存已经为0为0就跳过节省去update的性能 Long currentStock skuStock.getSkuStockSurplus(activity_sku_stock_count_key activitySkuStockKeyVO.getSku()); if (currentStock 0) { log.info(定时任务活动sku库存已经为0跳过该消息 sku:{} activityId:{}, activitySkuStockKeyVO.getSku(), activitySkuStockKeyVO.getActivityId()); continue; } log.info(定时任务更新活动sku库存 sku:{} activityId:{}, activitySkuStockKeyVO.getSku(), activitySkuStockKeyVO.getActivityId()); skuStock.updateActivitySkuStock(activitySkuStockKeyVO.getSku()); break; } } catch (Exception e) { log.error(定时任务更新活动sku库存失败, e); } } }