自研分布式任务调度引擎ax:设计、实现与线上踩坑
发布时间:2026/9/26 10:49:17来源:尧图网络
我把一个内部代号叫“ax”的任务调度引擎从设计到落地重新梳理了一遍。先说背景我们团队维护着几百个定时脚本和数据任务最早全靠 crontab 加上一批值班同学盯着后来脚本一多依赖一乱每天凌晨总有几个任务重复跑或者干脆不跑数据报表错了还得人工回溯。与其继续头疼不如做一个统一的调度模块把“时间触发、依赖编排、失败重试、分布式状态一致性”这些事情一次理清楚。这篇文章就把 ax 的完整设计思路、核心代码、以及线上跑了一年踩过的坑全部写出来后端研发、平台运维、还有被 shell 脚本和 crontab 折磨的同行都可以参考。1. 先想清楚ax 调度究竟在“调度”什么1.1 调度的本质时间、条件与执行者三方博弈调度器要解决的问题嘴上说叫“定时任务”实际上本质是三个维度的博弈什么时间执行、满足什么条件才能执行、由谁来执行。一个只支持 cron 表达式的定时器解决的是第一个维度到了依赖场景还要看上游任务有没有成功再到了多台机器部署还要决定任务到底派给哪个 worker。把一个简单 shell 脚本交给调度系统看起来就是把 exec 和 cron 组合一下但真正要设计的是一个“决策器”它拿着所有任务的时间、状态、资源信息在每一秒都回答“现在该派发谁”。我用医院叫号系统做个类比。普通排队是时间顺序急诊是优先级插队检查科室之间又有依赖先抽血才能拿化验单而调度器就像分诊护士——它不是埋头报号而是盯着各种条件变化做决策。任务调度的复杂度恰恰不在定时触发本身而在“条件什么时候满足”以及“条件满足了之后下一步怎么走”。所以 ax 从一开始就定位成任务调度引擎而不是“crontab 的替代品”。它需要支持定期执行、手动触发、上游完成后自动触发还要把执行记录、重试状态、机器心跳全部管起来。想清楚这一点后面数据模型和代码结构才不会跑偏。1.2 疯狂复用的四个核心模块我最终把 ax 拆成了四个模块调度器、执行器、任务仓库、状态存储。调度器负责扫描任务、计算下一次执行时间、抢占待执行任务、派发执行指令。它只做决策不执行具体业务。执行器常驻的 worker 进程从队列里领取任务用 subprocess 方式运行真实的命令或者脚本最后回传结果。任务仓库存放任务本身的信息比如任务名称、cron 表达式、执行命令、启用状态。状态存储记录每一次执行的过程数据包括触发时间、开始时间、结束时间、状态码、输出摘要也是做幂等和排查问题的依据。有人会问任务仓库和状态存储不能合并成一张表吗最好不要。任务元数据是低频变化的信息执行状态是高并发高频写入的信息放在一起会导致行锁竞争和查询性能互相干扰。ax 里任务表和日志表是严格分离的任务表负责“该跑什么”日志表负责“跑得怎么样”各自优化互不拖累。1.3 为什么我选了自研这条老路选型阶段我认真对比过 Quartz、APScheduler、Airflow、DolphinScheduler 这几个常见方案。Quartz 是 Java 生态的标杆功能全但太重团队以 Python 为主集成成本不低APScheduler 单机调度很舒服但跨进程和失败恢复都得自己补Airflow 适合做复杂的 ETL 工作流调度语义很强但部署和运维对一个小团队来说有点“杀鸡用牛刀”DolphinScheduler 功能完善不过它本身是一个平台接进来就得改自己的使用习惯。自研 ax 并不是因为“别人都做得不好”而是我们需要的刚好是一层“只做调度、不绑业务”的薄内核。给任务一个统一的调度入口把时间触发和依赖触发都收编进来然后对外暴露一个简单的 SDK 和接口。这样团队改动业务脚本时不用学习整套平台概念内部遗留脚本也能快速迁移。但要说清楚自研不是重复造轮子。ax 里cron 解析直接用 croniterHTTP 调用用 requests分布式锁就用数据库行锁或 Redis 的 SETNX。我们只写调度状态流转和分布式抢单这一层真正属于“调度引擎”的代码其余部分能用成熟库绝不动手。原则很简单核心逻辑必须自己掌控辅助能力全部拥抱开源。2. ax 的核心调度逻辑拆开揉碎讲2.1 时间驱动为主、事件驱动为辅的混合触发任务触发有两种典型模型。一种是时间驱动调度器定期扫描任务表凡是 next_run_time 小于当前时间的任务就进入待派发队列。它的优点是实现简单、状态容易恢复任何任务只要查表就知道下次什么时候跑缺点是实时性受扫描频率影响扫描间隔太短又会对数据库造成压力。另一种是事件驱动上游任务执行完发一条消息给下游让下游立刻进入可执行状态。它的优点是响应快、不依赖轮询但分布式环境下消息丢失、重复消费都是额外成本。为了追求“能跑就行”的稳定ax 采用了混合策略时间戳作为兜底骨架调度器每 0.5 秒扫一次库依赖关系则通过在状态存储里写入“依赖完成事件”同时把关联任务的 next_run_time 更新为当前时间用事件加速用时间兜底。这个设计的好处是哪怕事件丢失了下一次扫描也能把任务捞起来哪怕调度器重启也不会漏掉到点任务。工程上少赌概率多留保底才是调度系统的生存之道。2.2 下一次执行时间到底怎么算cron 表达式看起来简单真正算准并不容易。ax 直接用了 croniter 这个库去计算下一个执行时间代码很简单from croniter import croniter from datetime import datetime, timezone # 以当前时间为基准计算下一次执行时间 base datetime.now(timezone.utc) cron croniter(0 2 * * *, base) next_run_time cron.get_next(datetime)但有几个细节必须注意。第一任务表里存的应该是“下一次执行时间”而不是“cron 表达式本身”因为每次调度完成后要立刻算出并更新下一次时间避免同一任务被重复扫描。第二所有时间统一用 UTC 存储展示层再转换成本地时区。否则服务器时区设置不一致同一个任务在不同机器上可能触发两次。第三croniter 默认支持的就是标准 5 位和带秒的 6 位表达式ax 统一只用 5 位标准 cron秒级任务几乎没有实际意义反而会让日志膨胀和锁冲突加剧。还有一个很容易踩的坑像“每月最后一天执行”这类特殊需求标准 cron 表达不了。ax 的做法是在任务表里加一个 trigger_type 字段如果值是 last_day_of_month就用单独写了一个 helper 去计算每月最后一天而不是强行扩展 cron 语法。2.3 依赖关系从单任务到 DAG任务生命周期里最麻烦的是依赖。一个典型的业务链路可能是“数据采集 → 清洗 → 指标计算 → 报表推送”每一步都依赖上一步成功。如果只靠定时轮询上游延迟半小时下游要么空跑要么失败。ax 的任务表里加了一张 dag_edge 表专门记录边的依赖关系。原始设计也考虑过在 task_info 上用 parent_ids 数组字段但每次查下游任务都要扫全表。边表更灵活支持新增、删除和环检测CREATE TABLE dag_edge ( id BIGSERIAL PRIMARY KEY, parent_task BIGINT NOT NULL, child_task BIGINT NOT NULL, trigger_type VARCHAR(16) NOT NULL DEFAULT ALL_SUCCESS, UNIQUE (parent_task, child_task) );触发条件我用了一个状态机任务初始是 WAITING等所有上游都进入 SUCCESS 状态后变成 READY调度器扫描到 READY 就派发给执行器执行器把它置为 RUNNING最终到 SUCCESS 或 FAILED。多上游场景下支持两种策略ALL_SUCCESS 表示所有上游成功才触发ALL_DONE 表示任意上游成功都触发。实际业务里大多数是前者。注册 DAG 时还有一个隐形需求环检测。如果 A 依赖 B、B 依赖 A整个调度会死锁。ax 在保存边的时候会先跑一遍拓扑排序出现环就拒绝写入。这个校验虽然简单但没做之前我亲眼见过一个团队因为手工配置依赖把整个调度链路卡了一晚上。2.4 分布式环境下不抢活、不丢活调度器通常要部署至少两个实例来保证高可用。但多实例会带来一个经典问题同一时刻两个调度器都扫到了同一个到点任务如果都去派发执行就是双份。ax 解决这个问题用的是数据库乐观锁核心 SQL 是条件更新UPDATE task_info SET next_run_time :new_time, trigger_version trigger_version 1 WHERE id :task_id AND next_run_time :old_time;只有影响行数为 1 的实例才算抢占成功下一个调度器即使扫到了同一条记录也会因为 next_run_time 已经被改掉而放弃。这套逻辑比 select update 更安全因为检查和更新是在一条语句里完成的天然防并发。执行层面也要考虑幂等。每次派发一个任务我都生成一个全局唯一的 execution_id执行日志用它作为主键记录状态。就算某个极端情况下调度器重复派发worker 在写入结果时发现 execution_id 已经存在就直接忽略本次结果。调度系统至少要做到 at-least-once真正去重和幂等由任务本身配合完成。这个认知越早建立线上事故越少。3. 从零撸一个 ax 调度最小闭环3.1 最小数据模型与建表语句ax 的最简版本只需要三张表任务表、执行日志表、调度锁表。任务表负责描述“有什么任务”日志表负责记录“每次跑得怎么样”锁表负责分布式抢单。下面是 PostgreSQL 版本的建表示例你可以在自己的项目里直接抄CREATE TABLE task_info ( id BIGSERIAL PRIMARY KEY, name VARCHAR(128) NOT NULL UNIQUE, cron VARCHAR(64) NOT NULL, command TEXT NOT NULL, status VARCHAR(16) NOT NULL DEFAULT ENABLED, next_run_time TIMESTAMPTZ NOT NULL, trigger_version BIGINT NOT NULL DEFAULT 0, owner VARCHAR(64), created_at TIMESTAMPTZ NOT NULL DEFAULT now(), updated_at TIMESTAMPTZ NOT NULL DEFAULT now() ); CREATE TABLE task_log ( id BIGSERIAL PRIMARY KEY, execution_id VARCHAR(64) NOT NULL UNIQUE, task_id BIGINT NOT NULL, trigger_time TIMESTAMPTZ NOT NULL, start_time TIMESTAMPTZ, end_time TIMESTAMPTZ, state VARCHAR(16) NOT NULL, result_summary TEXT ); CREATE TABLE task_lock ( task_id BIGINT PRIMARY KEY, locked_by VARCHAR(64), lock_until TIMESTAMPTZ );这里有个细节值得解释为什么 task_log 里要单独放 execution_id 而不是直接用 id 做主键因为 execution_id 是调度器生成的worker 通过它来保证幂等数据库自增 id 只用于内部关联。两者分开职责清晰。3.2 调度器主循环与抢单逻辑调度器的主体就是一个循环扫描到期任务 → 尝试抢单 → 更新下一次执行时间 → 派发到执行队列。核心代码长这样import time from datetime import datetime, timezone def scheduler_loop(): while True: tasks fetch_due_tasks(limit20) for task in tasks: if claim_task(task): dispatch(task) time.sleep(0.5) def fetch_due_tasks(limit): return db.query( SELECT * FROM task_info WHERE status ENABLED AND next_run_time now() ORDER BY next_run_time ASC LIMIT :limit FOR UPDATE SKIP LOCKED , {limit: limit}) def claim_task(task): new_next croniter(task[cron], task[next_run_time]).get_next(datetime) result db.execute( UPDATE task_info SET next_run_time :new_time, trigger_version trigger_version 1 WHERE id :task_id AND next_run_time :old_time , { new_time: new_next, task_id: task[id], old_time: task[next_run_time] }) return result.rowcount 1fetch_due_tasks 里的FOR UPDATE SKIP LOCKED很关键。它的意思是只锁住那些“我先抢到”的行已经被其他调度器锁住的直接跳过而不是排队等待。这样两个调度器实例同时跑性能不会互相拖垮。接着要讲一个我最初写错的点。刚开始我没更新 next_run_time只把它置为一个“已派发”标记。结果调度器每 0.5 秒扫一次同一个任务永远处于“到期”状态worker 还没执行完第二次就派发出去了。后来改成上面这种“抢单时直接更新下一次执行时间”的写法才彻底解决重复派发问题。3.3 执行器 worker 的运行细节执行器从队列里拿到任务后会真正执行命令。这里不能简单地用 Python 的 os.system因为需要拿到退出码、输出内容、超时控制。我用 subprocess 来做进程隔离import subprocess import shlex import time def run_command(command, timeout_seconds3600): proc subprocess.Popen( shlex.split(command), stdoutsubprocess.PIPE, stderrsubprocess.PIPE, textTrue ) try: out, err proc.communicate(timeouttimeout_seconds) except subprocess.TimeoutExpired: proc.kill() out, err proc.communicate() return { state: TIMEOUT, output: out err } return { state: SUCCESS if proc.returncode 0 else FAILED, output: out err, exit_code: proc.returncode }每个 worker 在真正执行前还会先向 task_lock 表写入一条“租约”记录执行开始 lease定期续期。如果 worker 突然被 kill锁记录会在 lease 到期后自动失效调度器可以把任务重新派发给别的机器。这个机制我放到第 4 章细讲因为它是线上稳定性最关键的兜底。执行结果回传时worker 只更新 task_log 的对应记录不修改任务本身的 next_run_time。任务的下一次执行时间已经在调度器抢单那一刻算好了。所以调度器和执行器的职责是严格分离的调度器说“你可以跑”worker 说“我跑完了”。3.4 端到端验证一个任务从注册到执行完第一轮测试我只注册了一个最简单的任务每分钟执行一次 echo通过 SQL 手动插入INSERT INTO task_info (name, cron, command, next_run_time) VALUES (hello_test, * * * * *, echo hello ax scheduling, now());然后启动 scheduler 和 worker 两个进程盯数据库里的 task_log。正常情况下一分钟后 task_log 里会有一条 state 为 SUCCESS 的记录。接着我改了 cron 为*/5 * * * *观察 next_run_time 是否按 5 分钟跳变确认 cron 计算没问题。验证幂等的方法更简单同时启动两个 scheduler 实例观察同一个任务是否生成了两个 execution_id。只要任务表里的 next_run_time 更新逻辑没问题重复派发这件事基本不会发生。如果出现了重复优先怀疑是不是有人在事务里把 old_time 判断写丢了。4. 线上跑了一年最值得写的坑4.1 任务不触发和重复触发先查这三件事线上问题里出现频率最高的就是“该跑的任务没跑”。我的排查顺序永远是先看 task_log 有没有这个 execution_id如果有问题在 worker 层如果没有问题在调度器层。调度器层最常见的原因是时区。曾经有一台服务器系统时间是 CST另外一台是 UTC任务表里 next_run_time 存的是 UTC但调度器用本机时间做比较结果两台机器对“现在”的判断差了 8 小时凌晨两点半的任务一台提前跑、一台干脆不跑。这个问题的根子在于比较时间时应用代码必须显式使用 timezone.utc不能依赖系统默认时区。重复触发的原因更隐蔽。有一次排查发现同一个任务在凌晨两点被派发了三次日志看起来一模一样。最后定位到是“补偿扫描”逻辑写错了调度器重启后把过去一小时的任务全部补扫但我没判断任务是否已经被成功执行过结果把已经跑完的任务重新塞回队列。修复方案是补扫时只挑 state 是 WAITING 或者 FAILED 的任务已经成功的直接跳过。排查工具我用得最顺手的还是这条 SQLSELECT id, execution_id, task_id, trigger_time, start_time, end_time, state FROM task_log WHERE task_id :task_id ORDER BY id DESC LIMIT 20;执行记录一查是调度器问题还是 worker 问题基本一眼就能看清。4.2 执行器假死与僵尸任务处理worker 进程被杀掉正在执行的任务会被迫停在 RUNNING 状态这就是僵尸任务。如果不处理下游依赖会一直卡住整个 DAG 悬挂在那。ax 的解决思路是租约机制。worker 执行任务前在 task_lock 表里写入锁记录INSERT INTO task_lock (task_id, locked_by, lock_until) VALUES (:task_id, :worker_id, now() interval 2 minutes) ON CONFLICT (task_id) DO UPDATE SET locked_by EXCLUDED.locked_by, lock_until EXCLUDED.lock_until;worker 每 30 秒续一次租约把 lock_until 往后推。调度器扫描僵尸任务时只需要检查锁记录是否过期如果一个任务的 state 是 RUNNING但锁已经过期超过 60 秒就判定 worker 失联把任务重新放回待派发队列并把旧执行记录标记为 FAILED。这里有一个价值很高的经验不要把“任务重新派发”做成全自动至少要有失败次数上限。比如同一个任务连续重新派发 3 次仍然失败就停止重试并发送告警通知值班同学介入。否则一旦任务命令本身有问题比如数据库连不上系统会陷入无限重试把日志表写爆。4.3 凌晨两点批量任务把调度线程拖垮当几百个任务同时到点调度器很容易出现“派发不过来的”现象。最开始的实现里调度循环一次捞 20 个到期任务派发完立刻捞下一批。听起来很正常但凌晨高峰期一个任务从捞出来到彻底派发完成可能经历几十次数据库查询数据库连接池很快就满了其他正常任务也跟着排队。我后面加了两层限制才稳定下来。第一层是速率限制调度器每秒最多派发 50 个任务避免瞬时压力打车数据库MAX_DISPATCH_RATE 50 # 每秒钟最多派发任务数第二层是优先级。任务表里加了一个 priority 字段P0 任务最优先派发。调度器捞任务时按priority DESC, next_run_time ASC排序这样即使高峰期出现排队核心任务也不会被边缘任务挡住。实际压测下来这两层限制的效果非常明显。凌晨批量任务把数据库请求打爆的问题从“经常发生”变成“几乎消失”。如果你也在做类似的调度系统我建议一开始就把派发速率和队列长度监控加上别等出事故再补。4.4 从日志到看板给调度系统留好可观测性调度系统不像普通接口问题不是实时暴露的很多故障要等业务方反馈“数据没更新”才发现。为了快速定位问题我后来在 task_log 旁边加了一张 event_log 表记录调度过程中的关键事件CREATE TABLE event_log ( id BIGSERIAL PRIMARY KEY, execution_id VARCHAR(64) NOT NULL, event_type VARCHAR(32) NOT NULL, message TEXT, created_at TIMESTAMPTZ NOT NULL DEFAULT now() );event_type 包括 CLAIM_SUCCESS、DISPATCH_SUCCESS、EXEC_START、EXEC_SUCCESS、EXEC_TIMEOUT、RETRY_SCHEDULED 等。每次任务异常只需要按 execution_id 过滤 event_log整个生命周期就一目了然。这个表我会长期保留因为它对复盘线上问题帮助巨大。可视化部分我不建议一开始就上复杂平台。先用 Grafana 连接数据库直接统计 task_log 里的成功率、平均耗时、失败 Top 10 任务就足够日常值班用了。告警也只做一件事任务最终失败并且重试次数用完时通过钉钉或企微 webhook 推一条消息。告警通道稳定比告警系统智能更重要。写到最后说点真心话。ax 这一年多下来我最深的体会是调度器的难点从来不在“如何到点触发”而在“触发之后能不能保证状态一致、故障能不能恢复得干净”。设计任务表、日志表、锁表这三张表时多花的一小时往往能省下将来排查问题的一整天。如果你也要做一个类似的调度系统先把最小闭环跑通再用真实任务去压它最后才谈得上扩展和优化。状态必须持久化执行必须幂等监控必须留痕——这三点做到位调度系统就基本立住了。
网站建设高端定制企业官网