目录
一、业务背景
当前开发的模块需要一个高效的批量券码导入方案。初步方案处理10万耗时约360秒,请你提供一个更加快捷的方案。
二、分析
先让我们来看一下初步方案:

虽然DB层完全异步处理,但是任务更新完成状态还是在DB层处理完成后才发送的,所以实际任务完成时长是包括同步+异步全流程的,是准确的。
这里我们不会做DB层的优化,具体原因相信大家也已经知道了:
之所以当前方案总体执行时长长达六分钟左右,就是因为其每解析一行就至少要发送2次网络IO,而文件一共10w行,那么解析就要发送20w次网络IO!
而DB层是用MQ异步处理,每5000批次处理一次,其如下:
<insert id="insertBatch"> INSERT INTO user_coupon (id, user_id, coupon_template_id, receive_time, receive_count, valid_start_time, valid_end_time, use_time, source, status, create_time, update_time, del_flag) VALUES <foreach collection="list" item="item" separator=","> (#{item.id}, #{item.userId}, #{item.couponTemplateId}, #{item.receiveTime}, #{item.receiveCount}, #{item.validStartTime}, #{item.validEndTime}, #{item.useTime}, #{item.source}, #{item.status}, #{item.createTime}, #{item.updateTime}, #{item.delFlag}) </foreach> </insert>所以本次业务场景的重点耗时域肯定就在监听器逐行解析部分。
而对于这种优化,我们已经有了前车之鉴:
没错,只要使用管道批次发起请求即可。
但是这原方案之所以一行一行地发送请求是为了确保即使在应用发生宕机的情况下,也只会丢失一行数据而已,很容易就恢复了。如果我们改成管道,那么一旦发生宕机就会一下子损失一批数据,那么这该怎么恢复呢?后面再讲这个,现在先讲一讲如何优化:
首先对于脚本,其完成了三个动作:一是扣减库存,二是添加用户优惠券记录、三是更新进度。
那么添加记录这一步可以使用管道进行处理,那么管道批次设置多少合适呢?对于MQ批次我们采用的是5000行发一次,那么管道批次我们采用2500行即可,整个体系就是每发两次管道发一次MQ。
然后我们保证管道发送完成后,再扣减库存,注意此时扣减库存就要一次性扣减2500了。
虽然我们把一次Lua脚本的三个请求给拆分出来了,但是实际发起请求次数肯定比Lua脚本的要少得多,因为每2500行只发送3次redis请求,原方案需要发送2500次。
然后就是对于获取进度的请求,我们将其和管道批次捆绑,每发送一次管道后,再更新进度到redis。
剩下一个请求就是解析行的第一个请求:获取redis进度。这里我们直接大胆抛弃掉这个请求,至于为什么,请听我细细道来:
我们只需要在应用程序启动后获取一次进度即可,因为假设redis进度为5000的时候,我们已经到了redis批次准备发送管道的时候应用程序宕机了,这个时候重启后获取进度获取到5000,直接从5000开始解析即可,因为刚才的宕机之前并没有对redis层进行任何操作。
那么如果刚才宕机之前已经发送了一次redis管道呢,那么就重新发一次呗,我们的redis用户优惠券集合是Set结构,会自动去重,所以重新执行一次也完全没有任何影响。
万一连扣减库存的操作也做了呢?这个总不能重新扣减一次吧?是的,所以我们在获取进度的时候拿到优惠券模板的总库存,然后再减去redis当前进度的结果赋值给redis库存,算是一个回滚操作,这样的话就不会重复扣减库存了。
那如果已经发送了一次MQ的话该怎么办? 此时可能出现发送相同消息到MQ,那就是消费幂等性的问题。 可以选择将每条消息的唯一标识放入redis,设置TTL为5分钟。当然这样就多了一次获取和设置的网络IO次数。 也可以选择在SQL语句中加上IGNORE关键字,这样即使重复插入也会忽略报错,不中断程序继续运行,然后我们再获取这条语句影响的行数,如果为0但是又没有抛出异常,就证明是重复消息,直接返回继续消费下一条消息即可。
INSERT IGNORE INTO t_user_coupon (id, user_id, coupon_template_id, receive_time, receive_count, valid_start_time, valid_end_time, use_time, source, status, create_time, update_time, del_flag) VALUES 一键获取完整项目代码XML(#{item.id}, #{item.userId}, #{item.couponTemplateId}, #{item.receiveTime}, #{item.receiveCount}, #{item.validStartTime}, #{item.validEndTime}, #{item.useTime}, #{item.source}, #{item.status}, #{item.createTime}, #{item.updateTime}, #{item.delFlag})
因为宕机的概率非常小,所以我们其实是以牺牲宕机恢复的时间换取了正常运行的效率。
然后为了保证防止库存不足导致的管道发送成功库存扣减失败的场景,我们在发送管道之前应该发送一次redis请求检查库存是否充足。
这样一来,原本的20w次网络IO立马减少变成 10w/5000 * [(1次检查库存 + 1次管道 + 1次扣减库存 + 1次更新进度) * 2 + 1次MQ] + 1次宕机恢复进度 + 1次宕机恢复库存 = 182次,而且经测试,此方案的执行时间从原本的360s优化到了15 - 20s的区间! 可谓巨大提升!
具体如下:

然后我们还需要进一步完善细节,比如说错误数据记录、重试机制、日志记录等等。
详细请见下面的代码实现。
对了,因为此方案是IO集中型任务,需要经常阻塞,所以推荐使用虚拟线程防止占用平台线程。 具体使用请见下面这篇博客: 带你轻松学习虚拟线程和StructuredTaskScope-CSDN博客https://blog.csdn.net/2401_88959292/article/details/150288495?spm=1001.2014.3001.5502
三、代码实现
@Slf4j@RequiredArgsConstructorpublic class ReadExcelDistributionListener extends AnalysisEventListener<CouponTaskExcelObject> {
private final CouponTaskDO couponTaskDO; private final CouponTemplateDO couponTemplateDO; private final CouponTaskFailMapper couponTaskFailMapper;
private final StringRedisTemplate stringRedisTemplate; private final CouponExecuteDistributionProducer couponExecuteDistributionProducer;
// 计数器 private final CountDownLatch processingLatch; // 停止执行标记 private volatile boolean stopProcessing = false; // 确保回调只执行一次 private boolean isFinished = false; // 当前全局进度 private int rowCount = 1; // 当前内存进度 private int memoryCount = -1; // 5000发一次MQ private final static int BATCH_USER_COUPON_MQ_SIZE = 5000; // 2500发一次管道 private final static int BATCH_USER_COUPON_PIPELINE_SIZE = 2500;
// 最大重试次数 private static final int MAX_RETRY = 3; // 重试间隔时间 private static final int RETRY_INTERVAL_MS = 1000;
// 管道批次存储集合 private final List<Map<Object,Object>> batchRedisPinplineList = new ArrayList<>(BATCH_USER_COUPON_MQ_SIZE); // MQ批次存储集合 private final List<CouponTemplateDistributionSumEvent> batchMqList = new ArrayList<>(BATCH_USER_COUPON_MQ_SIZE + 1000);
// lua脚本 private static final DefaultRedisScript<Void> REBUILD_REDIS_PROGRESS_SCRIPT;
// 单例装载 static { REBUILD_REDIS_PROGRESS_SCRIPT = new DefaultRedisScript<>(); REBUILD_REDIS_PROGRESS_SCRIPT.setLocation(new ClassPathResource("lua/rebuild_redis_progress.lua")); REBUILD_REDIS_PROGRESS_SCRIPT.setResultType(Void.class); }
@Override public void invoke(CouponTaskExcelObject couponTaskExcelObject, AnalysisContext analysisContext) {
// 停止执行 if(stopProcessing){ return; }
String progressKey = String.format(DistributionRedisConstant.TEMPLATE_TASK_EXECUTE_PROGRESS_KEY, couponTaskDO.getId());
/* 如果发生批次宕机,只是重复消费幂等性问题,redis重复插入不会报错,可以忽略,MySQL的话使用IGNORE关键字忽略报错即可 */ // 在构造器初始化时读取进度,也就是内存进度为 -1 的时候 if(memoryCount == -1){ // 获取进度 String progress = stringRedisTemplate.opsForValue().get(progressKey); memoryCount = ObjectUtil.isNotNull(progress) ? Integer.parseInt(progress) : 0; rowCount = memoryCount + 1; // 回滚redis库存 int totalStock = couponTemplateDO.getStock(); int count = totalStock - memoryCount; String key = String.format(EngineRedisConstant.COUPON_TEMPLATE_KEY, couponTemplateDO.getId()); // 执行lua脚本,先删除hash库存字段,再新增即可 stringRedisTemplate.execute(REBUILD_REDIS_PROGRESS_SCRIPT, List.of(key), "stock", count ); }
// 转化成map Map<Object, Object> batchUserCoupon = MapUtil.builder() .put("userId", couponTaskExcelObject.getUserId()) .put("rowNum", rowCount) .build();
// 转化成MQ事件 CouponTemplateDistributionSumEvent event = CouponTemplateDistributionSumEvent.builder() .couponTaskId(couponTaskDO.getId()) .couponTaskBatchId(couponTaskDO.getBatchId()) .notifyType(couponTaskDO.getNotifyType()) .shopNumber(couponTaskDO.getShopNumber()) .couponTemplateId(couponTemplateDO.getId()) .validEndTime(couponTemplateDO.getValidEndTime()) .couponTemplateConsumeRule(couponTemplateDO.getConsumeRule()) .userId(couponTaskExcelObject.getUserId()) .phone(couponTaskExcelObject.getPhone()) .mail(couponTaskExcelObject.getMail()) .batchUserSetSize(BATCH_USER_COUPON_MQ_SIZE) .build();
// 放入批次集合 batchRedisPinplineList.add(batchUserCoupon); batchMqList.add(event);
// 当达到批次阈值时,执行一次管道 if(batchRedisPinplineList.size() >= BATCH_USER_COUPON_PIPELINE_SIZE){ executePipeline(couponTaskExcelObject); // 更新redis进度 stringRedisTemplate.opsForValue().set(progressKey, String.valueOf(memoryCount)); }
// 指向下一行 rowCount++; // 计数器++ processingLatch.countDown(); }
// 执行管道 private void executePipeline(CouponTaskExcelObject couponTaskExcelObject) {
// 停止执行 if(stopProcessing){ return; }
String key = String.format(DistributionRedisConstant.TEMPLATE_TASK_EXECUTE_BATCH_USER_KEY, couponTaskDO.getId());
// 是否发送成功 boolean isSuccess = false; // 当前重试次数 int retryCount = 0; // 批次写入数量 Long listLength = 0L;
// 检查库存 String stockKey = String.format(EngineRedisConstant.COUPON_TEMPLATE_KEY, couponTemplateDO.getId()); Object stockObj = stringRedisTemplate.opsForHash().get(stockKey, "stock"); Long currentStock = null; if (stockObj instanceof String) { currentStock = Long.valueOf((String) stockObj); }
if (currentStock == null || currentStock < batchRedisPinplineList.size()) { log.error("库存不足!任务ID={}, 可用库存={}, 需要={}", couponTaskDO.getId(), currentStock, batchRedisPinplineList.size()); recordFail(couponTaskExcelObject); // 此处不做回滚,只要有不能让失败的数据影响到成功的数据 // 停止后续执行 stopProcessing = true; // 清理批次数据 batchRedisPinplineList.clear(); batchMqList.clear(); return; }
// 开启管道 while(!isSuccess && retryCount < MAX_RETRY){ try { List<Object> result = stringRedisTemplate.executePipelined((RedisCallback<?>) connection -> { List<String> values = new ArrayList<>(batchRedisPinplineList.size()); for (Map<Object, Object> row : batchRedisPinplineList) { values.add(JSONUtil.toJsonStr(row)); } // 批量添加到 Set connection.setCommands().sAdd( key.getBytes(StandardCharsets.UTF_8), values.stream().map(s -> s.getBytes(StandardCharsets.UTF_8)).toArray(byte[][]::new) );
// 批量写入失败 return null; });
// 如果写入成功 if (!result.isEmpty()) { listLength = stringRedisTemplate.opsForSet().size(key); log.info("批量写入成功: 任务ID={}, 当前集合大小={}", couponTaskDO.getId(), listLength); isSuccess = true; }
}catch (Exception e){ // 重试次数++ retryCount++; log.error("批量写入redis失败,重试次数:{},错误信息:{}", retryCount, e.getMessage(), e); try { // 休眠一段时间预防网络波动 TimeUnit.MILLISECONDS.sleep(RETRY_INTERVAL_MS); } catch (InterruptedException ie) { // 重置标记防止后续休眠失效 Thread.currentThread().interrupt(); } } }
// 如果失败则记录 if(!isSuccess){ recordFail(couponTaskExcelObject); // 手动回滚 rollBack(); // 停止后续执行,因为一般失败是代码或者数据问题,一般一出问题就是都有问题 stopProcessing = true; // 清理集合 batchRedisPinplineList.clear(); batchMqList.clear(); return; }else{ // 扣减库存 Long newStock = stringRedisTemplate.opsForHash().increment(stockKey, "stock", -batchRedisPinplineList.size());
if (ObjectUtil.isNotNull(newStock) && newStock >= 0) { log.info("库存扣减成功: 任务ID={}, 新库存={}", couponTaskDO.getId(), newStock); } else { log.error("库存扣减异常!任务ID={}", couponTaskDO.getId()); recordFail(couponTaskExcelObject); // 停止后续执行,因为一般失败是代码或者数据问题,一般一出问题就是都有问题 stopProcessing = true; // 清理集合 batchRedisPinplineList.clear(); batchMqList.clear(); return; } // 更新内存进度 memoryCount = rowCount; }
// 进度到达阈值,执行一次MQ if(batchMqList.size() >= BATCH_USER_COUPON_MQ_SIZE){ CouponTemplateDistributionEvent event = CouponTemplateDistributionEvent.builder() .couponTaskId(couponTaskDO.getId()) .couponTaskBatchId(couponTaskDO.getBatchId()) .shopNumber(couponTaskDO.getShopNumber()) .couponTemplateId(couponTaskDO.getCouponTemplateId()) .batchUserSetSize(batchMqList.size()) .distributionEndFlag(false) // 深拷贝,避免引用共享引发事故 .couponTemplateDistributionSumEvents(new ArrayList<>(batchMqList)) .build();
couponExecuteDistributionProducer.sendMessage(event);
// 清空MQ批次集合 batchMqList.clear(); }
// 清空管道批次集合 batchRedisPinplineList.clear(); }
// 手动回滚 private void rollBack() { log.error("手动回滚: {}", couponTaskDO.getId()); // 逻辑暂时省略 }
@Override public void doAfterAllAnalysed(AnalysisContext analysisContext) { // 确保只执行一次 if (isFinished) return; isFinished = true; // 等待所有数据处理完成 try { processingLatch.await(); } catch (InterruptedException e) { // 此处不可能被打断,除非是也无宕机,但还是为了向后兼容,这里写一些逻辑比较好 log.error("处理数据时发生异常,异常信息:{}", e.getMessage()); // 重置标记防止后续waiting状态失效 Thread.currentThread().interrupt(); }
Long couponTaskId = couponTaskDO.getId();
// 处理剩余数据 if (!batchRedisPinplineList.isEmpty()) { // 使用最后一条数据 CouponTaskExcelObject lastObject; Map<Object, Object> last = batchRedisPinplineList.getLast(); lastObject = new CouponTaskExcelObject(); lastObject.setUserId(ObjectUtil.defaultIfNull(String.valueOf(last.get("userId")), "UNKNOWN")); lastObject.setPhone(ObjectUtil.defaultIfNull(String.valueOf(last.get("phone")), "UNKNOWN")); lastObject.setMail(ObjectUtil.defaultIfNull(String.valueOf(last.get("mail")), "UNKNOWN")); executePipeline(lastObject); }
// 将剩余数据放入MQ CouponTemplateDistributionEvent event = CouponTemplateDistributionEvent.builder() .couponTaskId(couponTaskId) .couponTaskBatchId(couponTaskDO.getBatchId()) .shopNumber(couponTaskDO.getShopNumber()) .couponTemplateId(couponTaskDO.getCouponTemplateId()) .batchUserSetSize(batchMqList.size()) .distributionEndFlag(true) .couponTemplateDistributionSumEvents(new ArrayList<>(batchMqList)) .build(); couponExecuteDistributionProducer.sendMessage(event);
// 清空集合 batchRedisPinplineList.clear(); batchMqList.clear(); }
// 记录失败 private void recordFail(CouponTaskExcelObject couponTaskExcelObject){ // 空指针防护 String userId = "UNKNOWN"; if (couponTaskExcelObject != null) { userId = couponTaskExcelObject.getUserId(); } else if (!batchRedisPinplineList.isEmpty()) { // 尝试从批次获取最后一条记录 Map<Object, Object> last = batchRedisPinplineList.getLast(); userId = String.valueOf(last.get("userId")); }
log.error("处理失败: 最后用户ID={}, 批次大小={}", userId, batchRedisPinplineList.size()); String json = "失败批次的最后一条数据的userId为:" + userId + ",失败批次数量:" + batchRedisPinplineList.size(); CouponTaskFailDO couponTaskFailDO = CouponTaskFailDO.builder() .batchId(couponTaskDO.getId()) .jsonObject(json) .build(); couponTaskFailMapper.insert(couponTaskFailDO); }
}码文不易、留个赞再走呗
原文链接: 360s→15s:近25倍的性能优化重构十万Excel券码批量导入方案 作者: Yilena
如果这篇文章对你有帮助,欢迎分享给更多人!
部分信息可能已经过时










