# tianji **Repository Path**: su27sk/tianji ## Basic Information - **Project Name**: tianji - **Description**: 本项目是一套基于Java微服务的一套在线学习教育&电商项目,基于SpringCloud,SpringCloudAlibaba,Nacos,OpenFegin,GateWay;使用RabbitMQ,Redis,DelayQueue,BitMap优化签到,数据定时持久化分表,XXL-job,异步线程池,乐观/悲观锁,分布式锁等技术 - **Primary Language**: Java - **License**: Not specified - **Default Branch**: dev_lw - **Homepage**: None - **GVP Project**: No ## Statistics - **Stars**: 0 - **Forks**: 0 - **Created**: 2026-06-26 - **Last Updated**: 2026-08-09 ## Categories & Tags **Categories**: Uncategorized **Tags**: 微服务, Java, 高并发 ## README # 教育电商 ## 项目简介 本项目是一套基于Java微服务的一套在线学习教育&电商项目,基于SpringCloud,SpringCloudAlibaba,Nacos,OpenFegin,GateWay;使用Swagger,RabbitMQ,Redis,DelayQueue,BitMap优化签到,数据定时持久化分表,XXL-job,异步线程池,乐观/悲观锁,分布式锁,设计模式 主要二次开发tj-trade(交易服务);新增tj-learning(学习服务),tj-promotion(促销服务),tj-remark(评价服务) ### 用户端页面 ![用户端首页.png](picture/%E7%94%A8%E6%88%B7%E7%AB%AF%E9%A6%96%E9%A1%B5.png) ### 管理端页面 ![管理端首页.png](picture/%E7%AE%A1%E7%90%86%E7%AB%AF%E9%A6%96%E9%A1%B5.png) ## 运维部署 启动docker系统 ```yaml systemctl restart docker docker.socket docker restart tj-gateway tj-user tj-auth tj-course tj-media tj-search tj-trade tj-pay tj-message tj-data tj-exam ``` 其次配置好虚拟机然后运行各个docker容器,之后下载SwitchHosts完成域名的配置 ```yaml 192.168.150.101 git.tianji.com 192.168.150.101 jenkins.tianji.com 192.168.150.101 mq.tianji.com 192.168.150.101 nacos.tianji.com 192.168.150.101 xxljob.tianji.com 192.168.150.101 es.tianji.com 192.168.150.101 api.tianji.com 192.168.150.101 www.tianji.com 192.168.150.101 manage.tianji.com 192.168.150.101 cpolar.tianji.com ``` 然后查看持续继承的环境 首先查看Gogs私域代码仓库:http://git.tianji.com/tjxt/tianji (tjxt/123321),然后查看Jenkins http://jenkins.tianji.com (账号:root/123),然后可以使用web钩子测试推送编译代码然后点击每个项目将其运行为docker容器 然后可以查看后端网关的swagger文档 http://192.168.150.101:10010/doc.html 。可以查看前端用户端的页面 http://www.tianji.com (jack/123 rose/123456),管理端 http://manage.tianji.com 下次开机可以重启java微服务容器:docker restart tj-gateway tj-user tj-auth tj-course tj-media tj-search tj-trade tj-pay tj-message tj-data tj-exam ```yaml Git 私服:域名git.tianji.com,账号 tjxt/123321,端口 10880 Jenkins 持续集成:域名jenkins.tianji.com,账号 root/123,端口 18080 RabbitMQ:域名mq.tianji.com,账号 tjxt/123321,端口 15672 Nacos 控制台:域名nacos.tianji.com,账号 nacos/nacos,端口 8848 Redis 密码 123321,端口 6379 xxl-job 控制台:域名xxljob.tianji.com,账号 admin/123456,端口 8880 ES 的 Kibana 控制台:域名es.tianji.com,无账号,端口 5601 微服务网关:域名api.tianji.com,无账号,端口 10010 用户端入口:域名www.tianji.com,jack/123 账号 rose/123456,端口 18081 管理端入口:域名manage.tianji.com,默认项目账号,端口 18082 ``` ## 修复Bug热身 首先我们发现一个Bug,发现如果是rose用户登录的情况下无法删除她相关订单下的订单 ![热身Bug.png](picture/%E7%83%AD%E8%BA%ABBug.png) ### 找到Bug 第一步是前端点击F12查看调用的前端请求,发现请求url是 http://api.tianji.com/ts/orders/1597502678241378305(DELETE) ,然后打开虚拟机服务器查看nginx的配置, 打开目录 cd /usr/local/src/nginx/conf,然后查看nginx.conf文件,看到 ```yaml server { listen 80; server_name api.tianji.com; location / { proxy_set_header Host $host; # 传递原始Host头 proxy_pass http://192.168.150.101:10010; # proxy_pass http://192.168.150.1:10010; # 由于要修改网关在本地测试,所以代理到本地IP } } ``` 发现前端的api.tianji.com路径被映射到10010的网关模块,然后打开IDEA查看网关的配置文件bootstrap.yml,发现 ```yaml - id: ts uri: lb://trade-service predicates: - Path=/ts/** ``` 发现这个路由的ts匹配到前端url的 http://api.tianji.com/ts/orders/1597502678241378305(DELETE) 的/ts,接下来就去查看trade-service的代码。 找到剩余路径为/order并且请求方式为DELETE的接口,并找到业务层代码如下: ```java @Override public void deleteOrder(Long id) { // 1.获取登录用户 Long userId = UserContext.getUser(); // 2.查询订单 Order order = getById(id); if (order == null) { return; } // 3.判断订单所属用户与当前登录用户是否一致 if(userId != order.getUserId()){ // 不一致,说明不是当前用户的订单,结束 throw new BadRequestException("不能删除他人订单"); } // 4.删除订单 boolean success = removeById(id); if (!success) { throw new DbException(OPERATE_FAILED); } } ``` 最后定位到是因为都是使用Long这个包装类,129的字面量不再享元范围内,所以要使用equles方法来进行比较,核心代码修改为 ```java @Override public void deleteOrder(Long id) { // 1.获取登录用户 Long userId = UserContext.getUser(); // 2.查询订单 Order order = getById(id); if (order == null) { return; } // 3.判断订单所属用户与当前登录用户是否一致 if(!userId.equals(order.getUserId())){ // 不一致,说明不是当前用户的订单,结束 // TODO 成功解决一个Bug throw new BadRequestException("不能删除他人订单"); } // 4.删除订单 boolean success = removeById(id); if (!success) { throw new DbException(OPERATE_FAILED); } System.out.println("删除订单"); } ``` ### 解决Bug 我们可以使用IDEA的远程调用功能,同时结合jenkins的配置 ![远程调试.png](picture/%E8%BF%9C%E7%A8%8B%E8%B0%83%E8%AF%95.png) 然后因为我们使用的是dev_lw分支下下开发,将jenkins的CICD分支调整为dev_lw分支,然后重新推送 ```yaml git push gitee dev_lw # 推送到开源码云 git push origin dev_lw # 推送到jenkins ``` 可以在Linux虚拟机下查看tj-trade这个镜像的运行日志 ```yaml dlog -f tj-trade # 持续跟踪查看镜像 ``` ## 学习服务模块 学习服务模块是一个全新的模块,从0到1开发 ### 创建课表 ![创建订单.png](picture/%E5%88%9B%E5%BB%BA%E8%AE%A2%E5%8D%95.png) 点击立刻报名,首先tj-trade(交易服务)创建订单,然后通过MQ发送数据给tj-learning(学习服务)创建课表 首先是tj-trade(交易服务)通过RabbitMQ发送消息 ```java @Override @Transactional public PlaceOrderResultVO enrolledFreeCourse(Long courseId) { ..... // 5.发送MQ消息,通知报名成功 rabbitMqHelper.send( MqConstants.Exchange.ORDER_EXCHANGE, MqConstants.Key.ORDER_PAY_KEY, OrderBasicDTO.builder() .orderId(orderId) .userId(userId) .courseIds(cIds) .finishTime(order.getFinishTime()) .build() ); ...... } ``` 然后tj-learning(学习服务)创建消费者方法接收消息 ```java @Slf4j @Component @RequiredArgsConstructor public class LessonChangeListener { private final ILearningLessonService lessonService; @RabbitListener(bindings = @QueueBinding( value = @Queue(value = "learning.lesson.pay.queue", durable = "true"), exchange = @Exchange(name = MqConstants.Exchange.ORDER_EXCHANGE, type = ExchangeTypes.TOPIC), key = MqConstants.Key.ORDER_PAY_KEY )) public void listenLessonPay(OrderBasicDTO order){ // 1.健壮性判断 if(order == null || order.getUserId() == null || CollUtils.isEmpty(order.getCourseIds())) { log.error("接收到的MQ消息有误,订单数据为空"); return; } // 2.添加课程 log.debug("监听到用户{}的订单{},需要添加课程{}到课表中",order.getUserId(),order.getOrderId(),order.getCourseIds()); lessonService.addUserLessons(order.getUserId(), order.getCourseIds()); } } ``` 然后因为涉及到课程的有效期,我们需要调用tj-course(课程服务)根据课程id去查询课程的详细信息,我们去tj-api(工具服务)找合适的client,tj-api专门负责这个 ```java @GetMapping("/courses/simpleInfo/list") List getSimpleInfoList(@RequestParam("ids") Iterable ids); ``` 然后完成tj-learning(学习服务)课表的批量新增 ```java @Service @RequiredArgsConstructor @Slf4j @Transactional public class LearningLessonServiceImpl extends ServiceImpl implements ILearningLessonService { private final CourseClient courseClient; @Override public void addUserLessons(Long userId, List courseIds) { // 1.查询课程有效期 List cInfoList = courseClient.getSimpleInfoList(courseIds); if(CollUtils.isEmpty(cInfoList)) { // 课程不存在,无法添加 log.error("课程信息不存在,无法添加到课表"); return; } // 2.循环遍历,处理LearningLesson数据 List list = new ArrayList<>(cInfoList.size()); for (CourseSimpleInfoDTO cInfo : cInfoList) { LearningLesson lesson = new LearningLesson(); // 2.1获取过期时间 Integer validDuration = cInfo.getValidDuration(); if(validDuration != null && validDuration > 0) { LocalDateTime now = LocalDateTime.now(); lesson.setCreateTime(now); lesson.setExpireTime(now.plusMonths(validDuration)); } // 2.2 填充userId和courseId lesson.setUserId(userId); lesson.setCourseId(cInfo.getId()); list.add(lesson); } // 3 批量新增 saveBatch(list); } } ``` 另外为了保证一个用户对于一个课程不会重复创建课表,在learning_lesson课程表里面创建了唯一联合索引依赖 ```SQL constraint idx_user_id unique (user_id, course_id) ``` ### 查询当前用户的课表 让当前登陆的用户默认根据最近学习时间进行降序排序,然后使用tj-course(课程服务)根据课程id去查询课程的详细信息,最后进行拼接成为vo ![查询课表.png](picture/%E6%9F%A5%E8%AF%A2%E8%AF%BE%E8%A1%A8.png) ```java @Override public PageDTO queryMyLesson(PageQuery query) { // 1.获取当前登录用户 Long userId = UserContext.getUser(); // 2.分页查询 Page page = lambdaQuery() .eq(LearningLesson::getUserId, userId) .page(query.toMpPage("latest_learn_time", false)); List records = page.getRecords(); if(CollUtils.isEmpty(records)){ return PageDTO.empty(page); } // 3.查询课程信息 // 3.1 获取课程id Set cIds = records.stream().map(LearningLesson::getCourseId).collect(Collectors.toSet()); // 3.2 查询课程信息 List cInfoList = courseClient.getSimpleInfoList(cIds); if(CollUtils.isEmpty(cInfoList)){ // 课程不存在,无法添加 throw new BadRequestException("课程信息不存在! "); } // 3.3 将课程信息集合处理为Map,key是courseId,value是course本身 Map cMap = cInfoList.stream(). collect(Collectors.toMap(CourseSimpleInfoDTO::getId, c -> c)); // 4.封装VO List list = new ArrayList<>(records.size()); // 4.1 循环遍历,po转vo for (LearningLesson r : records) { // 4.2 拷贝LearningLesson到LearningLessonVO里面 LearningLessonVO vo = BeanUtils.copyBean(r, LearningLessonVO.class); // 4.3 获取课程信息,填充到vo CourseSimpleInfoDTO cInfo = cMap.get(r.getCourseId()); vo.setCourseName(cInfo.getName()); vo.setCourseCoverUrl(cInfo.getCoverUrl()); vo.setSections(cInfo.getSectionNum()); list.add(vo); } return PageDTO.of(page,list); } ``` ### 修复Swagger请求头 因为我们使用Swagger进行接口测试需要user-info这个请求头,但是源代码没有配置,我自己修改了Knife4jConfiguration.java文件 ```java @Configuration @ConditionalOnProperty(prefix = "tj.swagger", name = "enable",havingValue = "true") @EnableConfigurationProperties(SwaggerConfigProperties.class) public class Knife4jConfiguration { @Resource private SwaggerConfigProperties swaggerConfigProperties; @Bean(value = "defaultApi2") public Docket defaultApi2(TypeResolver typeResolver) { // 构建 user-info 请求头 RequestParameter userInfoHeader = new RequestParameterBuilder() .name("user-info") .description("用户基础信息") .in(ParameterType.HEADER) // 指定为请求头 .required(true) .query(param -> param.model(model -> model.scalarModel(ScalarType.STRING))) .build(); // 1.初始化Docket Docket docket = new Docket(DocumentationType.SWAGGER_2); // 2.是否需要包装R if (swaggerConfigProperties.getEnableResponseWrap()) { docket.additionalModels(typeResolver.resolve(R.class)); } return docket.apiInfo(new ApiInfoBuilder() .title(this.swaggerConfigProperties.getTitle()) .description(this.swaggerConfigProperties.getDescription()) .contact(new Contact( this.swaggerConfigProperties.getContactName(), this.swaggerConfigProperties.getContactUrl(), this.swaggerConfigProperties.getContactEmail())) .version(this.swaggerConfigProperties.getVersion()) .build()) .select() //这里指定Controller扫描包路径 .apis(RequestHandlerSelectors.basePackage(swaggerConfigProperties.getPackagePath())) .paths(PathSelectors.any()) .build() .globalRequestParameters(List.of(userInfoHeader)); } } ``` 最后成功加载user-info请求头 ![修复Swagger请求头.png](picture/%E4%BF%AE%E5%A4%8DSwagger%E8%AF%B7%E6%B1%82%E5%A4%B4.png) ## 学习计划和进度 ![学习服务总流程图.png](picture/%E5%AD%A6%E4%B9%A0%E6%9C%8D%E5%8A%A1%E6%80%BB%E6%B5%81%E7%A8%8B%E5%9B%BE.png) ### 提交视频考试学习记录 我们开启一个学习计划的小结有看视频和考试两种,视频需要看之前这节有没有学习记录,有的话就是修改学习记录,如果是第一次学完(进度 >= 50% 且旧状态为未完成)就修改任务学习状态完成修改完成时间;如果没有就新增学习记录;考试就默认学习状态完成。 最后修改课表状态已学习小结数量,是否课程学习完结,最新学习小结和最新学习时间。 前端也要配合,每隔5s定时发送修改视频考试学习记录的接口请求,保证一边学一边修改学习状态,最多出现5s这个极短的用户可容忍误差。 ![提交学习记录-old.png](picture/%E6%8F%90%E4%BA%A4%E5%AD%A6%E4%B9%A0%E8%AE%B0%E5%BD%95-old.png) ```java @Override @Transactional public void addLearningRecord(LearningRecordFormDTO recordDTO) { // 1.获取登录用户 Long userId = UserContext.getUser(); // 2.处理学习记录 boolean finished = false; if(recordDTO.getSectionType() == SectionType.VIDEO) { // 2.1 处理视频 finished = handleVideRecord(userId, recordDTO); }else { // 2.2 处理考试 finished = handleExamRecord(userId, recordDTO); } // 3.处理课表数据 handleLearningLessonsChanges(recordDTO,finished); } private void handleLearningLessonsChanges(LearningRecordFormDTO recordDTO, boolean finished) { // 1.查询课表 LearningLesson lesson = lessonService.getById(recordDTO.getLessonId()); if (lesson == null){ throw new BizIllegalException("课表不存在,无法更新数据! "); } boolean allLearned = false; // 是否全部小结学完 // 2.判断是否有新的完成小结 if(finished) { // 3.如果有新完成的小结,则需要查询课程数据 CourseFullInfoDTO cInfo = courseClient.getCourseInfoById(lesson.getCourseId(),false,false); if(cInfo == null){ throw new BizIllegalException("课程不存在,无法更新数据!"); } // 4. 比较课程是否全部学完;已学习小结 >= 课程总小结 allLearned = lesson.getLearnedSections() + 1 >= cInfo.getSectionNum(); } // 5.更新课表 lessonService.lambdaUpdate() .set(lesson.getLearnedSections() == 0,LearningLesson::getStatus,LessonStatus.LEARNING.getValue()) .set(allLearned,LearningLesson::getStatus, LessonStatus.FINISHED.getValue()) .set(!finished, LearningLesson::getLatestSectionId, recordDTO.getSectionId()) .set(!finished,LearningLesson::getLatestLearnTime,recordDTO.getCommitTime()) .setSql(finished,"learned_sections = learned_sections + 1") .eq(LearningLesson::getId,lesson.getId()) .update(); } private boolean handleExamRecord(Long userId, LearningRecordFormDTO recordDTO) { // 1.转换DTO为PO LearningRecord record = BeanUtils.copyBean(recordDTO, LearningRecord.class); // 2.填充数据 record.setUserId(userId); record.setFinished(true); record.setFinishTime(recordDTO.getCommitTime()); // 3.写入数据库 boolean success = save(record); if(!success){ throw new DbException("新增考试记录失败! "); } return true; } private boolean handleVideRecord(Long userId, LearningRecordFormDTO recordDTO) { // 1.查询旧的学习记录 LearningRecord old = lambdaQuery() .eq(LearningRecord::getLessonId, recordDTO.getLessonId()) .eq(LearningRecord::getSectionId, recordDTO.getSectionId()) .one(); // 2.判断是否存在,存在就更新,不存在就新增 if (old == null){ // 3.不存在则新增 // 3.1 转换PO LearningRecord record = BeanUtils.copyBean(recordDTO, LearningRecord.class); // 3.2 填充数据 record.setUserId(userId); // 3.3 写入数据库 boolean success = save(record); if(!success){ throw new DbException("新增学习记录失败! "); } return false; } // 4. 存在则更新 // 4.1 判断是否是第一次完成(之前学过了但是没学完,finished为false;并且本次的学习进度超过50%) boolean finished = !old.getFinished() && recordDTO.getMoment() * 2 >= recordDTO.getDuration(); // 4.2 更新数据 boolean success = lambdaUpdate() .set(LearningRecord::getMoment, recordDTO.getMoment()) .set(finished, LearningRecord::getFinished, true) .set(finished, LearningRecord::getFinishTime, recordDTO.getCommitTime()) .eq(LearningRecord::getId, old.getId()) .update(); if(!success){ throw new DbException("更新学习记录失败! "); } return finished; } ``` ## 高并发性能升级 ### 高并发写优化 - MQ异步调用:可以用MQ做异步请求,做到流量削峰;缺点是没有降低写的次数,只是降低了频率 - Redis合并写请求:将多个服务的请求积攒到Redis,而且Redis请求IO速度很快,然后定义一起请求数据库,降低了写的次数;缺点是架构复杂而且不支持事务性 ### 总架构梳理 - 业务修改痛点: 我们通过梳理提交视频考试学习记录这个业务逻辑,发现考试学习场景,视频学习记录不存在新增记录以及视频学习记录第一次完成都属于低频行为, 最需要优化非第一次完成小节修改学习记录的情况(old != null && finished == false)。 - Redis结构设计: 因为Redis里面key需要单独的存储空间,要尽量节省使用,所以我们采用lessonId(课表id)作为key而不是sectionId(小节id)作为key,key的数量可以变少。然后我们查询修改sectionId(小节id),moment(播放进度),finished(是否学习完成)三个字段是最需要的。 所以我们采用hash数据结构,lessonId(课表id)作为key,sectionId(小节id)作为HashKey,sectionId(小节id),moment(播放进度),finished(是否学习完成)三个字段作为HashValue ![高并发写缓存数据结构.png](picture/%E9%AB%98%E5%B9%B6%E5%8F%91%E5%86%99%E7%BC%93%E5%AD%98%E6%95%B0%E6%8D%AE%E7%BB%93%E6%9E%84.png) - 保证缓存数据一致性: 首先我们看一下加入高并发写优化的逻辑流程图 ![提交学习记录-new.png](picture/%E6%8F%90%E4%BA%A4%E5%AD%A6%E4%B9%A0%E8%AE%B0%E5%BD%95-new.png) 分析一下为什么第一次学完后修改小节学习状态后要清除缓存。因为我们第一次学完要把finished改为true,如果不做任何操作,那么数据库里面finished为true,缓存里面finished为false,下次查询记录是否存在优先查缓存,发现finished为false,就出现了缓存数据不一致。 那我们把缓存的finished修改为true呢,也存在问题,如果修改课表的情况下发生了异常,代码回滚,数据库的finished为false,缓存依然为finished为true,依然出现数据不一致。所以我们干脆删除缓存。一切只看数据库,下次查询记录是否存在的时候无缓存,直接查数据库,避开缓存数据不一致问题 - 延时任务: 我们每隔15s提交一次学习记录,假如是视频有30分钟,肯定会有上百次提交,但是我们只想退出视频前的最后一次提交记录的moment(播放进度)才有意义,之前的都是被覆盖的过期数据。所以刚刚的流程图里面 涉及到一个延时任务。每次缓存Redis的时候我们将moment记录在本地,每隔15s提交一次,我们就20s后延时将本地的moment与缓存Redis的moment进行比较,如果与缓存Redis的moment数据不一致,那就说明 还有新的提交,这次的就先不将Redis写入数据库,继续等;如果一致那就说明用户退出视频了,那20s前提交的缓存数据就算最新的moment,那就将Redis写入数据库 ### 延时任务工具类 在高性能优化的时候需要用到延时队列,所以这里需要进行延时任务的选型,这里就用比较简单快捷的JDK自带的DelayQueue 首先要完成一个延时任务工具类 ```java // 延时队列工具类 @Data public class DelayTask implements Delayed { private D data; // 使用延时队列要传入的数据 private long deadlineNanos; // 延时队列要执行的时间,精确到纳秒 public DelayTask(D data, Duration delayTime){ this.data = data; this.deadlineNanos = System.nanoTime() + delayTime.toNanos(); } // 获取当前时间距离延时队列要执行时间的时间 @Override public long getDelay(TimeUnit unit) { // this.convert(时长, 该时长对应的原始单位) return unit.convert(Math.max(0,deadlineNanos - System.nanoTime()),TimeUnit.NANOSECONDS); } @Override public int compareTo(Delayed o) { long l = getDelay(TimeUnit.NANOSECONDS) - o.getDelay(TimeUnit.NANOSECONDS); if (l == 0){ return 0; } return l > 0 ? 1 : -1; } } ``` 然后在LearningRecordDelayTaskHandler这个总的Redis & DelayQueue工具类中进行使用 ```java // 延时任务进行比较 public void handleDelayTask(){ while (begin){ try { // 1.获取到期的延时任务 DelayTask task = queue.take(); RecordTaskData data = task.getData(); // 2.查询Redis缓存 LearningRecord record = readRecordCache(data.getLessonId(), data.getSectionId()); if(record == null){ continue; } // 3.比较数据,moment值 if(!Objects.equals(record.getMoment(), data.getMoment())) { // 不一致,说明用户提交了新的播放进度,放弃旧数据 continue; } // 4.一致,持久化播放进度数据到数据库 // 4.1 更新学习记录的moment record.setFinished(null); // mp设置非 NULL 才更新,所以我们设置为null,避免finished属性更新,只更新moment recordMapper.updateById(record); // 4.2 更新课表最近学习时间 LearningLesson lesson = new LearningLesson(); lesson.setId(data.getLessonId()); lesson.setLatestSectionId(data.getSectionId()); lesson.setLatestLearnTime(LocalDateTime.now().minusSeconds(20)); lessonService.updateById(lesson); } catch (Exception e) { log.error("处理延时任务发生异常",e); } } } ``` ### 最后业务代码完成修改 ```java @Override public LearningLessonDTO queryLearningRecordByCourse(Long courseId) { // 1.获取登录用户 Long userId = UserContext.getUser(); // 2.查询课表 LearningLesson lesson = lessonService.queryByUserIdAndCourseId(userId, courseId); if(lesson == null){ return null; } // 3.查询学习记录 List records = lambdaQuery().eq(LearningRecord::getLessonId, lesson.getId()).list(); // 4.封装结果 LearningLessonDTO dto = new LearningLessonDTO(); dto.setId(lesson.getId()); dto.setLatestSectionId(lesson.getLatestSectionId()); dto.setRecords(BeanUtils.copyList(records, LearningRecordDTO.class)); return dto; } @Override @Transactional public void addLearningRecord(LearningRecordFormDTO recordDTO) { // 1.获取登录用户 Long userId = UserContext.getUser(); // 2.处理学习记录 boolean finished = false; if(recordDTO.getSectionType() == SectionType.VIDEO) { // 2.1 处理视频 finished = handleVideRecord(userId, recordDTO); }else { // 2.2 处理考试 finished = handleExamRecord(userId, recordDTO); } if(!finished) { return; } // 3.处理课表数据 handleLearningLessonsChanges(recordDTO); } private void handleLearningLessonsChanges(LearningRecordFormDTO recordDTO) { // 1.查询课表 LearningLesson lesson = lessonService.getById(recordDTO.getLessonId()); if (lesson == null){ throw new BizIllegalException("课表不存在,无法更新数据! "); } boolean allLearned = false; // 是否全部小结学完 // 3.如果有新完成的小结,则需要查询课程数据 CourseFullInfoDTO cInfo = courseClient.getCourseInfoById(lesson.getCourseId(),false,false); if(cInfo == null){ throw new BizIllegalException("课程不存在,无法更新数据!"); } // 4. 比较课程是否全部学完;已学习小结 >= 课程总小结 allLearned = lesson.getLearnedSections() + 1 >= cInfo.getSectionNum(); // 5.更新课表 lessonService.lambdaUpdate() .set(lesson.getLearnedSections() == 0,LearningLesson::getStatus,LessonStatus.LEARNING.getValue()) .set(allLearned,LearningLesson::getStatus, LessonStatus.FINISHED.getValue()) .setSql("learned_sections = learned_sections + 1") .eq(LearningLesson::getId,lesson.getId()) .update(); } private boolean handleExamRecord(Long userId, LearningRecordFormDTO recordDTO) { // 1.转换DTO为PO LearningRecord record = BeanUtils.copyBean(recordDTO, LearningRecord.class); // 2.填充数据 record.setUserId(userId); record.setFinished(true); record.setFinishTime(recordDTO.getCommitTime()); // 3.写入数据库 boolean success = save(record); if(!success){ throw new DbException("新增考试记录失败! "); } return true; } private boolean handleVideRecord(Long userId, LearningRecordFormDTO recordDTO) { // 1.查询旧的学习记录 LearningRecord old = queryOldRecord(recordDTO.getLessonId(), recordDTO.getSectionId()); // 2.判断是否存在,存在就更新,不存在就新增 if (old == null){ // 3.不存在则新增 // 3.1 转换PO LearningRecord record = BeanUtils.copyBean(recordDTO, LearningRecord.class); // 3.2 填充数据 record.setUserId(userId); // 3.3 写入数据库 boolean success = save(record); if(!success){ throw new DbException("新增学习记录失败! "); } return false; } // 4. 存在则更新 // 4.1 判断是否是第一次完成(之前学过了但是没学完,finished为false;并且本次的学习进度超过50%) boolean finished = !old.getFinished() && recordDTO.getMoment() * 2 >= recordDTO.getDuration(); if(!finished) { LearningRecord record = new LearningRecord(); record.setLessonId(recordDTO.getLessonId()); record.setSectionId(recordDTO.getSectionId()); record.setMoment(record.getMoment()); record.setId(old.getId()); record.setFinished(old.getFinished()); taskHandler.addLearningRecordTask(record); return false; } // 4.2 更新数据 boolean success = lambdaUpdate() .set(LearningRecord::getMoment, recordDTO.getMoment()) .set(LearningRecord::getFinished, true) .set(LearningRecord::getFinishTime, recordDTO.getCommitTime()) .eq(LearningRecord::getId, old.getId()) .update(); if(!success){ throw new DbException("更新学习记录失败! "); } // 4.3 清除缓存 taskHandler.cleanRecordCache(recordDTO.getLessonId(), recordDTO.getSectionId()); return finished; } // 查询旧的数据 private LearningRecord queryOldRecord(Long lessonId, Long sectionId) { // 1.查询缓存 LearningRecord record = taskHandler.readRecordCache(lessonId, sectionId); // 2.查看是否命中 if (record != null){ return record; } // 3.未命中查询数据库 record = lambdaQuery() .eq(LearningRecord::getLessonId, lessonId) .eq(LearningRecord::getSectionId, sectionId) .one(); // 4.写入缓存 taskHandler.writeRecordCache(record); return record; } ``` ## 点赞功能 ### 前期配置 首先完成一张记录点赞信息的记录表 ```sql -- auto-generated definition create table liked_record ( id bigint auto_increment comment '主键id' primary key, user_id bigint not null comment '用户id', biz_id bigint not null comment '点赞的业务id', biz_type varchar(16) not null comment '点赞的业务类型', create_time datetime default CURRENT_TIMESTAMP not null comment '创建时间', update_time datetime default CURRENT_TIMESTAMP not null on update CURRENT_TIMESTAMP comment '更新时间', constraint idx_biz_user unique (biz_id, user_id) ) comment '点赞记录表'; ``` 然后在项目下新建一个tj-remark(评价服务)模块,具体配置参考tj-learning。使用Mybatis-plus生成相关层级代码。 然后补充一下git代码commit的术语 ```yaml feat:新增功能 fix:修复 bug docs:修改文档 style:调整代码格式,无逻辑改动 refactor:代码重构 perf:性能优化 test:新增 / 修改测试用例 chore:构建、依赖、工具配置调整 ci:修改 CI 流水线配置 build:修改打包构建相关代码 revert:回滚提交 ``` ### 点赞功能 首先判断是点赞操作还是取消点赞操作,然后不能重复点赞或者重复取消点赞。然后统计当前业务id的总点赞数量,通过MQ发送给学习服务,让学习服务下的评论修改自己的点赞数量 ![点赞业务流程图.png](picture/%E7%82%B9%E8%B5%9E%E4%B8%9A%E5%8A%A1%E6%B5%81%E7%A8%8B%E5%9B%BE.png) 点赞业务发送MQ消息 ```java private final RabbitMqHelper mqHelper; @Override public void addLikeRecord(LikeRecordFormDTO recordDTO) { // 1.基于前端的参数,判断是执行点赞还是取消点赞 boolean success = recordDTO.getLiked() ? like(recordDTO) : unlike(recordDTO); // 2.判断是否执行成功,如果失败直接结束 if(!success){ return; } // 3.如果执行成功,统计点赞总数 Integer likeTimes = lambdaQuery() .eq(LikedRecord::getBizId, recordDTO.getBizId()) .count(); // 4.发送MQ通知 mqHelper.send(LIKE_RECORD_EXCHANGE, StringUtils.format(LIKED_TIMES_KEY_TEMPLATE, recordDTO.getBizType()), LikeTimesDTO.of(recordDTO.getBizId(), likeTimes)); } private boolean unlike(LikeRecordFormDTO recordDTO) { return remove(new QueryWrapper().lambda() .eq(LikedRecord::getBizId, recordDTO.getBizId()) .eq(LikedRecord::getUserId, UserContext.getUser())); } private boolean like(LikeRecordFormDTO recordDTO) { Long userId = UserContext.getUser(); // 1.查询点赞记录 Integer count = lambdaQuery() .eq(LikedRecord::getUserId, userId) .eq(LikedRecord::getBizId, recordDTO.getBizId()) .count(); // 2.判断是否存在,如果已经存在,直接结束 if(count > 0){ return false; } // 3.如果不存在,直接新增 LikedRecord r = new LikedRecord(); r.setBizId(recordDTO.getBizId()); r.setBizType(recordDTO.getBizType()); r.setUserId(userId); save(r); return true; } ``` 然后学习服务消费MQ传送过来的点赞数量,修改评论的点赞数量 ![学习服务接受点赞服务消息.png](picture/%E5%AD%A6%E4%B9%A0%E6%9C%8D%E5%8A%A1%E6%8E%A5%E5%8F%97%E7%82%B9%E8%B5%9E%E6%9C%8D%E5%8A%A1%E6%B6%88%E6%81%AF.png) ```java @Component @Slf4j @RequiredArgsConstructor public class LikeTimesChangeListener { private final IInteractionReplyService replyService; @RabbitListener(bindings = @QueueBinding( value = @Queue(name = "qa.liked.times.queue", durable = "true"), exchange = @Exchange(name = LIKE_RECORD_EXCHANGE, type = ExchangeTypes.TOPIC), key = QA_LIKED_TIMES_KEY )) public void listenReplyLikedTimesChange(LikeTimesDTO likeTimesDTO){ log.debug("监听到了回答或者评论点赞数变更的消息:{}, 点赞数:{}", likeTimesDTO.getBizId(), likeTimesDTO.getLikeTimes()); InteractionReply reply = new InteractionReply(); reply.setId(likeTimesDTO.getBizId()); reply.setLikedTimes(likeTimesDTO.getLikeTimes()); replyService.updateById(reply); } } ``` ### 点赞功能优化 用户对某一条回答用set这个数据结构,回答的id作为key,用户id作为value,因为set数据结构具备唯一性所以每个用户只能点赞一次。然后根据业务类型建立key,然后用Zset数据结构存储每个回答id,点赞数作为score; 然后使用定时任务定期通过MQ持久到数据库。 ![点赞功能优化.png](picture/%E7%82%B9%E8%B5%9E%E5%8A%9F%E8%83%BD%E4%BC%98%E5%8C%96.png) ```java Override public void addLikeRecord(LikeRecordFormDTO recordDTO) { // 1.基于前端的参数,判断是执行点赞还是取消点赞 boolean success = recordDTO.getLiked() ? like(recordDTO) : unlike(recordDTO); // 2.判断是否执行成功,如果失败直接结束 if(!success){ return; } // 3.如果执行成功,统计点赞总数 Long likedTimes = redisTemplate.opsForSet() .size(RedisConstants.LIKES_BIZ_KEY_PREFIX + recordDTO.getBizId()); if(likedTimes == null){ return; } // 4.缓存点赞总数到Redis redisTemplate.opsForZSet().add( RedisConstants.LIKES_TIMES_KEY_PREFIX + recordDTO.getBizType(), recordDTO.getBizId().toString(), likedTimes ); } private boolean unlike(LikeRecordFormDTO recordDTO) { // 1.获取用户id Long userId = UserContext.getUser(); // 2.获取Key String key = RedisConstants.LIKES_BIZ_KEY_PREFIX + recordDTO.getBizId(); // 3.执行SADD命令 Long result = redisTemplate.opsForSet().remove(key, userId.toString()); return result != null && result > 0; } private boolean like(LikeRecordFormDTO recordDTO) { // 1.获取用户id Long userId = UserContext.getUser(); // 2.获取Key String key = RedisConstants.LIKES_BIZ_KEY_PREFIX + recordDTO.getBizId(); // 3.执行SADD命令 Long result = redisTemplate.opsForSet().add(key, userId.toString()); return result != null && result > 0; } ``` 定时任务定期持久化 ```java @Component @RequiredArgsConstructor public class LikedTimesCheckTask { private static final List BIZ_TYPES = List.of("QA","NOTE"); private static final int MAX_BIZ_SIZE = 30; private final ILikedRecordService recordService; @Scheduled(fixedDelay = 20000) public void checkLikedTimes(){ for (String bizType : BIZ_TYPES) { recordService.readLikedTimesAndSendMessage(bizType,MAX_BIZ_SIZE); } } } @Override public void readLikedTimesAndSendMessage(String bizType, int maxBizSize) { log.info("定时任务开始持久化数据了"); // 1.读取并移除Redis里面的缓存总数 String key = RedisConstants.LIKES_TIMES_KEY_PREFIX + bizType; Set> tuples = redisTemplate.opsForZSet().popMin(key, maxBizSize); if(CollUtils.isEmpty(tuples)) { return; } // 2.数据转换 List list = new ArrayList<>(tuples.size()); for (ZSetOperations.TypedTuple tuple : tuples) { String bizId = tuple.getValue(); Double likeTimes = tuple.getScore(); if(bizId == null || likeTimes == null){ continue; } list.add(LikeTimesDTO.of(Long.valueOf(bizId), likeTimes.intValue())); // 3.发送MQ消息 mqHelper.send( LIKE_RECORD_EXCHANGE, StringUtils.format(LIKED_TIMES_KEY_TEMPLATE, bizType), list ); } } ``` ### 使用批处理pipeline打包合并请求数量 避免了循环遍历bizIds重复多次建立RedisIO请求,直接一次IO将多个bizIds打包一次请求 ```java @Override public Set isBizLiked(List bizIds) { Long userId = UserContext.getUser(); // 2.查询点赞状态 使用批处理打包减少IO次数 List objects = redisTemplate.executePipelined(new RedisCallback() { @Override public Object doInRedis(RedisConnection connection) throws DataAccessException { StringRedisConnection src = (StringRedisConnection) connection; for (Long bizId : bizIds) { String key = RedisConstants.LIKES_BIZ_KEY_PREFIX + bizId; src.sIsMember(key, userId.toString()); } return null; } }); // 3. 返回结果 Set set = new HashSet<>(); for (int i = 0; i < objects.size(); i++) { Boolean o = (Boolean) objects.get(i); if(o){ set.add(bizIds.get(i)); } } return set; } ``` ## 签到积分功能 Redis 的 Bitmap 底层基于 String 实现,1 个 bit 位标记一天签到状态:1‑已签到,0‑未签到。 Key 设计:sign:{userId}:{yyyyMM}(按用户 + 年月隔离); offset:当月日期‑1(1 号对应 offset=0); 存储优势:一个用户一个月最多 31bit(不足 4 字节),内存占用相比 MySQL 节省 99% 以上;Redis 位运算全部内存执行,速度极快。 ![用户签到.png](picture/%E7%94%A8%E6%88%B7%E7%AD%BE%E5%88%B0.png) ### 使用BitMap做签到功能 首先在linux命令行里面尝试一下 ```yaml docker exec -it redis redis-cli # 进入docker容器 auth 123321 # 输入用户密码 setbit bm 0 1 # 将命名为bm的BitMap的第一位设置为1 bitfield bm get u3 0 # 从第零位开始将前三位以无符号的形式返回,例如111,返回7,7的二进制就是111,代表前三天都签到了 ``` 用户签到和计算连续签到天数 ```java @Service @RequiredArgsConstructor @Slf4j public class SignRecordServiceImpl implements ISignRecordService { private final StringRedisTemplate redisTemplate; private final RabbitMqHelper mqHelper; // 用户签到 @Override public SignResultVO addSignRecords() { // 1.签到 Long userId = UserContext.getUser(); LocalDate now = LocalDate.now(); // 1.1 组装key String key = RedisContants.SIGN_RECORD_KEY_PREFIX + userId + now.format(DateUtils.SIGN_DATE_SUFFIX_FORMATTER); // 1.2 计算offset int offset = now.getDayOfMonth() - 1; // 1.3 保存签到信息 Boolean exits = redisTemplate.opsForValue().setBit(key, offset, true); if(BooleanUtils.isTrue(exits)) { throw new BizIllegalException("不允许一天内签到多次! "); } // 2.计算连续签到天数 int signDays = countSignDays(key, now.getDayOfMonth()); // 3.计算签到得分 int rewardPoint = 0; if(signDays >= 7 && signDays < 14) { rewardPoint = 10; } if(signDays >= 14 && signDays < 28) { rewardPoint = 20; } if(signDays >= 28) { rewardPoint = 40; } // 保存积分 mqHelper.send(MqConstants.Exchange.LEARNING_EXCHANGE, MqConstants.Key.SIGN_IN, SignInMessage.of(userId, rewardPoint + 1)); SignResultVO vo = new SignResultVO(); vo.setSignDays(signDays); vo.setSignPoints(rewardPoint); return vo; } // 得到用户签到记录,最后用数组返回,签到为1,没签到为0 @Override public List getSignRecords() { LocalDate now = LocalDate.now(); String key = RedisContants.SIGN_RECORD_KEY_PREFIX + UserContext.getUser() + now.format(DateUtils.SIGN_DATE_SUFFIX_FORMATTER); List result = redisTemplate.opsForValue().bitField(key, BitFieldSubCommands.create() .get(BitFieldSubCommands.BitFieldType.unsigned(now.getDayOfMonth())).valueAt(0)); if(CollUtils.isEmpty(result)){ throw new BizIllegalException("用户签到数据不存在! "); } // 以内容0,1的list返回 List SignDaysResult = new ArrayList<>(now.getDayOfMonth()); int num = result.get(0).intValue(); for(int i = 1 ; i <= now.getDayOfMonth(); i++) { if((num & 1) == 1){ SignDaysResult.add(1); } else { SignDaysResult.add(0); } num >>>= 1; } Collections.reverse(SignDaysResult); return SignDaysResult; } // 计算连续签到天数 private int countSignDays(String key, int len) { // 1.获取本月从第一天开始到今天为止的全部签到记录 List result = redisTemplate.opsForValue().bitField(key, BitFieldSubCommands.create() .get(BitFieldSubCommands.BitFieldType.unsigned(len)).valueAt(0)); if(CollUtils.isEmpty(result)){ return 0; } int num = result.get(0).intValue(); log.info("全部签到记录:{}", num); // 2.定义一个计数器 int count = 0; // 3.循环,最后一个bit与1做与运算,判断是否为0,为0终止,为1继续 while ((num & 1) == 1) { // 4.计数器+1 count++; // 5.把数字右移一位,最后一位被舍弃,倒数第二位成为最后一位 num >>>= 1; } return count; } } ``` 用MQ监听签到以后保存签到积分 ```java @RabbitListener(bindings = @QueueBinding( value = @Queue(name = "sign.points.queue", durable = "true"), exchange = @Exchange(name = MqConstants.Exchange.LEARNING_EXCHANGE, type = ExchangeTypes.TOPIC), key = MqConstants.Key.SIGN_IN )) public void listenSignInMessage(SignInMessage message){ recordService.addPointsRecord(message.getUserId(), message.getPoints(), PointsRecordType.SIGN); } ``` 最后新增签到记录 ```java @Override public void addPointsRecord(Long userId, int points, PointsRecordType type) { // 1.判断当前方式有没有积分上限 int maxPoints = type.getMaxPoints(); int realPoints = points; if(maxPoints > 0) { // 2.有积分上限 LocalDateTime now = LocalDateTime.now(); LocalDateTime begin = DateUtils.getDayStartTime(now); LocalDateTime end = DateUtils.getDayEndTime(now); // 2.1 查询今天这个积分类型的已得到积分 int currentPoints = queryUserPointByTypeAndDate(userId, type, begin, end); // 2.2 查看是否超过积分上限,超过直接结束 if(currentPoints >= maxPoints) { return; } // 2.3 没有超过,保存积分记录 if(currentPoints + points > maxPoints){ realPoints = maxPoints - currentPoints; } } // 没有积分上限直接存积分 && 有积分上限但没有超过,保存积分记录 PointsRecord p = new PointsRecord(); p.setPoints(realPoints); p.setType(type); p.setUserId(userId); save(p); } private int queryUserPointByTypeAndDate( Long userId, PointsRecordType type, LocalDateTime begin, LocalDateTime end) { // 1.查询条件 QueryWrapper wrapper = new QueryWrapper<>(); wrapper.lambda() .eq(PointsRecord::getUserId, userId) .eq(type !=null, PointsRecord::getType, type) .between(begin != null && end != null, PointsRecord::getCreateTime, begin, end); // 2.调用mapper Integer points = getBaseMapper().queryUserPointsByTypeAndDate(wrapper); return points == null ? 0 : points; } ``` ## 排行榜功能 在用户新增积分的时候将每次同步的积分同步到Redis里面;采用Zset数据结构,使用当前赛季为key,用户的id作为member,用户当前赛季累计所有的为score。然后下个月新赛季根据赛季进行分表将数据同步到数据库里面 ![排行榜功能.png](picture/%E6%8E%92%E8%A1%8C%E6%A6%9C%E5%8A%9F%E8%83%BD.png) ### 排行榜添加数据 改造之前新增积分代码,每次新增积分以后将总的积分同步到Redis做积分排行榜使用 ```java @Override public void addPointsRecord(Long userId, int points, PointsRecordType type) { ...... // 4.累计积分到Redis的SortSet里面 String key = RedisContants.POINTS_BOAED_KEY_PREFIX + now.format(DateUtils.POINTS_BOARD_SUFFIX_FORMATTER); redisTemplate.opsForZSet().incrementScore(key, userId.toString(), realPoints); } ``` ### 查询排行榜 本质我们的操作是利用Zset数据结构的Score做一个排序查询,倒序查询,分数越高越靠前。同时要注意处理分页数据from = (pageNo - 1) * pageSize;end = from + pageSize - 1 ```java private List queryCurrentBoardList(String key, Integer pageNo, Integer pageSize) { // 1.计算分页 int from = (pageNo - 1) * pageSize; // 2.查询 Set> tuples = redisTemplate. opsForZSet().reverseRangeWithScores(key, from, from + pageSize - 1); if(CollUtils.isEmpty(tuples)) { return CollUtils.emptyList(); } // 3.封装 int rank = from + 1; List list = new ArrayList<>(tuples.size()); for (ZSetOperations.TypedTuple tuple : tuples) { String userId = tuple.getValue(); Double point = tuple.getScore(); if(userId == null || point == null){ continue; } PointsBoard p = new PointsBoard(); p.setUserId(Long.valueOf(userId)); p.setPoints(point.intValue()); p.setRank(rank++); list.add(p); } return list; } ``` ### 分表策略 创建一个定时任务在每个月的月初凌晨三点开始查询上个月的赛季信息,根据上个月赛季信息创建分表存储上个月赛季的积分数据 ```java @Component @RequiredArgsConstructor @Slf4j public class PointsBoardPersistentHandler { private final IPointsBoardSeasonService seasonService; private final IPointsBoardService pointsBoardService; @Scheduled(cron = "0 0 3 1 * ?") // 每月1号凌晨3:00执行 public void createPointsBoardTableOfLastSeason() { // 1.获取上月时间 LocalDateTime time = LocalDateTime.now().minusMonths(1); // 2.查询赛季id Integer season = seasonService.querySeasonByTime(time); if(season == null) { // 赛季不存在 log.info("赛季不存在"); return; } // 3.创建表 pointsBoardService.createPointsBoardTableBySeason(season); } } ``` ### 分布式任务调度 XXL-Job 将业务代码修改 ```java @XxlJob("createTableJob") public void createPointsBoardTableOfLastSeason() { ... } ``` 然后在创建任务管理,之后在操作里面点击执行一次测试效果 ![创建XXLJob任务管理.png](picture/%E5%88%9B%E5%BB%BAXXLJob%E4%BB%BB%E5%8A%A1%E7%AE%A1%E7%90%86.png) 然后完成所有的业务逻辑,首先梳理出流程图。 ![数据持久化.png](picture/%E6%95%B0%E6%8D%AE%E6%8C%81%E4%B9%85%E5%8C%96.png) - 1.分表 首先是创建出分表。我们采用了Mybatis-Plus的动态表名,创建表名的时候用ThreadLocal来存储,然后借助于DynamicTableNameInnerInterceptor插件将ThreadLocal的信息加载。 创建存储动态表名的ThreadLocal ```java public class TableInfoContext { private static final ThreadLocal TL = new ThreadLocal<>(); public static void setInfo(String info){ TL.set(info); } public static String getInfo() { return TL.get(); } public static void remove() { TL.remove(); } } ``` 创建动态表名的时候将表名加里面 ```java Integer season = seasonService.querySeasonByTime(time); TableInfoContext.setInfo("points_board_" + season); ``` 再使用DynamicTableNameInnerInterceptor插件加载 ```java @Configuration public class MybatisConfiguration { @Bean public DynamicTableNameInnerInterceptor dynamicTableNameInnerInterceptor() { Map map = new HashMap<>(1); map.put("points_board", new TableNameHandler() { @Override public String dynamicTableName(String sql, String tableName) { return TableInfoContext.getInfo(); } }); return new DynamicTableNameInnerInterceptor(map); } } ``` 最后在总配置里面按照DynamicTableNameInnerInterceptor的有无来加载配置类 ```java @Bean @ConditionalOnMissingBean public MybatisPlusInterceptor mybatisPlusInterceptor(@Autowired(required = false)DynamicTableNameInnerInterceptor innerInterceptor) { ..... if(interceptor != null){ interceptor.addInnerInterceptor(innerInterceptor); } return interceptor; } ``` - 2.分片读取Redis 考虑到多实例部署,我们采用了XXL-job的分片策略。实现多个实例分片读取Redis信息 ```java int index = XxlJobHelper.getShardIndex(); // 第几个实例 int total = XxlJobHelper.getShardTotal(); // 一共几个实例 // 分片操作 int pageNo = index + 1; int pageSize = 1000; while (true) { List boardList = pointsBoardService.queryCurrentBoardList(key, pageNo, pageSize); if(CollUtils.isEmpty(boardList)){ break; } boardList.forEach(b -> { b.setId(b.getRank().longValue()); b.setRank(null); b.setSeason(null); }); // 4.持久化 pointsBoardService.saveBatch(boardList); // 5.翻页继续查询 // pageNo++; pageNo += total; } ``` - 3.清除Redis数据 因为是一个大key,采用unlink删除性能更好,是异步删除,不阻塞性能。业务里清理缓存、定时删数据、分片批量清理 key,优先用 UNLINK。 ```java @XxlJob("clearPointsBoardFromRedis") public void clearPointsBoardFromRedis(){ /*// 1.获取上月时间 LocalDateTime time = LocalDateTime.now().minusMonths(1);*/ // 测试时间戳 2026-07-01 00:00:00 LocalDateTime time = LocalDateTime.of(2026, 7, 1, 0, 0); String key = RedisContants.POINTS_BOAED_KEY_PREFIX + time.format(DateUtils.POINTS_BOARD_SUFFIX_FORMATTER); redisTemplate.unlink(key); // 因为是一个大key,采用unlink删除性能更好 } ``` - 4.XXL—job配置 要注意三个方法的任务管理的执行链路,而且记得savePointBoardToDB的执行策略是分片。 总代码 ```java @Component @RequiredArgsConstructor @Slf4j public class PointsBoardPersistentHandler { private final IPointsBoardSeasonService seasonService; private final IPointsBoardService pointsBoardService; private final StringRedisTemplate redisTemplate; // @Scheduled(cron = "0 0 3 1 * ?") // 每月1号凌晨3:00执行 @XxlJob("createTableJob") public void createPointsBoardTableOfLastSeason() { /*// 1.获取上月时间 LocalDateTime time = LocalDateTime.now().minusMonths(1);*/ // 测试时间戳 2026-07-01 00:00:00 LocalDateTime time = LocalDateTime.of(2026, 7, 1, 0, 0); // 2.查询赛季id Integer season = seasonService.querySeasonByTime(time); if(season == null) { // 赛季不存在 log.info("赛季不存在"); return; } // 3.创建表 pointsBoardService.createPointsBoardTableBySeason(season); } @XxlJob("savePointBoardToDB") public void savePointBoardToDB(){ /*// 1.获取上月时间 LocalDateTime time = LocalDateTime.now().minusMonths(1);*/ // 测试时间戳 2026-07-01 00:00:00 LocalDateTime time = LocalDateTime.of(2026, 7, 1, 0, 0); // 2.计算动态表名 Integer season = seasonService.querySeasonByTime(time); TableInfoContext.setInfo("points_board_" + season); // 3.查询Redis String key = RedisContants.POINTS_BOAED_KEY_PREFIX + time.format(DateUtils.POINTS_BOARD_SUFFIX_FORMATTER); /*// 单机操作 int pageNo = 1; int pageSize = 1000;*/ int index = XxlJobHelper.getShardIndex(); // 第几个实例 int total = XxlJobHelper.getShardTotal(); // 一共几个实例 // 分片操作 int pageNo = index + 1; int pageSize = 1000; while (true) { List boardList = pointsBoardService.queryCurrentBoardList(key, pageNo, pageSize); if(CollUtils.isEmpty(boardList)){ break; } boardList.forEach(b -> { b.setId(b.getRank().longValue()); b.setRank(null); b.setSeason(null); }); // 4.持久化 pointsBoardService.saveBatch(boardList); // 5.翻页继续查询 // pageNo++; pageNo += total; } TableInfoContext.remove(); } @XxlJob("clearPointsBoardFromRedis") public void clearPointsBoardFromRedis(){ /*// 1.获取上月时间 LocalDateTime time = LocalDateTime.now().minusMonths(1);*/ // 测试时间戳 2026-07-01 00:00:00 LocalDateTime time = LocalDateTime.of(2026, 7, 1, 0, 0); String key = RedisContants.POINTS_BOAED_KEY_PREFIX + time.format(DateUtils.POINTS_BOARD_SUFFIX_FORMATTER); redisTemplate.unlink(key); // 因为是一个大key,采用unlink删除性能更好 } } ``` ## 优惠券促销 ### 优惠券 目前涉及的业务有优惠券和优惠券作用范围的业务,他们是多对多的关系,一个优惠券可以作用多个分类,一个分类也能有多个优惠券服务。涉及到优惠券表coupon,和优惠券作用范围信息coupon_scope ![优惠券服务页面.png](picture/%E4%BC%98%E6%83%A0%E5%88%B8%E6%9C%8D%E5%8A%A1%E9%A1%B5%E9%9D%A2.png) 新增优惠券操作 ```java @Override @Transactional public void saveCoupon(CouponFormDTO dto) { // 1.保存优惠券 Coupon coupon = BeanUtils.copyBean(dto, Coupon.class); save(coupon); if(!dto.getSpecific()){ // 没有限定范围就不用操作优惠券作用范围信息coupon_scope return; } // 2.保存限定范围 List scopes = dto.getScopes(); if(CollUtils.isEmpty(scopes)) { throw new BadRequestException("限定范围不能为空"); } Long couponId = coupon.getId(); List list = scopes.stream() .map(bizId -> new CouponScope().setBizId(bizId).setCouponId(couponId)) .collect(Collectors.toList()); // 3.保存 scopeService.saveBatch(list); } ``` ### 生成兑换码 因为发放优惠券的数量有可能很大,上千张。所以我们采用线程池提高并发量,异步处理这个任务。 线程池代码 ```java @Slf4j @Configuration public class PromotionConfig { @Bean public Executor generateExchangeCodeExecutor(){ ThreadPoolTaskExecutor executor = new ThreadPoolTaskExecutor(); executor.setCorePoolSize(2); executor.setMaxPoolSize(5); executor.setQueueCapacity(200); executor.setThreadNamePrefix("exchange-code-handler-"); executor.setRejectedExecutionHandler(new ThreadPoolExecutor.CallerRunsPolicy()); // 同步执行 executor.initialize(); return executor; } } ``` 然后使用Redis的自增方法生成兑换码id,借助兑换码算法生成兑换码 ```java @Override @Async("generateExchangeCodeExecutor") public void asyncGenerateCode(Coupon coupon) { Integer totalNum = coupon.getTotalNum(); // 1.获取Redis自增长序列号 Long result = serialOps.increment(totalNum); if(result == null){ return; } int maxSerialNum = result.intValue(); List list = new ArrayList<>(totalNum); for (int serialNum = maxSerialNum - totalNum + 1; serialNum <= maxSerialNum; serialNum++) { // 2.生成兑换码 String code = CodeUtil.generateCode(serialNum, coupon.getId()); // 3.保存数据库 ExchangeCode e = new ExchangeCode(); e.setId(serialNum); e.setCode(code); e.setExchangeTargetId(coupon.getId()); e.setExpiredTime(coupon.getIssueEndTime()); list.add(e); } saveBatch(list); // 批处理提高性能 } ``` ## 领取优惠券 ### 手动点击领取优惠券 完成手动领取优惠券的业务代码,将优惠券的发放数量进行变化,同时新增用户领取优惠券的记录 ![领取优惠券.png](picture/%E9%A2%86%E5%8F%96%E4%BC%98%E6%83%A0%E5%88%B8.png) ```java @Service @RequiredArgsConstructor public class UserCouponServiceImpl extends ServiceImpl implements IUserCouponService { private final CouponMapper couponMapper; @Override @Transactional public void receiveCoupon(Long couponId) { // 校验能否发放 Coupon coupon = couponMapper.selectById(couponId); if(coupon == null) { throw new BadRequestException("优惠券不存在"); } LocalDateTime now = LocalDateTime.now(); if(now.isBefore(coupon.getIssueBeginTime())) { throw new BadRequestException("优惠券发放还未开始"); } if(now.isAfter(coupon.getIssueEndTime())) { throw new BadRequestException("优惠券发放已经结束"); } if(coupon.getIssueNum() >= coupon.getTotalNum()) { throw new BadRequestException("优惠券库存不足"); } // 限制领取数量 Long userId = UserContext.getUser(); Integer count = lambdaQuery() .eq(UserCoupon::getUserId, userId) .eq(UserCoupon::getCouponId, coupon.getId()) .count(); if(count != null && count >= coupon.getUserLimit()) { throw new BadRequestException("超过每个用户的限制领取数量"); } // 领取优惠券后优惠券数量 +1 couponMapper.incrIssueNum(couponId); // 保存用户券信息 saveUserCoupon(coupon, userId); } private void saveUserCoupon(Coupon coupon, Long userId) { UserCoupon uc = new UserCoupon(); uc.setCouponId(coupon.getId()); uc.setUserId(userId); // 有效期 LocalDateTime termBeginTime = coupon.getTermBeginTime(); LocalDateTime termEndTime = coupon.getTermEndTime(); if(termBeginTime == null){ termBeginTime = LocalDateTime.now(); termEndTime = termBeginTime.plusDays(coupon.getTermDays()); } uc.setTermBeginTime(termBeginTime); uc.setTermEndTime(termEndTime); save(uc); } } ``` ### 兑换码兑换领取优惠券 我们使用BitMap来存储每个兑换码的兑换状态(1代表已兑换,0代表未兑换),首先直接set这个兑换码的key,看看返回值old值是什么,判断这个key是否被兑换。然后后续相关DB的操作使用事务包裹,然后使用try-catch包裹。 如果出现异常事务机制会让DB回滚,try-catch并且抛出异常会让事务生效,再回滚BitMap的数据。总结就是BitMap使用直接一步完成,避免先getbit再业务最后setbit导致多线程情况下线程冲突导致卡死,在自己完成一个Redis的手动业务回滚 ![兑换码兑换领取优惠券.png](picture/%E5%85%91%E6%8D%A2%E7%A0%81%E5%85%91%E6%8D%A2%E9%A2%86%E5%8F%96%E4%BC%98%E6%83%A0%E5%88%B8.png) SETBIT xxxKey offset 1 :将offset修改为1,并返回old值实现 ```java // 检验该兑换码是否被兑换 @Override public boolean updateExchangeMark(long serialNum, boolean mark) { // SETBIT xxxKey offset 1 :将offset修改为1,并返回old值 Boolean isUsed = redisTemplate.opsForValue().setBit(COUPON_CODE_MAP_KEY, serialNum - 1, mark); return isUsed != null && isUsed; } ``` ```java private void checkAndCreateUser(Coupon coupon, Long userId) { // 限制领取数量 Integer count = lambdaQuery() .eq(UserCoupon::getUserId, userId) .eq(UserCoupon::getCouponId, coupon.getId()) .count(); if(count != null && count >= coupon.getUserLimit()) { throw new BadRequestException("超过每个用户的限制领取数量"); } // 领取优惠券后优惠券数量 +1 couponMapper.incrIssueNum(coupon.getId()); // 保存用户券信息 saveUserCoupon(coupon, userId); } @Override @Transactional public void exchangeCoupon(String code) { // 校验并解析兑换码 long serialNum = CodeUtil.parseCode(code); // 校验是否已经兑换 boolean exchanged = codeService.updateExchangeMark(serialNum, true); if(exchanged) { throw new BizIllegalException("兑换码已经被兑换过了"); } try { // 查询出兑换码 ExchangeCode exchangeCode = codeService.getById(serialNum); if(exchangeCode == null){ throw new BizIllegalException("兑换码不存在!"); } // 是否过期 if(LocalDateTime.now().isAfter(exchangeCode.getExpiredTime())) { throw new BizIllegalException("兑换码已经过期!"); } Long userId = UserContext.getUser(); Coupon coupon = couponMapper.selectById(exchangeCode.getExchangeTargetId()); checkAndCreateUser(coupon, UserContext.getUser()); // 更新兑换码状态 codeService.lambdaUpdate() .set(ExchangeCode::getUserId, userId) .set(ExchangeCode::getStatus, ExchangeCodeStatus.USED) .eq(ExchangeCode::getId, exchangeCode.getId()) .update(); } catch (Exception e) { // 如果发生错误兑换码状态回滚 codeService.updateExchangeMark(serialNum, false); throw e; } } ``` ## 超卖问题 - 什么是超卖: 商品库存只剩 1 件,但并发下单同时抢到库存,最终卖出超出 1 件,实际发货不足,就是超卖。 - 产生原因(核心): - 多线程 / 大量并发请求同时查询库存: - A、B 同时查库,都查到库存 = 1; - 两者都判断:库存 > 0,可以下单; - 先后修改卖出量,卖出量变成 2; - 实际卖了 2 件,库存数据错乱,多出订单。 - 典型场景 秒杀、优惠券兑换、限时抢购 ### 使用Jmeter压测 ![超卖问题压测数据.png](picture/%E8%B6%85%E5%8D%96%E9%97%AE%E9%A2%98%E5%8E%8B%E6%B5%8B%E6%95%B0%E6%8D%AE.png) ### 使用乐观锁 乐观锁流程图 ![超卖问题乐观锁.png](picture/%E8%B6%85%E5%8D%96%E9%97%AE%E9%A2%98%E4%B9%90%E8%A7%82%E9%94%81.png) 如果发生了并发不安全的情况要修改表的时候比较当前的购买量和优惠券总量的关系,是不是当前购买量小于优惠券总量(issue_num < total_num) ```java private void checkAndCreateUser(Coupon coupon, Long userId) { // 限制领取数量 Integer count = lambdaQuery() .eq(UserCoupon::getUserId, userId) .eq(UserCoupon::getCouponId, coupon.getId()) .count(); if(count != null && count >= coupon.getUserLimit()) { throw new BadRequestException("超过每个用户的限制领取数量"); } // 领取优惠券后优惠券数量 + 1 int r = couponMapper.incrIssueNum(coupon.getId()); if (r == 0) { throw new BizIllegalException("优惠券库存不足"); } // 保存用户券信息 saveUserCoupon(coupon, userId); } ``` 数据层具体改表数据 ```java public interface CouponMapper extends BaseMapper { @Update("UPDATE coupon SET issue_num = issue_num + 1 WHERE id = #{couponId} AND issue_num < total_num") int incrIssueNum(@Param("couponId") Long couponId); } ``` ### 锁失效和锁边界问题 首先设定一个场景,一个用户用自己的单个id高并发去领取优惠券,可能出现并发问题,同时抢到多张券不满足一人一券。这里是先查再新增,只能用悲观锁。将锁的粒度小一点,只锁userId;这样不同用户不会被悲观锁锁住,同一个用户会锁住。 并且要先提交事务再释放锁资源,等事务落表以后再释放,新的线程获取锁才能根据上次落表的数据做逻辑判断。 ![锁失效和锁边界问题.png](picture/%E9%94%81%E5%A4%B1%E6%95%88%E5%92%8C%E9%94%81%E8%BE%B9%E7%95%8C%E9%97%AE%E9%A2%98.png) ```java @Override public void receiveCoupon(Long couponId) { // 校验能否发放 Coupon coupon = couponMapper.selectById(couponId); if(coupon == null) { throw new BadRequestException("优惠券不存在"); } LocalDateTime now = LocalDateTime.now(); if(now.isBefore(coupon.getIssueBeginTime())) { throw new BadRequestException("优惠券发放还未开始"); } if(now.isAfter(coupon.getIssueEndTime())) { throw new BadRequestException("优惠券发放已经结束"); } if(coupon.getIssueNum() >= coupon.getTotalNum()) { throw new BadRequestException("优惠券库存不足"); } Long userId = UserContext.getUser(); // 校验并生成用户券 synchronized (userId.toString().intern()){ // 将将这个String放入常量池,保证每次同一个对象 checkAndCreateUser(coupon, userId); } } @Transactional public void checkAndCreateUser(Coupon coupon, Long userId) { // 限制领取数量 Integer count = lambdaQuery() .eq(UserCoupon::getUserId, userId) .eq(UserCoupon::getCouponId, coupon.getId()) .count(); if(count != null && count >= coupon.getUserLimit()) { throw new BadRequestException("超过每个用户的限制领取数量"); } // 领取优惠券后优惠券数量 + 1 int r = couponMapper.incrIssueNum(coupon.getId()); if (r == 0) { throw new BizIllegalException("优惠券库存不足"); } // 保存用户券信息 saveUserCoupon(coupon, userId); } ``` ### 事务失效场景 - 1.方法不是public修饰 - 2.非事务方法调用本类中的事务方法 - 3.异常被捕获但是没有抛出 - 4.回滚异常类型不匹配,SpringBoot默认只在 RuntimeException、Error 才回滚。指定 @Transactional(rollbackFor = Exception.class) - 5.被调用方法的事务传播类型跟调用方法类型的事务传播类型不一样 - 6.没有被Spring管理,比如没加@Service 目前存在第二个隐患,所以要进行修改。首先引入aspectjweaver然后生成当前类的代理对象调用事务方法 ```yaml org.aspectj aspectjweaver ``` ```java // 校验并生成用户券 synchronized (userId.toString().intern()){ // 将将这个String放入常量池,保证每次同一个对象 // 当前类的 AOP 代理对象 IUserCouponService userCouponService = (IUserCouponService) AopContext.currentProxy(); userCouponService.checkAndCreateUser(coupon, userId); } ``` ## 分布式锁 ### 分布式锁原理 一般原理是基于SETNX 来实现 ```yaml # 设置setnx SETEX key thread # 给key设置过期时间 EXPIRE 10 # 删除key DEL key # 一步命令结合设置key和设置过期时间 SET lock thread NX EX 20 ``` ![分布式锁原理.png](picture/%E5%88%86%E5%B8%83%E5%BC%8F%E9%94%81%E5%8E%9F%E7%90%86.png) 具体实现 ### 个人实现简易分布式锁 ```java @AllArgsConstructor public class RedisLock { private String key; private StringRedisTemplate redisTemplate; public boolean tryLock(long leaseTime, TimeUnit unit) { String name = Thread.currentThread().getName(); Boolean success = redisTemplate.opsForValue().setIfAbsent(key, name, leaseTime, unit); return BooleanUtil.isTrue(success); } public void unlock(){ redisTemplate.delete(key); } } ``` 业务具体使用 ```java // 自定义分布式锁 String key = "lock:coupon:uid:" + userId; RedisLock lock = new RedisLock(key, redisTemplate); try { boolean isLock = lock.tryLock(5 , TimeUnit.SECONDS); // 获取锁失败 if(!isLock) { throw new BizIllegalException("请求太频繁"); } IUserCouponService userCouponService = (IUserCouponService) AopContext.currentProxy(); userCouponService.checkAndCreateUser(coupon, userId); }finally { lock.unlock(); } ``` ### Redisson ***Redisson初步使用*** 单机 synchronized / Lock 只在一台服务器生效;多台服务集群下会并发超卖、重复请求,Redisson 提供分布式锁,让所有机器共用同一把锁。封装 Redis 的高级工具库,主打分布式锁,一站式解决集群环境并发控制问题。 ```java // 使用分布式锁Redisson String key = "lock:coupon:uid:" + userId; RLock lock = redissonClient.getLock(key); boolean isLock = lock.tryLock(); if(!isLock) { throw new BizIllegalException("请求太频繁"); } try { IUserCouponService userCouponService = (IUserCouponService) AopContext.currentProxy(); userCouponService.checkAndCreateUser(coupon, userId); } finally { lock.unlock(); } ``` ***AOP注解简化操作*** 后续可以外接AOP环绕增强简化这个操作,首先需要定义@MyLock注解来定义分布式锁需要的各个参数 ```java @Retention(RetentionPolicy.RUNTIME) // 运行时生效 @Target(ElementType.METHOD) // 作用在方法上 public @interface MyLock { String name(); long waitTime() default 1; long leaseTime() default -1; TimeUnit unit() default TimeUnit.SECONDS; } ``` 然后设置AOP环绕增强,ProceedingJoinPoint pjp代表切面原本的业务,在此基础做环绕业务,要注意实现Ordered接口,优先级调高,保证先执行锁,再执行事务,事务结束了再释放锁。 所有注解默认优先级都是Ordered.LOWEST_PRECEDENCE(Integer.MAX_VALUE),所以我们小于Integer.MAX_VALUE优先级就能调高。 ```java @Aspect @RequiredArgsConstructor public class MyLockAspect implements Ordered { private final RedissonClient redissonClient; @Around("@annotation(myLock)") public Object tryLock(ProceedingJoinPoint pjp, MyLock myLock) throws Throwable { // 创建锁对象 RLock lock = redissonClient.getLock(myLock.name()); boolean isLock = lock.tryLock(myLock.waitTime(), myLock.leaseTime(), myLock.unit()); if(!isLock) { throw new BizIllegalException("请求太频繁!"); } try { return pjp.proceed(); } finally { lock.unlock(); } } @Override public int getOrder() { return 0; } } ``` ***结合工厂模式 && 策略模式*** 首先是业务方法调用我们的分布式锁 ```java @MyLock(name = "lock:coupon:uid:") @Transactional @Override public void checkAndCreateUser(Coupon coupon, Long userId) { ... } ``` 将分布锁的不同类型用枚举收集 ```java public enum MyLockType { RE_ENTRANT_LOCK, // 可重入锁 FAIR_LOCK, // 公平锁 READ_LOCK, // 读锁 WRITE_LOCK; // 写锁 } ``` 然后使用工厂模式,核心是使用Map数据结构,key是锁的枚举类型,value是具体锁的实现api。最后使用getLock方法通过key拿方法的方式返回给调用方 ```java @Component public class MyLockFactory { private final Map> lockHandlers; public MyLockFactory(RedissonClient redissonClient) { this.lockHandlers = new EnumMap<>(MyLockType.class); this.lockHandlers.put(MyLockType.RE_ENTRANT_LOCK, redissonClient::getLock); this.lockHandlers.put(MyLockType.FAIR_LOCK, redissonClient::getFairLock); this.lockHandlers.put(MyLockType.READ_LOCK, name -> redissonClient.getReadWriteLock(name).readLock()); this.lockHandlers.put(MyLockType.WRITE_LOCK, name -> redissonClient.getReadWriteLock(name).writeLock()); } public RLock getLock(MyLockType lockType, String name) { return lockHandlers.get(lockType).apply(name); } } ``` 分布锁有几种获取锁失败的策略,这是使用枚举结合策略模式将其实现 ```java public enum MyLockStrategy { SKIP_FAST() { @Override public boolean tryLock(RLock lock, MyLock prop) throws InterruptedException { return lock.tryLock(0,prop.leaseTime(), prop.unit()); } }, // 快速结束 FAIL_FAST() { @Override public boolean tryLock(RLock lock, MyLock prop) throws InterruptedException { boolean isLock = lock.tryLock(0, prop.leaseTime(), prop.unit()); if(!isLock) { throw new BizIllegalException("请求太频繁!"); } return true; } }, // 快速失败 KEEP_TRYING() { @Override public boolean tryLock(RLock lock, MyLock prop) throws InterruptedException { lock.lock(prop.leaseTime(), prop.unit()); return true; } }, // 无限重试 SKIP_AFTER_RETRY_TIMEOUT() { @Override public boolean tryLock(RLock lock, MyLock prop) throws InterruptedException { return lock.tryLock(prop.waitTime(), prop.leaseTime(), prop.unit()); } }, // 重试超时后结束 FAIL_AFTER_RETRY_TIMEOUT() { @Override public boolean tryLock(RLock lock, MyLock prop) throws InterruptedException { boolean isLock = lock.tryLock(prop.waitTime(), prop.leaseTime(), prop.unit()); if(!isLock) { throw new BizIllegalException("请求太频繁!"); } return true; } }; // 重试超时后抛异常 public abstract boolean tryLock(RLock lock, MyLock prop) throws InterruptedException; } ``` 最后是@MyLock和其AOP切面 ```java @Retention(RetentionPolicy.RUNTIME) // 运行时生效 @Target(ElementType.METHOD) // 作用在方法上 public @interface MyLock { String name(); long waitTime() default 1; long leaseTime() default -1; TimeUnit unit() default TimeUnit.SECONDS; MyLockType lockType() default MyLockType.RE_ENTRANT_LOCK; MyLockStrategy lockStrategy() default MyLockStrategy.FAIL_AFTER_RETRY_TIMEOUT; } ``` ```java @Component @Aspect @RequiredArgsConstructor public class MyLockAspect implements Ordered { // private final RedissonClient redissonClient; private final MyLockFactory lockFactory; @Around("@annotation(myLock)") public Object tryLock(ProceedingJoinPoint pjp, MyLock myLock) throws Throwable { // 创建锁对象-普通实现 // RLock lock = redissonClient.getLock(myLock.name()); // boolean isLock = lock.tryLock(myLock.waitTime(), myLock.leaseTime(), myLock.unit()); // 获取工厂模式的分布锁 RLock lock = lockFactory.getLock(myLock.lockType(), myLock.name() + UserContext.getUser().toString()); // 获取分布锁的失败策略 boolean isLock = myLock.lockStrategy().tryLock(lock, myLock); if(!isLock) { return null; } try { return pjp.proceed(); } finally { lock.unlock(); } } @Override public int getOrder() { return 0; } } ``` 最后用SpEL优化了一下,固定套路,具体看代码,不赘述了 ## 优惠券折扣 ### 基本业务 该段代码是电商课程订单‑多张优惠券叠加最优抵扣方案的业务实现。支持通用券、指定课程分类限定券;枚举全部优惠券组合,依次叠加抵扣优惠券,优惠金额按课程价格权重分摊,最后选出优惠最大的优惠券搭配方案。 ![使用优惠券.png](picture/%E4%BD%BF%E7%94%A8%E4%BC%98%E6%83%A0%E5%88%B8.png) ### CompletableFuture并发 优惠券最优方案依赖异步线程的计算结果,主线程不能直接向下执行,必须等待异步任务; 可用优惠券一多,全排列组合数量指数暴涨,计算耗时不可控; Tomcat 请求线程资源有限,主线程不能无限阻塞等待,否则会耗尽服务线程、引发接口超时; 设置 1 秒超时就是一种熔断策略:最多等待一秒,超时舍弃未完成任务,优先保障接口快速响应,保护服务器。 ```java // 4.计算方案的优惠明细 List list = Collections.synchronizedList(new ArrayList<>(solutions.size())); CountDownLatch latch = new CountDownLatch(solutions.size()); for (List solution : solutions) { /* 同步运算 list.add(calculateSolutionDiscount(availableCouponMap, orderCourses, solution)); */ /* 异步运算 */ CompletableFuture .supplyAsync(() -> calculateSolutionDiscount(availableCouponMap, orderCourses, solution)) .thenAccept(dto -> { list.add(dto); latch.countDown(); }); } // 等待运算结果 try { latch.await(1, TimeUnit.SECONDS); } catch (InterruptedException e) { log.info("优惠方案计算被中断:{}", e.getMessage()); } ```