凌晨两点半的告警把我从床上拽起来手机屏幕上全是同一个错误堆栈数据库连接池耗尽沉睡的客户端连接触发了一连串雪崩式的失败重试。我登录跳板机看了一眼任务执行记录发现同一个跑批任务在凌晨两点的窗口期被触发了七次锁竞争、重复消费、死信堆积全凑齐了。那段时间我们正在做核心链路的重构代码里到处是拼出来的分布式任务Cron表达式的写法五花八门没有统一调度没有执行记录没有失败重试策略谁都不敢碰那套定时任务代码改了就要出事。后来我决定不再修修补补而是写一个足够干净、足够专注的调度内核代号就叫“ax”。它的核心工作很朴素统一接收任务注册按触发规则算出下一次执行时间把到期的任务投递到执行线程池跟踪执行结果在异常时触发重试或告警。整个系统没有复杂的编排界面没有自带的一站式监控大屏就是一个可以嵌入业务进程的调度库同时也支持独立部署成调度服务。“ax调度”这个词后来在我们团队内部就成了默认说法凡是涉及异步任务、定时跑批、延迟消息的地方都会先用它来兜底。这篇文章是“ax调度”从设计到落地过程的完整复盘包含核心抽象、触发引擎、并发控制、崩溃恢复和压测调参记录适合正在做调度系统、或者被乱七八糟的定时任务折磨过的后端工程师阅读。1. 先说说我为什么没有被现成框架收编动手之前我花了整整三天时间试用市面上主流的调度框架包括上到企业级调度平台下到各种轻量级的分布式任务库。看起来功能都很全但落到我们的场景里总有几个地方对不上而恰恰是这几个对不上的点决定了最终我选择自研。1.1 我们真正需要的功能优先级我相信每个团队自研基础设施的原因都差不多百分之八十的场景用现成工具能覆盖但那百分之二十的业务特殊性才是核心诉求。我们当时的业务形态是这样的大量短平快的异步任务比如异步通知、文档转码、数据回传执行时间通常在几十毫秒到几百毫秒之间同时还有一批定时执行的扫表任务比如凌晨的仓库对账、每小时的指标聚合、周五晚上的积分结算。任务量不是特别大但要保证极低的失败率并且要能非常清晰地说清楚“某个任务到底执行到哪一步了”。这个需求听起来很简单但现成框架往往把调度、执行、编排、界面绑定得很紧密我无法只取其中一块。我的功能优先级排序是这样的第一是确定性同一个任务不会因为重复触发而重复执行第二是隔离性某个任务的OOM或者线程阻塞不能拖垮调度主进程第三是可观测性每个触发动作都有完整的事件流可以追溯第四是部署形态的灵活性既能和业务进程同进程运行也能独立部署。现成框架要么把调度器和执行器捆得太死要么强调高可用却牺牲了单机场景的轻量性要么操作界面很漂亮但扩展点离业务太远。反正绕了一圈结论就是我们需要的其实只是一个严格定义了边界的小内核其余的能力全部通过接口开放给业务侧自己实现。1.2 这个调度器到底想解决什么不想解决什么写代码之前我做了个很克制的限定ax调度只负责“什么任务在什么时间点触发”不负责分布式任务的分片计算不负责工作流编排不负责可视化平台。边界画在这里是有道理的因为调度器的核心心智模型可以简化成事件流——时间到了就产生一个TriggerEvent执行器负责接收事件并处理。一旦把执行器的分片、排队、编排逻辑全部塞进调度器状态就爆炸了。我一直跟团队说调度系统最忌讳的是成为上帝服务。如果一个调度器既管时间触发又管任务分片还管执行进度和状态同步那等于把业务逻辑全部耦合到一个进程里任何一次常规发布都可能引发不可预期的事件。所以我们坚持把业务逻辑放进“执行器”调度器只做两件事算出触发时间然后在触发时间把执行上下文投递出去。这样即使未来我们要切换底层触发策略或者接入不同的执行环境调度器本身都不用动。2. ax调度的核心抽象任务、触发器与执行器我写代码的习惯是先定义名词再定义关系最后写实现。ax调度的领域模型一共三个角色Task、Trigger、Executor。Task是被调度的目标Trigger决定任务什么时候被触发Executor决定任务被触发之后怎么运行。三者之间的关系是一个Task对应多个Trigger一个Trigger触发后会通过Router绑定到一个Executor。2.1 任务注册时的命名约束与分组策略任务注册是每个接入方第一件要做的事情。注册接口长这样task : ax.NewTask(order.reconcile.daily). WithTrigger(ax.CronTrigger(0 0 2 * * *)). WithExecutor(executor). WithRetry(3, 5*time.Second) ax.Register(task)这里我强制要求每个任务必须有全局唯一的命名命名格式建议是“领域.子域.动作”比如“order.reconcile.daily”就是订单域对账子域的每日动作。为什么这么严格因为任务的监控、告警、日志检索全部基于这个命名来聚合。如果一个任务叫“job1”另一个叫“job2”出问题的时候根本无法在监控系统里快速定位归属团队。为了强制这个规范注册中心里还加了一个插件机制在注册阶段校验任务名是否符合正则不符合直接拒绝注册。分组策略也很简单每个任务声明自己的分组名默认分组是“default”。调度器启动时会为每个分组创建独立的管理器实例这个管理器自己持有触发器和执行队列。分组的价值在于故障隔离——如果“数据仓库”分组的执行器集体卡死不能影响“支付回调”分组的调度。用分组的方式实现物理隔离比在任务级别做并发限制要干净得多因为任务数量大了以后你根本没法逐一给出合理的并发参数但分组数量是可控的。2.2 触发器扩展出来的几种触发语义Trigger接口是整个调度器里最核心的抽象它的定义非常简单type Trigger interface { Next(prev time.Time) (time.Time, error) }只需要给定上一次触发时间返回下一次触发时间。就这么一个接口Cron触发器、固定间隔触发器、延迟单次触发器、按日历触发的触发器全部都能表达。这也是整个系统最让我满意的地方——原本以为会很复杂的调度语义被压缩到了一个方法里。CronTrigger的实现引擎用的是标准的带秒字段的Cron表达式支持“0 0 2 * * *”这样的六位表达式。FixedRateTrigger表示每隔固定毫秒触发一次适合异步轮询场景。OneShotTrigger只在注册后的某个延迟时间触发一次适合延迟通知。DelayTrigger和OneShotTrigger很像但它支持多次执行——每次触发完成后自动计算下一次触发时间为下次延迟加上固定间隔。一个干净的语义模型带来的好处远超预期后来我们接入了很多奇怪的触发需求“每个工作日上午十点半触发”本质上就是往这个接口里塞一个自定义实现代码量基本不超过二十行。2.3 执行器与执行结果回传的边界Executor接口是这样的type Executor interface { Execute(ctx context.Context, task *ax.Task, payload []byte) error }所有业务逻辑都放在Execute方法里。调度器调用Execute时传入一个带超时的context执行器自行决定如何处理payload返回error表示执行失败可被重试。我不允许Executor反向调用调度器的方法避免循环依赖所有状态回传都通过事件总线广播。执行结果是成功还是失败、执行耗时多少、重试了几次这些数据全部由调度器里的“执行记录器”统一采集。这里需要特别强调一点调度器的核心职责是触发而不是执行所以Executor的超时时间、并发限制、重试参数都不应该写死在调度器里。我们的做法是让每个Executor实现一个Descriptor接口描述自己的超时阈值、最大并发数和失败重试策略调度器在执行前先读取这些元数据再做相应的限流和重试。这样接入方可以精确控制自己业务的执行语义又不破坏调度器的统一管理。3. 触发引擎不吃内存的秘密最小堆和时间轮的取舍触发引擎是所有调度器的心脏它决定时间到达精度、内存占用以及高负载下的排序效率。我第一版实现用的是标准库里的timer堆后来经过压测和GC分析做了一次大的重构换成了分层时间轮配合溢出最小堆的混合方案并且把内存占用降到了原来的三分之一。3.1 时间轮和最小堆到底该选谁如果你只在教科书上见过时间轮和最小堆没有实际对比过性能那很容易掉进“时间轮一定更快”的直觉陷阱。时间轮适合大量短周期、高频率的定时任务比如网络框架里的超时检测毫秒级毫秒级地滚动CPU局部性好。最小堆则适合任务数量大但每个任务的下次触发时间跨度很大的场景每次取堆顶就是绝对最早需要触发的任务。ax调度的场景是“任务数量上万个、触发周期从几秒到几天都有”如果用纯时间轮必须开一个非常大的时间轮数组才能覆盖几天的跨度内存浪费惨不忍睹。如果用纯最小堆每秒可能会有大量任务的堆顶更新和重新平衡锁竞争严重。最终方案是组合式一个秒级时间轮负责驱动“最近一分钟”的触发一个大根堆负责保存“超过一分钟”的后续触发点。时间轮转一圈之后从堆里把最近一分钟的新任务拉出来重新填充到时间轮槽位里。3.2 堆节点的更新与新触发时间的计算每个任务会维护一个“当前触发状态”状态里保存了Trigger实例和上一次触发时间。当触发完成后调度器调用Trigger.Next(prev)计算出下一次时间用这个时间更新堆节点。这里有一个非常隐蔽的坑如果任务执行时间很长而Trigger.Next计算出的下一次时间已经过期了说明调度器积压了太多任务此时不能傻乎乎地把过期时间继续往后推而是需要立即触发一次同时记录一次“时钟漂移”警告。我写了一个快速trick堆节点不仅保存时间还保存一个任务版本号。每次任务重新入堆时版本号加一触发引擎从堆里拿到节点后会用版本号校验是否还是最新状态以避免并发情况下把旧节点的执行信息误派发出去。这个版本号方案比加锁简单得多也让触发引擎几乎是无锁读。时间轮槽位里的任务节点同样保存版本号插入槽位前也要做一次conflict检测确认任务没有被取消或更新。3.3 毫秒级触发丢数据的错觉做调度系统的人最忌讳的就是“定时任务没触发”。但分享一个我的观察超过一半的所谓“没触发”其实是客户端时间源和调度服务器时间源不一致造成的。ax调度内部的所有时间基准统一使用单调时钟monotonic clock只在对外展示和日志记录时才转换为墙上时钟。这样做有两个原因第一服务器时间被NTP校准往前跳或往后跳时单调时钟保证调度排序不乱第二任务日志里记录的是墙上时间方便人读和分析。但调度器不可能做到绝对精确我们对精度的定义是“一秒一次扫描提前量控制在正负100毫秒内”。触发时间到了但任务没有立即执行的情况是允许存在的因为执行本身需要排队、需要等待线程池有空位。真正不能接受的是“引擎认为已经触发了但执行记录里没有”这个要么是执行器崩溃要么是状态丢失。我们通过双重写入来规避这个在第五章会详细展开。4. 执行线程池的流量管控从失控到可控调度器把触发事件投递出来后真正的不可控阶段才开始。很多团队做调度器只关心“点”的问题——这个任务触发没有而不关心“面”的问题——触发之后系统扛不扛得住。ax调度在流量管控上踩过一个大坑就是跑批任务把数据库连接池直接拖垮的线上事故。4.1 跑批任务把连接池拖垮的那次故障那个故障是很经典的场景凌晨有一个数据同步的任务正常耗时三分钟触发后拉取订单表、回填指标、清理临时表。有一天上游的数据量突然翻了三倍任务从三分钟变成三十分钟于是同时有多个任务实例堆积在执行队列里。业务团队为了加快处理速度直接提高了执行器的最大并发数从5调到20。结果二十个任务实例几乎同时执行全表扫描和批量写库数据库连接池被打满其他核心业务的查询全部被阻塞最终雪崩。那次故障让我想明白一个道理调度器必须主动管理执行流量而不是把并发数当作一个交给业务方随意调大的旋钮。实现上我们做了三层防护。4.2 信号量与限流器两道防护网第一层是执行池的信号量控制。每个分组有一个独立的执行线程池池子的核心线程数和最大线程数可配置同时池子里内置了一个信号量用来限制同时运行的执行器数量。注意这里信号量和线程池大小是两码事线程池大小决定“能用多少线程”信号量决定“允许多少任务同时跑”。如果信号量满了任务直接进入等待队列不会创建额外线程。第二层是限流器。信号量解决的是“并发上限”限流器解决的是“触发速率上限”。比如某个数据同步任务的触发频率是每分钟一次但任务有时候会自己占用过多资源即使并发为1也会把数据库拖垮。我们用令牌桶做了一层速率限制每执行完一个任务就会消耗一个令牌令牌按固定速率恢复这样即使上游触发再快执行侧也能按照系统能承受的节奏消费。4.3 慢任务隔离的工程实践在执行流量管理里最容易被忽略的是慢任务隔离。一个慢任务长期占用线程资源会导致同一个分组里的快任务得不到执行机会这就是所谓的“队头阻塞”。我们的解决方案是给每个分组内的执行器做二级队列一个优先处理短任务的“快路径队列”一个处理慢任务的“慢路径队列”根据执行历史耗时自动判断任务应该进入哪个队列。首次执行的任务先按默认级别进入快路径队列如果执行耗时超过阈值比如5秒就把该任务标记为慢任务下次触发直接进慢路径队列。这样即使在凌晨跑批高峰期实时性要求较高的轻量任务也能迅速得到线程资源。这个机制上线后任务平均执行延迟降低了70%。5. 宕机之后的事故兜底幂等、恢复与补偿调度系统的可靠性不体现在正常运行期而是体现在故障之后。进程崩溃、机器断电、网络分区这些事故发生时调度器如果没能保证“任务不丢、不重、可恢复”那前面的所有设计都白搭。ax调度在这块的方案是事件落盘、状态机校对、分布式锁容灾。5.1 触发记录与执行记录的两次落库为了防止“触发完成后进程崩溃导致事件丢失”每条触发事件在投递到执行队列之前会先写入一张触发记录表。触发记录表里的字段包括任务名、触发时间、执行状态、入队时间、唯一ID。执行器执行完成后无论成功失败都会回写执行记录表。两张表通过唯一ID关联这样如果我们发现某条记录的触发状态是“已触发”但执行状态是“缺省”就能明确判断出“任务在入队之后丢失了”从而触发补偿恢复。第二次落库发生在执行开始前。执行器真正调用业务代码前会先把执行记录从“排队中”更新为“执行中”同时更新开始时间和节点IP。这一步非常关键因为只有标记了“执行中”崩溃恢复组件才能准确判断哪些执行是“中断的”。在实现上这两次落库是通过本地事务消息表完成的避免同步写入数据库造成性能损耗。5.2 Redis锁和数据库锁的取舍实践分布式调度器的高可用通常需要自动故障转移也就是多个调度实例之间选一个主节点主节点挂了再从节点顶上。ax调度支持多副本部署节点之间通过一个简单的选主算法决定谁是主调度器。我们最初用Redis分布式锁来做选主但很快发现Redis锁在GC暂停时间较长时会出现失效误判导致两个节点同时认为自己是主节点重复触发任务。后来我们把选主改成了数据库的唯一约束心跳续约方案每个节点启动时尝试在调度节点表里插入一条记录带有“实例标识任期编号”的唯一索引只有插入成功的节点成为主节点。从节点持续监控主节点的心跳主节点超过一定时间没有续约从节点删除主节点的记录然后重新尝试插入。整个过程选主是公平的而且不存在锁过期导致的双主问题因为唯一约束本身不允许两行记录同时存在。这个方案缺少Redlock的华丽但在实际生产环境里比Redis锁稳定得多。5.3 崩溃恢复补偿的一个隐含Bug功能做完以后我们做了一个混沌测试把主调度器进程随机kill观察任务恢复情况。测试结果发现了补偿逻辑里的一个Bug当调度器恢复并扫描触发记录表时常常会把“已经由另一个节点执行完成”的任务误判为“丢失任务”导致重复执行。根因是补偿扫描前只检查执行记录是否存在但没有验证执行记录的归属节点。修复方案是给补偿扫描加一个前置条件执行记录表里的“开始时间”必须晚于触发记录里的“入队时间”并且执行节点的心跳在最近一个周期内是存活的。如果执行记录存在且执行节点存活说明任务真的在跑只是还没结束如果执行记录不存在或执行节点已经下线才可以安全补偿。这行判断逻辑至今还在每次回想都觉得它是崩溃恢复设计里最容易写错、也最不该出错的一环。6. 压测数据和真实场景下的调参记录写到这里已经接近一个完整调度器的核心轮廓了最后分享一次有代表性的压测过程和几个关键的调参结论。数据不一定适用于所有业务但思路可以作为参考特别是当你也准备把一个调度内核嵌入到高并发业务里的时候。6.1 单机跑两万任务的压测过程压测环境是8核16G的普通虚拟机目标是验证ax调度在单机场景下的吞吐量上限。测试任务分三类一个是模拟快速IO的空任务执行大约10ms一个是模拟数据库查询的任务执行大约200ms一个是模拟外部接口调用的任务执行大约800ms。任务总数两万个触发规则是所有任务在启动后一小时内按各自频率触发数量分布为每秒触发2000次左右。压测结果记录如下指标空任务组轻查询组外呼组触发吞吐量2100次/秒2050次/秒1980次/秒执行失败率0%0.02%0.05%触发延迟P9923ms25ms31ms执行线程池等待时间P9916ms150ms1200ms内存占用280MB290MB285MB从数据可以看出瓶颈完全不在触发引擎而在执行侧资源。外呼组的线程池等待时间暴涨到1200ms是因为800ms的外部调用吃满了所有线程。这也验证了前面说的慢任务隔离的必要性——如果不隔离这几个外呼任务会把整个执行池堵死。6.2 几个调参结论和它们背后的逻辑关于线程池参数我最终确定的经验值是核心线程数等于CPU核数的1.5倍最大线程数等于核心线程数的2倍队列容量控制在5000以内。这个组合在IO密集型任务居多时表现最好。不要盲目相信“把线程数调大就能提高吞吐”线程数和故障半径是正相关的。关于时间轮的槽位数我最终设置为1024个槽位每一格表示一秒。也就是说时间轮只能精确表示1024秒以内的时间超过这个时间跨度的任务全部放入溢出堆。之前尝试过4096个槽位发现内存占用涨了四倍但性能收益不到5%明显不划算。关于重试参数我们的默认值是三次重试、指数退避并且只在任务失败且Executor没有标记为“不可重试”时才进行。但有一类任务我强烈建议宁可多等一天也不能重试涉及外部支付回调的任务重复通知会造成用户重复扣款的严重事故。所以每次接入新的执行器时我都会手动检查一遍它的重试策略确保语义安全。关于持久化一开始我们尝试落Redis后来发现一个尴尬的问题Redis重启后如果RDB没刷全触发记录会丢。最终选择了本地嵌入式存储引擎来保存触发器状态配合数据库记录事件日志这样既不拖慢主流程又能保证恢复时有足够的信息。最后说说日志和可观测性。调度器必须为每个任务生成完整的事件链从trigger到queued到started到finished每个节点的时间差都要能查得到。我们为此引入了结构化日志关键节点打点输出到一个独立的metrics文件Grafana里面直接拉一张时间线面板。遇到线上问题的时候这张时间线面板是最快定位故障原因的利器比翻各种业务系统日志高效得多。如果你正在被自己的调度任务坑害我建议你先从“完整的事件链日志”开始做再考虑更复杂的高可用和限流设计因为它是所有排查工作的地基。