1. 为什么要自己造一套叫 AX 的调度系统去年年中我接了一个让我头疼到失眠的活把分散在二十多个微服务里的定时任务全部收拢到统一平台支持动态调整执行时间、失败自动重试、按业务方隔离资源还要能随时查看到底哪个任务卡住了。第一反应是直接用现成的分布式任务调度框架但深入调研一圈后我发现团队的真实需求比能跑 cron 就行复杂得多任务要支持灰度发布、执行器要按租户隔离、部分任务对触发时间的精确性要求极高还有一小部分任务根本不能简单按 cron 表达式触发而是需要事件驱动 定时兜底相结合。这套组合需求让我最后决定自己写一个调度内核代号就叫 AX核心含义是 Async eXecution——异步执行与分发调度。这篇文章会把 AX 调度的完整设计思路、核心代码实现、部署和排障过程整理出来。如果你也在纠结要不要自研定时任务系统或者正在搭建一套类似的任务调度平台这篇内容应该能帮你少踩掉大半的坑。我不打算把它写成一本处处正确的教科书而是尽量还原当时做技术决策时真实发生的取舍过程——哪一步被现实教育过哪一步是拍脑袋后面又改的都会说清楚。1.1 不要一上来就自研先用现成组件把问题描述清楚很多团队看到自研调度系统这几个字就兴奋觉得这是做架构的大好机会。我的建议是先把现成组件在你团队里跑一个最小验证再决定要不要自研。道理很简单调度系统的核心难点从来不是写一个定时循环而是如何保证在分布式环境下不重复、不丢失、不阻塞地执行任务这三个问题无论用什么框架都绕不开。我当时拉了一张对比表把常见的方案按团队实际情况过了一遍方案单机任务支持分布式协调动态修改任务多租户隔离二次开发成本Linux crontab支持无需要登录机器改无极低但无法规模化Quartz很好需配合数据库锁一般弱中等集群模式有坑XXL-JOB好自带调度中心支持有基础分组较低但调度中心较重自研 AX按需定制自研实现支持按业务方隔离耗时约 6 人周表格里最扎眼的是 Quartz 的集群模式有坑。之前我们团队就在生产环境吃过一次亏多个节点同时调度同一个任务数据库锁在高峰期出现了明显的竞争任务触发时间被拉长了将近 8 秒。后来排查发现Quartz 默认按数据库记录行锁做分布式互斥这个机制在任务量小的时候看不出问题一旦任务列表上千、单个任务执行超过 30 秒锁冲突就会变得非常敏感。所以与其说是自研不如说是被现实逼出来的方案。AX 的设计目标从一开始就不是替代 Quartz而是聚焦解决前面提到的四个问题动态配置、触发精确度、租户隔离、失败自愈。1.2 AX 调度要解决的四个主要问题第一动态配置问题。定期任务如果每次调整时间都要改代码、发版本、重启进程那这个平台对业务方来说就是不可用的。AX 必须把任务配置变成数据存在数据库里通过管理接口随时改。第二触发精确度问题。有些任务虽然看着像定时任务但业务方要求秒级触发比如库存预占、优惠券过期扫描。传统的每分钟扫一次表方案在数据量上来之后扫描本身就会成为瓶颈而且触发时间可能偏差几十秒。AX 需要在内存里维护触发时间索引而不是每次去数据库里查到点了没。第三租户隔离问题。不同业务方共用同一套执行节点时某个业务方的慢任务可能拖垮所有任务。AX 要把 worker 池按租户分组每个租户有独立的并发度上限和任务队列。第四失败自愈问题。任务的执行环境千奇百怪可能是 HTTP 调用、RPC 调用、跑一段 Python 脚本甚至执行一条 SQL。AX 不应该假设任务一定会成功而是要提供可配置的重试次数、退避策略、超时控制和执行记录让失败的任务能自我修复或者至少被快速发现。这四个问题拆清楚之后架构才有的放矢。下一步就是设计整个系统怎么分层、模块怎么拆。2. 整体架构AX 调度到底在调什么任务调度系统的架构本质上回答一个问题从任务注册到任务执行结果返回数据和控制信号是怎么流转的AX 的架构围绕三条主线展开配置管理线、定时触发线、执行分发线。三条线各自独立又通过一张核心的任务表形成状态闭环。2.1 三个核心概念Job、Trigger、TaskInstance在动手写代码之前我先把概念模型定了。很多人设计调度系统时容易把概念搞混比如把任务既当成定义的配置又当成某一次具体的执行记录这会给数据结构设计带来很大麻烦。AX 里拆成三个模型Job一个任务的定义。包含任务名称、归属租户、执行器类型、执行参数、超时时间、重试策略、当前启用状态。Trigger一个任务被触发的方式。AX 支持 cron 触发、固定间隔触发、立即执行触发以及事件驱动触发。一个 Job 可以绑定多个 Trigger但当时为了简化设计成一对一。TaskInstance某一次实际执行的记录。每次触发都会生成一条 TaskInstance记录执行的开始时间、结束时间、状态、错误信息、执行节点标识。这三个模型的边界非常重要。Job 是常态配置Trigger 是触发策略TaskInstance 是一次瞬时快照。把三者分开后后续做灰度发布、任务历史查询、失败补偿都会简单很多。为了让概念更直观我打一个比方Job 相当于你要做的事项Trigger 相当于闹钟的设置TaskInstance 则相当于闹钟响之后你实际完成这项事项的记录。闹钟可以反复设置但每次响铃留下的记录是独立的。2.2 分层模块Registry、Timer、Dispatcher、ExecutorAX 整体由五个模块组成Registry 配置注册中心负责任务 Job 和 Trigger 的 CRUD统一暴露管理 API并且把配置变更通过版本号同步给调度节点。Timer 定时引擎调度节点的核心模块维护每一秒的待触发任务索引到点后把 Job ID 和 Trigger 时间推送给 Dispatcher。Dispatcher 分发器负责把一个待执行任务发送到对应的租户队列中并生成 TaskInstance 初始状态。Executor 执行器实际运行任务的模块运行方式可以是本进程内调用、HTTP 调用、RPC 调用或者拉起一个子进程。执行器运行结束后回报结果给 Dispatcher。Monitor 监控告警模块订阅任务状态变更统计执行耗时、失败率、队列积压情况超过阈值时发出告警。这五个模块中Registry 和 Monitor 可能是独立的服务而 Timer、Dispatcher、Executor 则常驻在同一组调度节点里方便共享内存状态。多个调度节点之间通过 Redis 进行 leader 选举避免多个节点同时触发同一个任务。需要特别说明的是我不是一上来就搞微服务拆分的。第一版把 Timer、Dispatcher、Executor 放在同一个进程里部署为无状态节点配置信息从数据库拉取后全量加载到内存。这样的好处是网络开销小、调试简单坏处是节点挂了会影响它所在租户的任务执行。因此这个方案依赖了后面要讲的任务表状态机兜底机制保证节点宕机后任务可以被其他节点接管。2.3 定时触发引擎的数据结构设计定时触发引擎是 AX 最核心的部分也是我反复重写次数最多的模块。第一版我采用每分钟扫一次数据库的朴素方案让调度节点每分钟查一次任务表找出 next_run_at 小于当前时间的任务。第一版看着没问题实际上跑起来就发现不对劲。问题出在两个方面。第一几十多个调度节点同时去查任务表哪怕有索引数据库连接也被打满而且还可能出现同一任务被多个节点同时查出来的情况。第二分钟级扫描的触发误差非常大任务可能被延迟到下一分钟才执行这对秒级任务完全不可接受。所以第二版我放弃治标直接重写为内存时间堆 两级时间轮第一级是秒级时间轮数组长度 3600 格每格代表未来某一秒存储该秒需要触发的 Job ID 集合。第二级是分钟级时间轮数组长度 60 格每格代表未来某一分钟存储该分钟需要被展开到秒级时间轮的任务列表。这个设计参考了经典的时间轮算法但细节做了调整调度节点启动时把未来 2 小时内的触发任务加载到时间轮里每 30 秒做一次补仓把接下来 10 分钟的新任务从数据库加载进来。这样一个任务触发的过程就变成了纯内存操作不需要在触发那一刻查数据库。如果你觉得两级时间轮复杂其实也可以先用最小堆来实现。最小堆的好处是实现简单、提前量精确缺点是插入和删除都是 O(logN)。AX 最终选择时间轮主要考虑的是高吞吐场景下时间轮入桶和出桶的时间复杂度都是 O(1)更加稳定。3. 核心实现定时触发、分发、执行的落地代码这一部分我把 AX 的代码细节展开。为避免读者看到一大坨代码晕掉我会按调度循环、时间计算、分发、执行池 4 个阶段来讲贴的代码都是可以单独运行的最小可落地版本。3.1 技术底座Go Redis PostgreSQL技术选型上我用了三件套Go 1.20、Redis 6.x、PostgreSQL 13。选择 Go 是因为调度器对并发模型要求高goroutine 和 channel 写 worker 池非常顺手编译产物又是单一二进制文件部署简单。Redis 用来做节点 leader 选举和租户队列PostgreSQL 用来存任务定义和 TaskInstance 记录。你可能问为什么不直接用 MySQL。原因主要是 PostgreSQL 对 JSON 类型的支持更好Job 的执行参数可以存成 JSONB查询和更新都灵活。另外 PostgreSQL 的 SKIP LOCKED 特性在任务补偿扫描中发挥很大作用这条后面排障部分会专门讲。3.2 调度循环从时间轮到触发任务AX 的调度循环会启动一个常驻 goroutine每 100 毫秒从秒级时间轮里拉取一次当前时间到期的 Job ID。func (s *Scheduler) loop(ctx context.Context) { ticker : time.NewTicker(100 * time.Millisecond) defer ticker.Stop() for { select { case -ctx.Done(): return case now : -ticker.C: s.dueJobs(now) } } } func (s *Scheduler) dueJobs(now time.Time) { // secondWheel 是秒级时间轮当前桶里存着所有计划在此秒触发的 jobID bucket : s.secondWheel.GetBucket(now) if len(bucket) 0 { return } for jobID : range bucket { triggerRunAt : now s.dispatch(jobID, triggerRunAt) // 触发成功后重新计算下一次触发时间 next, ok : s.calcNextRunAt(jobID, triggerRunAt) if ok { s.secondWheel.Add(next, jobID) } } }这里有两个关键点。第一时间轮的粒度设定为 1 秒已经足够满足日常场景如果后续遇到需要毫秒级触发的任务可以把 ticker 周期改成 10 毫秒但随之而来的是调度循环 CPU 开销变大。第二重新计算下一次触发时间是在 dispatch 成功之后才做的避免出现任务还没投递成功时间轮却已经排入下一次的情况。有人会问为什么触发之后才算下一次而不是任务注册时就算好未来若干次原因是任务触发后的执行结果会影响任务是否继续调度。例如一个任务是每 5 分钟执行若执行失败则暂停调度就只有在触发结束后才能决定是否更新 next_run_at。所以在 AX 里调度循环只负责触发触发之后的动作由 Dispatcher 去完成二者通过内部 channel 异步衔接。3.3 Cron 表达式解析与触发时间推算Cron 解析是第一版里最容易被低估的模块。最开始我天真地想cron 不就是五个字段嘛直到遇到每月的最后一个工作日这种需求才发现标准五个字段根本表达不了。AX 的 cron 解析我用了经典的做法解析成若干字段集合然后逐级进位查找下一个匹配时间点。为降低复杂度当时没有全量支持 cron 的所有特殊符号支持范围做了取舍特性支持情况分 时 日 月 周 五字段标准 cron支持支持*/n步进支持支持?模糊匹配支持支持L最后一天/最后一个工作日仅支持周字段的L支持秒级 cron六字段暂不支持支持时区指定支持按任务配置时区进行本地化计算推算下一次触发时间的核心函数是这样一个增量查找func (c *CronExpr) Next(after time.Time) time.Time { t : after.Add(time.Second) // 从秒开始逐级对齐 for { if !c.minuteContains(t.Minute()) { t t.Add(time.Minute - time.Duration(t.Second())*time.Second) continue } if !c.hourContains(t.Hour()) { t t.Add(time.Hour - time.Duration(t.Minute())*time.Minute - time.Duration(t.Second())*time.Second) continue } if !c.dayAndMonthMatch(t) { t t.AddDate(0, 0, 1) t time.Date(t.Year(), t.Month(), t.Day(), 0, 0, 0, 0, t.Location()) continue } // week 匹配也通过后就找到了 if c.weekContains(int(t.Weekday())) { return t } t t.AddDate(0, 0, 1) t time.Date(t.Year(), t.Month(), t.Day(), 0, 0, 0, 0, t.Location()) } }这段代码在逻辑上是对的但它极度依赖逐级进位时的边界处理尤其是月末 21:00 这种会发生跨月跳转的场景。我踩过一个坑cron 表达式0 0 0 1 * *表示每月 1 日 0 点执行但每次推算都会通过加一天的方式往后试逻辑虽正确性能不佳。后来做了一层缓存对同一个 Job 计算未来最近 20 次触发时间并缓存之后直接从缓存里取直到缓存耗尽才重新计算。这样把 cron 推算的开销从每次触发一次降为20 次触发一次。3.4 分发任务与执行器 Worker 池任务到点后Dispatcher 从 Job 配置中读取执行器类型、参数、超时阈值构造一个 TaskInstance然后投递到对应租户的 channel 队列中。每个租户有独立的 worker 池。func (d *Dispatcher) dispatch(job *Job, runAt time.Time) { task : TaskInstance{ JobID: job.ID, RunAt: runAt, Status: pending, CreatedAt: time.Now(), } d.db.Create(task) d.enqueue(task) } func (d *Dispatcher) enqueue(task *TaskInstance) { ch, ok : d.tenantChannels[task.TenantID] if !ok { ch make(chan *TaskInstance, d.tenantQueueSize) d.tenantChannels[task.TenantID] ch d.startTenantWorkers(task.TenantID, ch) } ch - task }租户 worker 池的大小、queue 长度和任务复杂度相关。我当时给每个租户设置了默认并发 8队列长度 1024。超过队列长度后新任务不再进入内存 channel而是直接标记为调度成功但执行排队中,由数据库 executor 补偿线程去捞。这一步很关键因为它避免了内存队列堆积导致进程整体 OOM。Worker 侧的代码非常直观func (d *Dispatcher) startTenantWorkers(tenantID string, ch -chan *TaskInstance) { for i : 0; i d.TenantWorkerConcurrency(tenantID); i { go func() { for task : range ch { d.execute(task) } }() } } func (d *Dispatcher) execute(task *TaskInstance) { ctx, cancel : context.WithTimeout(context.Background(), task.Timeout) defer cancel() err : d.invokeExecutor(ctx, task) if err ! nil { d.handleRetry(task, err) return } d.markSuccess(task) }这里要特别讲解一个隐蔽的问题TaskInstance 的执行状态不能只存在内存里。如果 worker 执行到一半节点崩溃内存里的状态就丢了任务会被永久遗忘。所以 AX 的做法是在进入执行前先把状态从 pending 更新为 running执行完成后立即更新为 success 或 failed。每次状态变更都伴随一个版本号通过乐观锁防止并发覆盖。4. 常见问题与排查实录这一部分是 AX 上线后我和团队与各种故障搏斗出来的经验。每一个问题都真实发生过而且都能在网上搜到相似的求助帖。我把排查思路和最终处理方式一并整理成速查表。4.1 任务被重复执行了源头到底在哪里上线第三周业务方反馈某个定时任务一天内被执行了两次影响了一部分数据的正确性。这种问题在分布式调度系统中非常经典原因一般有三个调度节点重复触发、分发网络重试、执行节点超时导致的重试叠加。排查时我先看任务状态表发现确实有两条成功的 TaskInstance间隔只有 1 秒这基本可以确定是触发阶段重复。进一步追日志发现当时两个调度节点都以为自己抢到了触发权——因为我早期多节点触发没有做好互斥两个节点同时从时间轮里取到任务。修复方式有两种一种是只在 leader 节点上运行 Timer 循环非 leader 节点禁用时间轮触发另一种是任务触发前在 Redis 里加一个短时锁锁的 key 是 jobID:nextRunAtTTL 设为 10 秒。最终我两种都做了leader 选举负责整体调度的唯一性Redis 锁作为极端情况下的兜底。你以为这就完了没有。后来我又遇到一种更隐蔽的重复发生在执行阶段Worker 在调用执行器时HTTP 客户端因为读超时返回了错误但服务端其实已经处理完成于是重试逻辑又执行了一次。这个问题无解的地方在于超时不确定到底成功没有只能靠任务幂等性来兜底。所以我从架构层面要求所有接入 AX 的业务方提供幂等键并把这个要求写进了接入规范。4.2 任务积压导致队列堆满如何快速恢复某天早上业务方来反馈财务对账任务比平时晚跑了 40 分钟。我一看监控面板Redis 里的租户队列深度打到了 9000 多明显是上一个长任务把 worker 全部占满了。这个场景暴露了我在 2.1 概念模型里没有设计好的地方一个 Job 在同一个租户内不能无限并发执行。虽然每个租户有固定 worker 数但超时时间太长比如某个 HTTP 任务设置了 5 分钟超时它就能占住一个 worker 很久后续任务只能排队。修复方式我给 TaskInstance 增加了同一租户下同一 Job 最多允许 N 个并发实例的限制并在入队前检查。具体实现很简单在数据库里加一个计数判断SELECT COUNT(*) FROM task_instance WHERE tenant_id $1 AND job_id $2 AND status IN (pending, running) AND created_at now() - INTERVAL 10 minutes;如果数量大于等于并发上限就不再入队而是把任务标记为 skipped 并记录原因。这样慢任务不会无限挤占队列同时业务方也能在监控里看到因为并发超上限被跳过的事件。4.3 节点宕机后任务消失了补偿扫描怎么设计调度节点宕机内存时间轮直接清空这时候会面临一个问题宕机瞬间内存里尚有 2 小时内需要触发的任务重启后这些任务不会自动出现。单纯依赖启动时加载未来 2 小时任务还不够因为宕机期间可能已经错过多次触发机会。我的解决思路是周期性补偿扫描。调度节点每 5 分钟执行一次 SQL查找所有 next_run_at 在过去 5 分钟内、状态仍为 enabled 的 Job然后重新触发一次。为了避免多节点同时扫描到相同任务造成重复扫描带上FOR UPDATE SKIP LOCKED子句让每个 Job 在同一轮扫描中只被一个节点拿到。PostgreSQL 的 SKIP LOCKED 在这个场景下非常好用。它和普通FOR UPDATE的区别是遇到被其他事务锁定的行不是等待而是直接跳过。这样多个调度节点并发扫描同一张表时既不会互相等待也不会拿到重复任务。补偿扫描也会引入另一个问题任务确实触发了但执行失败了下一次补偿又把它重新拉起来造成反复执行。因此我在任务表上加了 last_attempt_at 字段补偿扫描只会拉起那些距离上次触发时间超过 10 分钟且当前没有 running 实例的任务避免无限追打。4.4 时钟偏移导致任务提前触发原本以为是极小概率的事件结果在压测环境真出现了某调度节点系统时钟比真实时间快了 3 秒导致所有任务提前触发。调度系统对时钟的敏感度极高因为时间轮的所有判断都基于当前时间。排查方式是给每个节点增加了一个时钟服务校验周期性地从 NTP 服务器获取标准时间如果偏差超过 500 毫秒就把该节点标记为 unhealthy暂停调度权限。另外在节点启动时也会先校时再加载时间轮。在不能修改宿主机时钟的容器环境下这个校验逻辑尤其重要。容器内系统时间不准是一个常态问题不能只靠date命令肉眼判断。4.5 常见问题排查速查表现象排查方向推荐解决方式任务重复执行检查 task_instance 是否有两条近似时间的成功记录确认 leader 选举 Redis 触发锁任务延迟触发查看租户队列深度、worker 占用率设置同 Job 最大并发调大 worker 数任务丢执行检查是否存在节点宕机、补偿扫描未触发缩短补偿扫描周期确保 at least once任务一直被重试检查执行器超时设置是否过短根据执行器 P95 耗时调整超时时间负责人未收到告警检查监控模块订阅是否绑定 job owner告警按租户和 job 双维度配置数据库表锁冲突查看调度节点数量、扫描 SQL使用 SKIP LOCKED 或轮流抢占5. 一些没有写在文档里的经验最后这部分说点代码之外的东西。自研调度系统不是一个纯粹的技术项目它牵涉到团队协作、运维规范和业务方预期管理。以下几条经验是我在 AX 落地后复盘总结的严格来说不属于任何文档章节但价值不亚于上面的架构实现。5.1 先保证 at least once再考虑 exactly once这是整个系统设计中最重要的一条原则。分布式环境下的任务执行要做到 exactly once 极其困难尤其在任务执行结果可能丢失、网络可能抖动的情况下。AX 一开始就明确承诺的是 at least once也就是任务一定不会丢但极端情况下可能重复执行然后用业务层面的幂等来消除重复的影响。这条承诺让整个系统设计大大简化。我不需要去实现复杂的分布式事务来保证执行结果的一致性只需要保证触发不丢、状态可查、结果可追踪。和业务方沟通时也直接说明接入 AX 的任务请自行提供幂等键。团队磨合了两周后大家才意识到这其实是最省事的方案。5.2 监控和告警比调度本身更值得投入坦白说我最早把百分之八十精力放在调度引擎上监控告警只做了个简单日志。结果上线后第一个发现任务异常延迟的人不是我而是业务方的运营同学。那时候我才后知后觉调度系统做得再好如果出了问题不能在 30 秒内被发现对业务来说就是不可用的系统。后来我单独做了一个任务健康度面板核心指标就四个调度延迟从计划触发到实际派发的毫秒数、执行超时率、失败重试次数、队列积压深度。告警阈值参考了行业实践并结合我们自身业务情况做了调整调度延迟超过 2 秒就通知失败率超过 5% 就通知队列积压超过 500 就通知。这套面板在后续排查中帮了大忙。5.3 自建调度器的边界是什么如果你现在问我团队到底该不该自研调度系统我的回答会非常谨慎。这里给出几个我认为适合自研的判断条件你们对触发时间的精确度要求是秒级甚至毫秒级你们有多租户隔离的硬性要求并且现成框架难以满足你们有专门的平台团队能承诺为期一个月的开发周期你们已经很清楚任务执行环境的异构程度准备好统一执行器规范。如果以上四条一条都不满足我强烈建议优先使用成熟的开源调度平台。调度器是一个需要长期演进、持续压测的系统业务团队如果在快速发展期自研的沉没成本会非常高。回到 AX 本身这个项目最终确实扛住了每天几百万次触发支撑业务方顺利推进了灰度发布和租户隔离的改造。但真正让我满意的并不是代码写得多优雅而是整个团队终于对任务为什么会被执行执行到哪一步了失败了该怎么办有了统一的认知。这种确定性才是调度系统价值的核心。