
上一篇博客写完学习进度统计功能留了个尾巴前端每 15 秒提交一次播放进度直接写数据库高并发下扛不住。这篇记录我是怎么用合并写思路把数据库写频率降下来的Redis Hash 缓存播放进度 JDK 的 DelayQueue 做延迟检测 只在用户停止播放时才落库。核心思想是——95% 的提交都是覆盖式的中间进度只有最后一次提交才有必要写数据库。一、高并发优化的三个方向和写优化的三种方案先建立宏观认知。碰到高并发问题从系统层面看只有三条路水平扩展加服务器、分片、负载均衡。运维的事。服务保护限流、熔断、降级。也是运维和架构层面的手段。提高单机并发能力缩短单个请求的响应时间RT让一台机器能处理更多请求。这才是程序员写代码要解决的问题。要提高单机并发就要看瓶颈在哪。对绝大部分业务来说瓶颈都在数据库。数据库操作分读和写两类读优化大家都很熟——缓存、SQL 优化、索引。写优化很多人不熟悉这篇就是这个。高并发写有三种方案方案核心思路优点缺点适用场景代码/SQL 优化减少单次操作的耗时简单效果有限通用变同步为异步MQ 缓冲先返回再处理削峰、缩短 RT没减少写次数只是降低频率业务链长、多次写库合并写请求Redis 缓存 批量落库既降频率又降次数实现复杂、依赖 Redis高频覆盖式写入异步写之前领券、点赞那几篇都用过了。这篇的重点是合并写——一个用得少但威力大的方案。二、异步写还是合并写判断依据是数据可不可覆盖回到播放进度业务用户看一个 20 分钟的视频每 15 秒提交一次总共提交 80 次。但下一次续播时只有最后一次提交的进度有意义——中间那 79 次全是过时的、被覆盖的数据。这就是决策依据如果业务是每次都要独立处理比如领券、下单用异步写。MQ 只是让请求快点返回每个消息还是要消费落库。如果业务是后续覆盖前面的比如播放进度、位置上报、计数器用合并写。Redis 缓存最新值中间过程完全不需要落库数据库写次数直接减少一到两个数量级。播放进度属于后者选合并写。不一致还在播放一致停止播放 前端心跳每15秒一次 写RedisHash 缓存进度⏱️ 提交延迟任务20秒后检测 20秒后缓存值 任务值?忽略等下次心跳 落库UPDATE 学习记录课表这张图说的是整个合并写的核心流程前端提交先写 Redis 缓存不碰数据库同时提交一个 20 秒的延迟任务。任务到期后对比缓存值和任务记录值——如果变了说明用户还在看忽略如果没变说明用户停止播放了这时才把最新进度写到数据库。原本 80 次心跳 80 次写库变成了 80 次写 Redis 1 次写数据库。三、Redis 数据结构设计为什么按课程分组要缓存的信息有三个字段记录 id用于反查数据库、播放进度 moment、是否学完 finished。三个字段一个对象Redis 里怎么存最直觉的方案是一个小节一个 Hash KeyKEY: learning:record:{sectionId}:{userId} FIELD: id / moment / finished但这样有两个问题。第一一个课程几十个章节用户看几个视频就创建几个 KeyRedis 的 Key 本身有内存开销。第二用户在几个视频之间来回切换每个 Key 单独设置过期时间会造成缓存抖动——刚建好就过期过期又建。最终采用按课程分组的方案KEY: learning:record:{lessonId} FIELD: {sectionId} VALUE: {id:xxx, moment:242, finished:false}一个课表lessonId一个 Hash Key所有小节作为 Hash 的 Field 存进去。 Key: learning:record:1001一门课程 Field: 5001{moment:242, finished:true} Field: 5002{moment:20, finished:false} Field: 5003{moment:121, finished:false} Field: 5004{moment:80, finished:false}这样设计带来两个好处Key 数量少从小节数降到课程数TTL 统一续期用户在同一课程的不同视频间跳转整个 Key 的过期时间都会被续上避免频繁创建销毁。写入代码publicvoidwriteRecordCache(LearningRecordrecord){StringkeyStringUtils.format(RECORD_KEY_TEMPLATE,record.getLessonId());StringjsonJsonUtils.toJsonStr(newRecordCacheData(record));redisTemplate.opsForHash().put(key,record.getSectionId().toString(),json);redisTemplate.expire(key,Duration.ofMinutes(1));}只缓存三个字段id、moment、finished。字段少意味着网络传输少Redis 内存占用少。缓存结构应该按业务需求裁剪不要把整个 PO 一股脑塞进去。四、DelayQueue 原理PriorityQueue 阻塞队列合并写的核心难点是什么时候落库。最直觉的方案是定时任务每 N 秒扫一遍 Redis把该落库的数据写到数据库。但定时任务有个致命问题——时效性和压力不可兼得。产品要求续播误差控制在 30 秒内。定时任务如果设 20 秒一次压力大数据库每秒都在被批量刷设 2 分钟一次压力大是解决了但误差可能到 2 分钟产品不接受。换个角度想这个问题用户每 15 秒提交一次进度如果 Redis 中的进度值 15 秒内没变说明用户已经停止播放了——这时候把最新进度落库就行了。不需要定期刷只需要检测到用户不再提交时刷一次。这就是延迟任务的思路。项目用的是 JDK 自带的DelayQueue。看它的源码publicclassDelayQueueEextendsDelayedextendsAbstractQueueEimplementsBlockingQueueE{privatefinaltransientReentrantLocklocknewReentrantLock();privatefinalPriorityQueueEqnewPriorityQueueE();// ...}拆开来看DelayQueue 内部是三个组件的组合组件作用PriorityQueue存储任务按到期时间排序——最早到期的排最前ReentrantLock保证队列线程安全BlockingQueue语义take()阻塞等待直到有到期任务存入的元素必须实现Delayed接口两个方法publicinterfaceDelayedextendsComparableDelayed{longgetDelay(TimeUnitunit);// 剩余延迟时间intcompareTo(Delayedo);// 任务之间比大小}PriorityQueue 靠compareTo排序越靠近到期的越靠队首。take()方法内部会检查队首元素的getDelay如果大于 0 就继续等待等于 0 就弹出。项目里定义了一个通用延迟任务类DatapublicclassDelayTaskDimplementsDelayed{privateDdata;// 携带业务数据privatelongdeadlineNanos;// 到期时间纳秒publicDelayTask(Ddata,DurationdelayTime){this.datadata;this.deadlineNanosSystem.nanoTime()delayTime.toNanos();}OverridepubliclonggetDelay(TimeUnitunit){returnunit.convert(Math.max(0,deadlineNanos-System.nanoTime()),TimeUnit.NANOSECONDS);}OverridepublicintcompareTo(Delayedo){longdiffgetDelay(TimeUnit.NANOSECONDS)-o.getDelay(TimeUnit.NANOSECONDS);returnLong.compare(diff,0);}}用泛型D data携带业务数据是这里的一个小技巧。这样 DelayTask 就成了一个通用容器任何延迟执行的东西都能装进去——本项目装的是{lessonId, sectionId, moment}。延迟任务方案有四种横向对比一下方案原理优点缺点DelayQueueJDK 阻塞队列 优先级队列无第三方依赖、单机可用占 JVM 内存RedissonRedis SortedSet 发布订阅分布式、不占 JVM依赖 RedisMQ 死信队列TTL 到期转入死信队列分布式、不占 JVM依赖 MQ时间轮Kafka/Netty 的时间轮算法高精度、大量任务实现复杂本项目任务存储时间只有 20 秒任务量可控DelayQueue 简单直接够用。如果生产环境集群化严重、任务量大换成 Redisson 或 MQ 死信也很容易——把工具类里的DelayQueue替换掉其他逻辑不用改。五、延迟检测持久化一次巧妙的值对比DelayQueue 用起来了但还有一个业务细节要想清楚延迟任务到期后怎么判断用户是不是停止播放了方案是这样的每次心跳提交进度时同时提交一个 20 秒的延迟任务任务里记录这次提交的 moment 值。20 秒后任务到期取出 Redis 中当前的 moment 值和任务里记的值对比不相等说明这 20 秒内又有新的提交用户还在看任务直接放弃相等说明这 20 秒没有任何新提交用户已经关了视频/切走了才落库代码里核心逻辑publicvoidhandleDelayTask(){while(begin){try{// 1.获取到期的延迟任务queue.take 会阻塞等待DelayTaskRecordTaskDatataskqueue.take();RecordTaskDatadatatask.getData();// 2.查询 Redis 缓存LearningRecordrecordreadRecordCache(data.getLessonId(),data.getSectionId());if(recordnull)continue;// 3.比较 moment 值if(!Objects.equals(data.getMoment(),record.getMoment())){// 不一致还在持续提交忽略这个任务continue;}// 4.一致说明是最后一次落库record.setFinished(null);recordMapper.updateById(record);// 更新学习记录LearningLessonlessonnewLearningLesson();lesson.setId(data.getLessonId());lesson.setLatestSectionId(data.getSectionId());lesson.setLatestLearnTime(LocalDateTime.now());lessonService.updateById(lesson);// 更新课表}catch(Exceptione){log.error(处理延迟任务异常,e);}}}“设置 finished null 后再 updateById”这个小技巧值得说一下。MyBatis Plus 的updateById会忽略 null 字段——这里手动把 finished 设为 null避免在延迟任务里意外把已经 true 的 finished 又写回 true 触发额外的更新。延迟任务处理器本身是一个 Spring 组件用PostConstruct在启动时开一个线程持续消费队列PostConstructpublicvoidinit(){CompletableFuture.runAsync(this::handleDelayTask);}PreDestroypublicvoiddestroy(){beginfalse;// 优雅关闭让 while 循环退出}这里练习部分还留了一道题目前是单线程消费队列生产环境要改成线程池模式。把CompletableFuture.runAsync换成一个大小固定的线程池每个线程跑一份handleDelayTask循环即可因为queue.take()本身是阻塞的多线程并发 take 是安全的。六、改造后的完整业务流程把缓存、延迟任务、原有逻辑串起来改造后的提交流程长这样️ 数据库⏱️ DelayQueue Redis 服务端 前端️ 数据库⏱️ DelayQueue Redis 服务端 前端20 秒后...提交进度 moment100查缓存命中旧值85判断是否首次学完进度未达50%,不是更新缓存 moment100提交延迟任务(moment100, 20s)立即返回 ✅任务到期查缓存 moment100缓存值 任务值说明用户已停止播放UPDATE 学习记录UPDATE 课表最近学习信息整条链路里用户请求处理只涉及 Redis 读写和延迟任务提交完全不碰数据库。数据库的写入被用户停止播放这个事件触发从每次心跳一次变成每次播放会话一次。首次学完是特殊路径——不能只写缓存因为课表的learned_sections 1、状态从学习中到已学完这些是不可覆盖的状态变更必须同步落库。改造后的代码在检测到finished true时会同步写数据库并清缓存if(!finished){// 只更新进度走缓存 延迟任务taskHandler.addLearningRecordTask(record);returnfalse;}// 首次学完同步落库lambdaUpdate().set(LearningRecord::getMoment,recordDTO.getMoment()).set(LearningRecord::getFinished,true).set(LearningRecord::getFinishTime,recordDTO.getCommitTime()).eq(LearningRecord::getId,old.getId()).update();// 清缓存因为状态已经从 false 变 true缓存里的是脏数据taskHandler.cleanRecordCache(recordDTO.getLessonId(),recordDTO.getSectionId());returntrue;清缓存的动作很重要。因为延迟任务里读缓存时会拿到finishedfalse缓存里存的是旧状态跟数据库里的finishedtrue不一致会导致后续逻辑判断错误。所以首次学完落库后要立刻清缓存让下次读的时候重新回源。七、整体感受这篇的东西我觉得是这几天学到的最有工程味的一段。合并写这个思路一旦掌握能马上联想到很多应用场景位置上报外卖、打车、共享单车、订单计数器、在线人数统计、股票价格刷新——所有高频覆盖式写入的业务都适合。最大的收获是对**“数据的可覆盖性”** 这个判断标准的理解。以前碰到高并发写就是上 MQ 异步没想过还可以更进一步做合并。异步写只是延迟了写写次数没变合并写是消除中间态的写写次数直接砍掉一到两个数量级。另一个收获是 DelayQueue 的用法。之前只在面试八股里见过阻塞队列 优先级队列这种描述这次真的用起来才发现它的泛型设计Delayed接口有多巧妙——你只需要提供剩余时间和排序规则队列帮你搞定排序和阻塞等待。业务代码里一个take()就完事。面试时这个功能我可以聊很久视频续播需求 → 15 秒心跳 → 高频写库压力 → 为什么不用异步写 → 合并写方案 → Redis 结构选型 → DelayQueue 原理 → 值对比判定最后提交 → 首次学完特殊处理 → 缓存清理。一条链子串下来能覆盖缓存、并发、数据结构、JUC 四个知识领域。下一个模块打算看看互动问答评论相关的功能那里应该有树形结构、热点数据缓存这些技术点可以挖。