分布式定时任务框架的自研实践:锁、负载均衡与OpenAPI异步回调
有段时间我们生产环境的定时任务几乎每两周就要闹一次妖。不是同一批数据被重复处理就是某一个节点挂了以后任务静默消失最离谱的一次凌晨批次跑完以后下游告诉我收到了三份内容相同、任务ID也相同的回调。当时项目里用的是开源的分布式调度中间件理论上不该出现这种问题。我们排查了很久最后发现并不是中间件本身的缺陷而是我们对“定时任务”在分布式环境下的语义理解出了偏差。也就是从那时候起我决定自己写一个轻量级的分布式定时任务框架把负载均衡和OpenAPI异步调用这两件事揉进去彻底搞明白里面每个环节的取舍。这篇文章就是我整理出来的完整复盘架构怎么设计、分布式锁怎么写、任务路由怎么做、外部系统怎么通过OpenAPI异步拿结果以及上线三个月后踩过的那些坑。如果你正在分布式系统里做定时任务为重复执行、任务丢失、回调乱序这些问题头疼或者纯粹想搞清楚一个分布式定时任务框架内部是怎么工作的这篇文章应该能给你不少可以直接落地的思路。1. 动手之前的那笔账为什么不继续依赖现成调度中间件1.1 现成方案解决不了的三个问题先说结论我不反对用现成的分布式调度中间件。Quartz、XXL-Job、Elastic-Job在绝大多数场景下是够用的尤其XXL-Job带控制台、告警和日志对中小团队相当友好。但当时我遇到的是主流方案没打算在框架层面解决的三个问题路由策略不够灵活。我需要按业务维度把任务路由到特定节点而不是简单在所有worker之间轮询。比如处理大客户的数据任务我希望每次调度都落到“持有对应分片”的节点上轮询或随机都做不到。回调链路是封闭的。任务执行完之后需要通过OpenAPI把结果推给外部系统而这个外部系统和中间件生态没有任何关系。中间件虽然有执行日志但让我把“执行结果单独摘出来以标准HTTP抛给下游”就得改框架源码或做一大圈二次开发。依赖太重。有些方案要依赖数据库表来维护节点注册和心跳有些要把一堆class打进执行器。我们只是想用“定时触发分发回调”这个轻量能力不想被框架重度约束。当然我也不是直接拍脑袋决定自己写的。当时我列了一张需求表把诉求分成“必须做”和“可以不做”两类。真正促使我动手的核心诉求只有一个调度器只负责按时触发、按规则分发、把结果安全送出去业务逻辑全部留在业务系统内部。这个诉求在现成框架里往往被包装得很重光是想把一个简单任务变成“可被框架管理、可路由、可回调”的东西就得付出不小的改造代价。1.2 自研边界只做调度、分配和开放调用不做业务编排自己写框架最重要的一件事是先划清边界。我见过很多自研项目翻车不是技术不够而是什么都想塞进去。我做这个分布式定时任务框架时给自己定了三条铁律框架不感知业务。它只管“任务何时触发、触发后交给哪个节点执行、执行结果怎么返回”。任务内部逻辑是一段可执行代码或一个HTTP地址由业务侧注册进来。状态和元数据必须可持久化、可恢复。可以用Redis做秒级调度但任务注册信息、执行记录、结果状态不能只存在内存里否则节点一重启就全部消失。对外只暴露两样东西执行器SDK/通用API和回调/状态查询OpenAPI。回调给下游系统状态查询给调用方监控用。这个边界一旦划清楚后面写代码就不会拉扯。你不需要实现“万能调度引擎”只需要实现“知道任务何时该跑、该跑在哪、跑完怎么告诉别人”的协调器。其余的事包括cron解析细节、线程池调优、业务幂等都可以在使用侧做掉。很多时候我们高估了写框架的难度其实就是把任务拆成几个子系统注册、调度、分发、执行、回调。一个能跑的最小闭环远远比架构图上画得好看重要。2. 框架的骨架任务模型、模块边界和生命周期2.1 核心模块划分当时拆的模块很少但边界非常明确Scheduler解析cron表达式到点触发任务。它不关心怎么执行只负责把“该执行了”这个事件发出去。Registry保存任务定义和执行节点的实时状态。任务定义包含任务ID、cron、路由策略、回调地址、超时时间节点状态包含节点ID、存活时间、当前并发、持有分片。Dispatcher把触发的任务按策略路由到具体执行节点。负载均衡逻辑集中在这里。Executor运行在业务节点上接收调度指令、执行代码、上报结果、维护执行中的任务状态。Callback/OpenAPI把执行结果通过HTTP异步推送给外部系统同时提供状态查询接口。这个拆分看起来普通好处在于Scheduler和Dispatcher都是无状态的可以随时水平扩展Executor挂在业务节点上是唯一带状态的部分用Redis分布式锁做兜底。这样框架里最容易出事的“重复调度”和“任务丢失”被限定在一个我们可以完全掌控和观测的范围里。2.2 待触发队列为什么用Redis ZSet任务注册以后Scheduler如何知道下一个时刻该触发谁这里我没有用传统的“定时扫描数据库表”方案而是用一个Redis ZSet来存放待触发事件。ZSet的score存的是下一次触发时间戳member存的是“任务ID 计划触发序号”的拼接。Scheduler每隔1秒执行一次范围查询String JOB_PENDING_KEY job:pending; long now System.currentTimeMillis(); SetString dueJobs redisTemplate.opsForZSet() .rangeByScore(JOB_PENDING_KEY, 0, now, 0, 100); for (String member : dueJobs) { // 原子地将member移出ZSet能移出的实例才真正抢到这个时间片 if (redisTemplate.opsForZSet().remove(JOB_PENDING_KEY, member) 1) { dispatcher.dispatch(member); } }查到到点任务后把它从ZSet中移除交给Dispatcher分发然后算出下一个触发时间重新写入ZSet。这个设计的精妙之处在于多个Scheduler实例可以同时扫ZSet但只有成功把member从ZSet中移出的那个实例才算真正抢到了这个时间片。ZSet的删除操作是原子的天然充当了一把“时间维度锁”。有朋友问我为什么不用DelayedQueue或者内存队列来替代。很简单DelayedQueue是单机的节点一重启就全丢内存队列也一样。用Redis ZSet之后哪怕Scheduler全部重启待触发的任务事件还在Redis里恢复就是重新扫一遍的事调度不会丢。2.3 任务定义和状态流转一个任务定义长这样{ jobId: job_20240601_001, name: order-sync-to-erp, cron: 0 0/5 * * * ?, routeStrategy: consistency_hash, routeKey: customer_id9527, executor: http://10.0.0.1:8080/run, timeoutMs: 30000, retryTimes: 3, callbackUrl: https://openapi.example.com/v1/callback/order, status: ENABLED }值得展开的字段有两个。routeStrategy和routeKey是一对组合routeStrategy决定用哪种路由算法routeKey给一致性哈希这类算法提供输入用于支持“同一客户的任务落到同一节点”。executor支持两种执行方式一种是业务侧通过SDK注册本地执行实例另一种是直接指定一个HTTP地址让框架充当“通用调度HTTP触发”这种方式对非Java场景尤其友好。任务触发后状态大致经历TRIGGERED、DISPATCHED、RUNNING、SUCCESS/FAILED、CALLBACK_DONE。唯一的问题出在“TRIGGERED到RUNNING”这个区间如果Dispatcher推送失败或者Executor在RUNNING阶段宕机任务可能挂死。我为此加了两个兜底调度超时重试、执行心跳检查。这两个机制在协同工作时也踩过一个坑放到后面的排查部分详细说。3. Redis分布式锁保证任务不被重复执行的关键一笔3.1 锁的语义锁住的是“执行动作”不是“调度动作”提到分布式定时任务第一反应通常是“用分布式锁防止重复执行”。但很多人忽略一个问题分布式锁到底锁的是什么这里有一个常见误区让Scheduler去抢锁抢到锁才分发任务。如果任务执行时间很长Scheduler抢到锁之后中途宕机锁过期另一个Scheduler又抢到锁再次分发同一个任务——重复执行还是存在。所以在我的设计里真正抢锁的是Executor在开始执行业务之前抢锁的key是jobId 计划调度时间窗口粒度细化到“某一次调度”而不是“整个任务”。这样即使调度端重复触发执行端也只允许一个实例真正跑起来。锁的粒度越细业务越安全。用jobId做锁一个任务还没跑完下一个周期又到了两个周期互相阻塞调度就死了用jobId 调度序号做锁每个周期之间的执行互不干扰重复的触发也进不来。3.2 加锁和解锁的工程实现Redis锁我用Lua脚本实现原因只有一个原子性。不能用“先SETNX再单独设置过期时间”这种两步操作中间一旦宕机就变成没有过期时间的死锁。加锁脚本if redis.call(set, KEYS[1], ARGV[1], NX, PX, ARGV[2]) then return 1 else return 0 end对应的Java侧调用boolean locked redisTemplate.execute( new DefaultRedisScript(LOCK_SCRIPT, Long.class), Arrays.asList(job:lock: jobId : scheduleNo), requestId, lockTimeoutMs ) 1L;解锁脚本if redis.call(get, KEYS[1]) ARGV[1] then return redis.call(del, KEYS[1]) else return 0 end每个请求生成一个唯一requestId作为锁的value解锁时先校验再删除。这个“校验删除”几乎是分布式锁最值得重视的细节宁可解锁失败也不能删掉别人的锁。一旦删掉别人的锁紧接着就是双跑。3.3 看门狗续期一个简化但可用的实现有过期时间就必须考虑续期否则长任务会丢锁。业界方案叫“看门狗”核心思想是业务没执行完就持续给锁续期。我的实现思路Executor开始执行时启动一个ScheduledExecutorService定时任务每隔lockTimeoutMs / 3的时间给锁执行一次续期Lua脚本。业务执行结束后先关闭续期线程再执行解锁。节点宕机续期线程随进程消失锁在过期时间后自动释放不会死锁。这个方案比“不设过期时间”稳也比“手动把锁超时设成业务最大耗时”更健壮。上线后确实遇到过业务代码偶发卡顿原本15秒的任务跑了3分钟如果没有续期早就双跑了。后来我又踩过一个更隐蔽的续期跟丢问题放在第6章说。3.4 关于RedLock的取舍这里顺便说下RedLock。很多人一上来就问“要不要用RedLock”。我的观点是绝大多数业务场景不需要。RedLock解决的是多个独立Redis节点之间的分布式共识问题代价是要部署至少5个节点、性能下降、复杂度上升。公司内部Redis用的是主从加哨兵架构单节点锁足够用锁过期风险已经被看门狗兜住了。玩Redis锁先把单节点的原子操作、value校验、续期这三件事做对远比你引入RedLock更有价值。4. 负载均衡任务分配不能只靠“轮询”4.1 任务负载均衡和流量负载均衡的区别这里的负载均衡和Nginx/网关那套不一样。网关的负载均衡对象是请求每个请求的消耗相对均等定时任务的对象是任务有的跑1秒有的跑10分钟。如果纯轮询一个长任务节点会被拖垮其他短任务节点却闲着。所以我一开始就排除了纯轮询。我把几种策略串起来对比过策略核心逻辑适合场景主要缺点轮询依次分配短任务、等权重长任务拖垮节点加权轮询按权重分配节点性能差异大权重需人工维护最小活跃数选当前并发任务最少的节点任务时长不确定需要实时统计并发数一致性哈希对routeKey哈希映射节点任务与数据/客户强相关需处理虚拟节点最终选型是“最小活跃数 一致性哈希”的组合并且给任务留了routeStrategy字段让业务方自己选。框架不替你决定用哪种策略而是把选择权交给最了解任务的业务方。有的任务必须按客户分片路由有的任务只要“尽快跑完”两者诉求天然不同。4.2 最小活跃数策略的实现要点最小活跃数的实现依赖每个Executor上报的“当前并发任务数”。我用Registry聚合心跳数据Executor每次开始任务activeCount1结束任务activeCount-1同时随心跳上报。Dispatcher拿到任务时先过滤掉不健康节点再从剩余节点中选activeCount最小的那个。这里有个细节并发数小不代表真快可能那个节点的机器CPU已经打满但任务还没起来。所以我还加了一个辅助指标最近1小时平均执行耗时。平均耗时越短的节点在同样“最小活跃数”的情况下会被优先选择。这样一来“更快”和“更闲”两个维度都被照顾到了。4.3 一致性哈希在热点任务上的实战一致性哈希是本次项目里最让我长经验的部分。需求是“客户分片任务”同一个客户的多个任务希望每次都路由到持有对应数据分片的节点上。简化实现public class ConsistentHashRouter { private final TreeMapLong, String ring new TreeMap(); private final int virtualNodeCount 200; public void addNode(String nodeId) { for (int i 0; i virtualNodeCount; i) { ring.put(hash(nodeId # i), nodeId); } } public String route(String routeKey) { long h hash(routeKey); Map.EntryLong, String entry ring.ceilingEntry(h); if (entry null) { entry ring.firstEntry(); } return entry.getValue(); } }两个坑必须提醒。坑一虚拟节点数不能太少。节点只有两三个时不引入虚拟节点哈希环上的分布会严重倾斜一个节点可能扛60%以上的任务。我一开始用100个虚拟节点后来在生产里调到200才比较均匀。坑二哈希函数的均匀性远比想象中重要。我最初直接用hashCode()对环取模结果发现“customer_id9527”这种前缀高度相似的字符串hashCode计算后分布非常集中大量任务落到了同一节点。后来换成FNV1_32_HASH分布立刻正常了。写一致性哈希一定要先在本地用真实routeKey集合做压测看分布曲线不要想当然。4.4 动态上下线和慢节点剔除路由到不健康节点是最浪费的事故。我们有一次发布新版本服务启动成功但依赖的中间件还没连上节点处于“假活”状态任务全往这个节点送结果全部失败。这个问题靠心跳感知解决每个Executor每5秒上报心跳节点ID、当前并发数、最近任务执行耗时、错误计数。Dispatcher路由前先从Registry拉健康节点列表连续3个心跳周期没上报的节点直接剔除。连续错误超过阈值的节点标记为“慢节点”最小活跃数策略会把它的并发数人为放大让新任务绕开。这套机制上线后救了我很多次。发布高峰期节点假活、磁盘被日志打满、线程池打爆都被“心跳慢节点剔除”挡在了任务分发之前。5. OpenAPI异步调用让外部系统安全地接收任务结果5.1 为什么用OpenAPI而不是内部RPC任务执行完结果怎么给外部系统外部系统和我们的技术栈完全无关甚至不在同一内网RPC是没法用的。最通用的就是HTTP/HTTPS接口也就是标题里的OpenAPI异步调用。这里我说的OpenAPI不是特指OpenAPI Specification那套规范而是指“对外提供的HTTP API”。当然接口文档你可以用Swagger/OpenAPI来写那也说得通。设计上核心是两件事回调Callback任务出结果后框架主动调下游提供的callbackUrl把结果推出去。适合事件驱动场景比如订单同步完成、批量报表生成完毕。状态查询Query下游不想被动接收推送提供一个GET /api/v1/tasks/{jobId}返回任务状态和结果摘要。适合调用方自己控制轮询节奏。两者配合起来就是典型的异步调用模型提交任务、处理完成通知你、你随时可以来问进度。这和现在很多异步OpenAPI调用的设计是一脉相承的。5.2 回调签名与防重放回调接口必须解决三个问题身份确认、防重放、防重复入账。身份确认我用HMAC-SHA256签名约定appId和appSecretappSecret只在服务端和下游保存。回调HTTP Header带上X-Timestamp、X-Nonce、X-Signature。签名内容用“原始请求体字节 timestamp nonce”拼接后计算。注意不要用JSON反序列化之后的对象再拼接字段顺序一变就验不过了。防重放靠timestamp和noncetimestamp与当前时间差超过5分钟直接拒绝。nonce在Redis里保存5分钟遇到重复的nonce直接拒绝。防重复入账靠幂等键每次调度生成一个全局唯一的scheduleId回调体里原样带上。下游必须用scheduleId做唯一索引或幂等判断。框架重试时scheduleId不变。同一个任务无论重试多少次下游看到的幂等键都是同一个。我一开始以为协议层的签名防重放足够安全后来还是在下游重复入账上栽了跟头第6章详细说。5.3 回调失败时的重试策略回调不是发出去就完事。下游接口可能5xx、可能超时、可能正好在发布窗口。我用的重试策略是“指数退避 最大次数”第一次失败后延迟1秒重试。之后分别延迟5秒、30秒、2分钟、10分钟最多5次。5次都失败任务标记为CALLBACK_FAILED保留原始结果摘要到Redis并打一条明显告警。有一个值得强调的细节回调重试和任务执行重试是两套独立机制。任务执行失败框架在Executor里重跑业务逻辑业务成功但回调失败框架只在回调模块里重新发送结果。如果不区分这两件事你会看到“同样一个任务业务执行结果被下游重复计算了好几次”。5.4 状态查询API的存储设计状态查询看起来简单实际要考虑“结果存在哪、放多久”。我的做法是每次执行结束把执行摘要写入Rediskey是jobStatus:{scheduleId}TTL设7天。查询接口先读Redis读不到再查数据库归档表。摘要内容包含jobId、scheduleId、状态、开始/结束时间、错误码、结果摘要、回调状态。调用方优先用“回调 超时兜底”组合而不是纯轮询。一次任务要跑10分钟纯轮询会浪费大量请求。更优雅的异步调用方式是短任务同步等待回调长任务先返回“已受理”给调用方等回调同时保留一个“超过5分钟没收到回调就查一下状态”的兜底逻辑。这个组合非常贴合真实工程。6. 上线三个月后踩过的坑和完整的排查链路6.1 锁自动过期导致任务双跑的定位过程上线初期有一周我们连续收到下游重复数据投诉。一开始怀疑回调重试但查日志发现回调次数正常于是把怀疑目标转向任务执行本身。完整排查链路取重复数据对应时间点在Redis执行日志里找到该时间窗内的scheduleId。拿scheduleId去Executor日志搜索发现两个不同节点都打印了“开始执行业务”的日志。对比两个节点的锁状态发现第一个节点在执行业务过程中锁就过期了第二个节点成功抢到锁。追查第一个节点为什么没续期发现看门狗续期线程只对“内存中正在执行的任务”追踪而任务执行到某个外部调用时线程被线程池切换了续期逻辑跟丢了。根因不在锁算法而在续期任务跟踪的线程模型。修复方式是续期依据从“当前线程”改为“任务在内存中的注册状态”。执行器用一个ConcurrentHashMap维护“正在执行中的任务集合”只要任务还在集合中续期线程就持续续期不关心它跑在哪个线程上。这个坑让我明白分布式锁的生命周期必须和业务真实执行周期绑定不能和语言层面的线程绑定。后者很容易在线程池、异步化改版时被悄悄破坏。6.2 节点宕机后任务静默丢失的修复另一个让人通宵的坑是任务丢失某个处理报表的节点OOM挂掉了这个节点上的凌晨任务全部没有输出。按理说节点挂了任务应该被重新调度到别的节点但现实是没有。因为执行记录只在Executor内存里Dispatcher根本不知道“某个任务已经接管但没执行完”。修复分两步。第一步把“执行中任务”的状态从内存挪到Redis。Executor收到任务时在Redis写executing:{scheduleId}标记RUNNING并带上节点IDDispatcher定期扫描。第二步引入“孤儿任务回收”。如果任务进入RUNNING状态超过maxExecuteTime比如30分钟且没有心跳续期Dispatcher判定为死任务重新走分布式锁逻辑分配给其他节点重试。由于同一时刻仍然只有一把锁即使原节点没有完全死掉、事后恢复执行也不会造成双跑。这个机制本质上把调度可靠性从“单节点进程的稳定性”转移到了“分布式状态管理”。改完以后节点宕机不再等于任务消失最坏情况是任务延迟完成而且因为锁的存在它不会重复完成。6.3 回调重放导致下游重复入账的补救这个坑最值得写。我一开始认为回调有了签名和防重放足够安全。但忽略了一个场景下游系统自己也有超时重试机制。我们的框架第一次回调下游处理成功了但响应超时下游网关层自动重试又转发了一次同样的回调第二次回调过来下游业务模块没做幂等数据就重复入账了。问题本质是安全协议层能防网络重放防不了业务层的重复投递。幂等必须做在业务层。补救措施在回调协议里强制要求scheduleId作为业务幂等键并在文档里给出下游参考实现——用scheduleId建唯一索引重复请求时返回“已处理”而不是报错。这也是整个项目里我最大的认知升级你替调用方考虑得再周全都不如把“幂等键”这个契约焊死在协议里。因为总会有你想不到的中间环节在悄悄重试。7. 如果重新让我写一遍我会改这些地方最后做点个人复盘。这个框架不复杂核心代码不到三千行但覆盖的问题很真实。如果再让我写一遍有三个地方我一定会动。第一Scheduler的触发模型。当时用单线程cron扫描器保证触发顺序和触发时间的单调性但吞吐有限。再来一遍我会把“扫描触发”和“事件分发”彻底隔离触发器只负责产生时间事件通过消息队列把事件发给Dispatcher这样Scheduler可以更轻松地水平扩展。第二负载均衡的健康判断从“心跳”升级为“真实成功率”。心跳只能说明节点还活着不能说明节点还能干活。下次我会把最近1小时执行成功率、平均耗时、错误码分布都纳入路由决策让“慢节点”的判定更准确而不是用一个简单的错误次数阈值。第三回调模块加慢调用熔断。虽然重试有指数退避但系统故障时重试本身也会形成回调风暴。下一次我会给回调地址增加“连续失败熔断”失败超过阈值进入半开状态只放少量探测请求。这样下游即使彻底挂了也不至于被我们的重试压垮。这个项目的代码被内部持续用了一年多替换了散落在各项目里的临时调度脚本。它不算完美但让我把分布式定时任务、负载均衡、OpenAPI异步调用这几件事的底层逻辑彻底理清了。我的建议是先小范围做一个最小闭环去解决一个真实痛点别一开始就奔着通用平台去。能跑通一个真实场景比架构图上画得再漂亮都有用。之后再根据真实业务反馈慢慢把分布式锁、路由策略、回调机制一层层完善这个框架才会真正变成你自己的东西。