ARTICLE DETAIL

资讯详情

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

基于库表分段扫描与 Redis 数据预热,构建低延迟的分布式延迟任务触达方案

基于库表分段扫描与 Redis 数据预热,构建低延迟的分布式延迟任务触达方案 基于库表分段扫描与 Redis 数据预热构建低延迟的分布式延迟任务触达方案【免费下载链接】CodeGuide:books: 本代码库是作者小傅哥多年从事一线互联网 Java 开发的学习历程技术汇总旨在为大家提供一个清晰详细的学习教程侧重点更倾向编写Java核心内容。如果本仓库能为您提供帮助请给予支持(关注、点赞、分享)项目地址: https://gitcode.com/gh_mirrors/code/CodeGuide分布式延迟任务延迟队列是互联网业务中高频使用的通用能力本文基于 CodeGuide 仓库中的方案设计文档《基于库表分段扫描和数据Redis预热优化分布式延迟任务触达时效性》完整还原任务库表 门牌号分段扫描 Redis 预热 二阶段消费的整体设计思路并结合仓库内 Redis 实战文档 与 分布式任务调度 的工程化佐证讲解如何用 Redisson 延迟队列快速落地一个可运行的延迟任务触达原型读者读完可以掌握延迟任务场景的选型依据、库表设计要点与 Redis 低延迟消费队列的完整实现路径。一、前言为什么延迟任务要单独做设计在真实业务里很多时间到了才需要动作的需求都会落到延迟任务上。最朴素的做法是一台机器对着数据库表轮询条件满足了就变更状态或插入新数据。当数据量小、使用方少时这种方案足够用但一旦业务体量上来问题就会立刻暴露——比如贷款单息费的产生如果第二天用户还没看到自己的息费信息或者还款后的对账迟迟没有重新执行客诉马上就来了。因此延迟任务方案的演进本质上是围绕两个目标展开的扫描能力海量任务数据在分库分表下需要被快速扫描触达时效部分任务体系要求低延迟处理不能等待过长的扫库间隔。参考仓库内 《第 15 章 分布式任务调度》 中对需求背景的描述扫描库表待结算日息、扫描待开始活动状态、扫描用户会员过期时间、处理异常流程的补偿动作都是典型的延迟任务诉求。单机任务处理能力有限时每天 0 点到 3 点扫描贷款日息可能到了第二天 0 点还没处理完第一批数据——这就是本文要解决的时效性问题。二、延迟任务场景与通用处理流程典型的延迟任务业务场景包括活动开始前的状态变更活动时间到达后任务状态需要从待开始切换到进行中订单结算后的 T1 对账订单结算次日触发对账动作贷款单息费的产生到了计息时间点需要生成当期的息费数据。通用的任务中心处理流程是定时任务扫描任务库表 → 把即将达到超时时间的任务信息扫描到处理队列内存/MQ 消息→ 业务系统消费处理 → 处理完成后更新库表中的任务状态。这个流程在落地时会遇到三个核心问题海量数据规模任务列表数据量较大在分库分表下需要快速扫描单表轮询难以支撑耦合问题任务扫描服务与业务逻辑处理耦合在一起不具有通用性和复用性时效问题细分任务体系中有些是需要低延迟处理的不能等待过长时间的扫库间隔。三、任务表方式通用延时系统的库表设计除了一些较小的状态变更场景业务自己的库表中就带一个状态字段既有程序逻辑变更状态也有到达指定到期时间后由任务服务自动变更的操作对于较大且频繁使用的场景如果每个系统的 N 多张表都各自维护这类字段会非常冗余且不易维护。因此更适合抽出一个通用的任务延时系统各业务系统把需要被延时执行的动作提交到延时系统延时系统在指定时间进行回调回调动作可以是接口或 MQ 消息触达。任务调度表的设计核心有三点职责单一任务调度表只负责拿到什么任务、在什么时间发起动作具体的动作处理仍交给业务工程处理分库分表大批量、多业务的任务集中处理需要设计分库分表满足后续业务体量的增长门牌号设计针对一张表的扫描如果数据量较大又不希望只是一个任务扫描一张表可以让多个任务扫描一张表来提升扫描体量。此时需要一个**门牌号范围标识**来隔离不同任务扫描的范围避免扫描出重复的任务数据。门牌号本质上是给同一张物理表内的数据打上谁负责扫描我的归属标记。多个扫描任务并行推进时各自只认自己的门牌号范围从而在不拆分表的前提下并行提升扫描吞吐。四、低延迟方式Redis 预热 二阶段消费任务表方式可以解决海量 可扩展的扫描问题但扫库存在固有间隔无法满足低延迟触达的要求。低延迟方案是在任务表方式的基础上增加时间把控处理把即将到期的前一段时间内的任务放置到 Redis 集群队列中消费时再从队列中 pop 出来从而更接近任务的处理时效避免因扫库间隔较大而延迟任务执行。整体策略是在接收业务系统提交进来的延迟任务时按照执行时间的长短决定放置位置执行时间较晚的任务先放到任务库再通过扫描的方式添加到超时任务执行队列执行时间临近的任务则同步到 Redis 集群中等待快速消费该方案的设计核心在于Redis 队列的使用同时为了保证消费的可靠性需要引入二阶段消费并可注册 ZK 注册中心至少保证一次消费的处理。1. Redis 消费队列的数据结构设计Redis 消费队列采用槽位Slot化设计槽位计算按照消息体计算对应数据所属的槽位index CRC32 7将消息均匀散列到 8 个槽位中StoreQueue存储队列采用 Slot 方式SlotKey 形如#{topic}_#{index}底层使用Sorted Set有序集合按执行任务分数排序存放任务执行信息分数规则定时消息将时间戳作为分数消费时每次弹出分数小于当前时间戳的消息天然满足到点才能被取出的延迟语义。Sorted Set 是这套设计的关键以执行时间戳为 scorezrangebyscore即可拿到当前时刻之前应该执行的任务集合Redis 单线程模型下的有序结构也保证了并发消费场景下弹出元素的原子性。2. 二阶段消费保障至少一次消费为了保证每条消息至少可消费一次消费者不是直接 pop 有序集合中的元素而是将元素从StoreQueue 移动到 PrepareQueue预备队列并把消息返回给消费者消费成功后从 PrepareQueue 中删除该元素消费失败时从 PrepareQueue 重新移动回 StoreQueue等待下一轮重新消费。这套先移后删、失败回退的处理就是二阶段消费它把取走消息和确认消息拆成两个动作既避免了直接 pop 后消费失败导致消息丢失也通过回退机制兜底了至少一次消费的语义。本文重点放在 Redis 队列的设计上ZK 注册中心、消费进度管理等其他逻辑可以按业务需求扩展完善。3. 简单案例Redisson 延迟队列原方案文档给出了一个可直接运行的 Redisson 延迟队列测试案例使用RBlockingQueueRDelayedQueue组合实现写入后等待消费时间再进行 POP 消费Test public void test_delay_queue() throws InterruptedException { RBlockingQueueObject blockingQueue redissonClient.getBlockingQueue(TASK); RDelayedQueueObject delayedQueue redissonClient.getDelayedQueue(blockingQueue); new Thread(() - { try { while (true){ Object take blockingQueue.take(); System.out.println(take); Thread.sleep(10); } } catch (InterruptedException e) { e.printStackTrace(); } }).start(); int i 0; while (true){ delayedQueue.offerAsync(测试 i, 100L, TimeUnit.MILLISECONDS); Thread.sleep(1000L); } }测试数据输出如下消费线程按延迟时间从队列中取出消息2022-02-13 WARN 204760 --- [ Finalizer] i.l.c.resource.DefaultClientResources : io.lettuce.core.resource.DefaultClientResources was not shut down properly, shutdown() was not called before its garbage-collected. Call shutdown() or shutdown(long,long,TimeUnit) 测试1 测试2 测试3 测试4 测试5 Process finished with exit code -1关键点说明getBlockingQueue(TASK)创建阻塞队列getDelayedQueue(blockingQueue)将其包装为延迟队列offerAsync(测试 i, 100L, TimeUnit.MILLISECONDS)表示消息延迟 100ms 后可见消费线程通过blockingQueue.take()阻塞等待消息到期后自动被取出案例使用了 Redisson 的 DelayedQueue 作为消息队列写入后等待消费时间进行 POP 消费。4. 仓库内的工程化佐证xfg-dev-tech-redis 延迟队列测试该方案的设计思路在仓库配套的 Redis 实战工程中得到了落地验证。在 Redis 缓存、加锁、发布/订阅常用特性的使用 中作者明确标注延迟队列场景即本文方案并在工程测试类RedisTest中提供了同名测试test_getDelayedQueue/** * 延迟队列场景应用https://mp.weixin.qq.com/s/jJ0vxdeKXHiYZLrwDEBOcQ */ Test public void test_getDelayedQueue() throws InterruptedException { RBlockingQueueObject blockingQueue redissonService.getBlockingQueue(xfg-dev-tech-task); RDelayedQueueObject delayedQueue redissonService.getDelayedQueue(blockingQueue); new Thread(() - { try { while (true){ Object take blockingQueue.take(); log.info(测试结果 {}, take); Thread.sleep(10); } } catch (InterruptedException e) { e.printStackTrace(); } }).start(); int i 0; while (true){ delayedQueue.offerAsync(测试 i, 100L, TimeUnit.MILLISECONDS); Thread.sleep(1000L); } }从工程实践看有两处与本方案紧密相关Redisson 客户端的统一封装该工程在app模块的config下通过RedisClientConfigProperties配置类与RedisClientConfig客户端启动类创建 Redisson 连接客户端配置项包括 host、port、password、pool-size、min-idle-size、connect-timeout、retry-attempts 等redissonService封装了getBlockingQueue、getDelayedQueue等能力便于延迟队列在生产配置下直接复用延迟队列与任务场景的结合工程使用xfg-dev-tech-task作为队列名与任务调度的业务语义对应且该文档明确将基于 Redis 队列实现的低延迟任务调度列为 Redis 的核心应用场景之一。此外仓库 Redis 实战文档 还提供了与延迟任务方案配套的 Redis 能力参考独占锁与分段锁无锁化的对比实现、基于自定义注解动态注入RTopic的发布/订阅高级编码。其中**分段锁无锁化**与本文门牌号 多任务并行扫描的设计思想一脉相承——都是通过拆细资源粒度来换取并发吞吐可以在实现任务扫描的并行消费时作为参考。五、总结与选型建议调度任务的使用在实际场景中非常频繁xxl-job、大厂自研的分布式任务调度组件很多都源于很小很简单的功能经过抽象、整合、提炼变成了核心通用的中间件服务设计要匹配体量无论哪种方式的设计和实现都需要考虑该功能后续的迭代和维护性。如果只是一个非常小的场景、没多少人使用那么在自己机器上折腾即可——过度设计和使用有时候会把研发资源拖入泥潭方案分级任务表方式解决海量任务可扫描、可扩展的问题适合对时效要求不苛刻的批量任务Redis 预热 二阶段消费在任务表基础上追加时间把控适合需要低延迟触达的任务体系。实际落地时可按执行时间长短对任务做分级路由让两类存储各司其职技术点即工具库表分段扫描、门牌号隔离、Sorted Set 时间戳排序、CRC32 槽位散列、二阶段消费这些知识点像一件件兵器把它们按各自特点组合起来才能形成一套真正抗用的分布式延迟任务方案。【免费下载链接】CodeGuide:books: 本代码库是作者小傅哥多年从事一线互联网 Java 开发的学习历程技术汇总旨在为大家提供一个清晰详细的学习教程侧重点更倾向编写Java核心内容。如果本仓库能为您提供帮助请给予支持(关注、点赞、分享)项目地址: https://gitcode.com/gh_mirrors/code/CodeGuide创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
返回列表