恒美微站
首页
关于我们
建站服务
主题模板
案例展示
资讯中心
联系我们
自研轻量任务调度器adwawd:状态机与DAG编排实践
首页
资讯中心
/
自研轻量任务调度器adwawd:状态机与DAG编排实践
自研轻量任务调度器adwawd:状态机与DAG编排实践
发布时间:2026/9/10 7:00:25
1. 项目整体思路一个调度器的诞生如果你维护过超过几十条定时任务的调度系统八成会遇到这么几个尴尬时刻开源调度工具功能确实强但改一个细节需要通读几万行源码团队内部几十个业务方都来申请任务权限值班的同学光审批就累够呛最头疼的是跨部门调度A部门想触达B部门的数据两边为了任务优先级能拉扯半天。我最近把内部这套任务编排调度器重构后给它起了个名字叫“adwawd”。说白了就是一个自带依赖关系的任务调度和编排系统核心能力是让用户用一份精简的配置文件描述任务之间的上下游关系、执行时间、重试策略、资源限制然后交给调度中心去统一管理。整个项目没有依赖重型外部组件部署起来十几分钟就能跑通特别适合中小团队或者单部门内部使用也是我写这篇文章想介绍的核心内容。先说结论这套系统解决的核心问题有三个第一把任务编排从代码里抽出来业务方不用为了调整执行顺序去动核心代码第二把调度权限下放到各业务方谁的任务谁自己管结果第三通过可视化和状态机把整个执行过程变得可观测出问题能快速定位。如果你手头正好有类似的需求不管用的是 Python、Java 还是 Go这套设计思路都能直接参考。1.1 当时的痛点开源方案为什么不够用项目启动前我对市面上几类方案都做了评估。Airflow 名气最大DAG 概念成熟生态丰富但它在部署层面依赖 Python 环境和元数据库在团队运维基础偏弱的情况下光调通组件就花了不少时间。XXL-Job 在 Java 生态里很流行调度能力直观但它的定位更偏向定时任务调度对任务之间的依赖关系支持比较弱。DolphinScheduler 功能完整度很高但架构相对重审批流、租户体系这些在企业场景里有用对一个小团队来说反而成了负担。另一个关键触动点是任务配置的可移植性。原来用 shell 脚本串任务参数全靠环境变量新同学接手后很难快速看清全局。用现成调度工具如果想要一套配置搞定多环境测试、预发、生产迁移通常得开发额外的适配层。后来我想清楚了一件事与其在一堆开源框架之间做取舍不如直接抽一个符合团队习惯的轻量内核把核心调度模型做干净其他都交给约定。1.2 方案选型背后的取舍自研调度器最常见的质疑是“重复造轮子”但实际做下来我更愿意把它看作“造适合自己的轮子”。核心选型决策有这么几条语言用了 C 实现调度调度核心Python 作为用户的配置和扩展接口层为什么这么干后面会讲到任务消息分发用 Redis 的 Stream 结构天然支持消费组不会丢消息元数据选了 SQLite单机部署零配置后续数据量大可以平滑切到 MySQL。依赖管理这块我完全放弃了 DSL 式的 Crontab 表达式改成了基于时间戳的前驱后驱关系。用户定义每个任务在上游任务成功后再启动上游是多个任务的话可以做 AND 或 OR 的依赖组合。这个设计初期看起来实现更麻烦但使用体验比纯 Crontab 灵活得多特别是数据管道类场景时间不重要数据到达顺序才重要。在《人月神话》里 Brooks 讲过软件工程没有银弹。选型也一样没有哪个方案是完美适配所有场景的关键在于认清自己的边界。adwawd 的核心假设是绝大多数任务调度需求可以归约为“下游等待上游 重试 超时”没必要把所有业务规则都塞进调度器里。这个假设一直是我架构取舍的准绳。2. 核心机制调度器内部设计与关键细节一个调度器最核心的部分不是“到点触发任务”这种基础功能而是三个问题任务状态怎么管理、依赖关系怎么解析、执行结果怎么反馈。这一节我拆开讲。2.1 任务状态机为什么是七种状态状态机是整个系统的地基。adwawd 把任务生命周期定义为七种状态READY待调度、WAITING等待上游、RUNNING执行中、SUCCEEDED成功、FAILED失败、TIMEOUT超时、SKIPPED跳过。这几个状态覆盖了所有我需要考虑的情况。状态流转的核心规则是所有任务初始为READY当依赖的上游任务没有全部成功时任务进入WAITING满足依赖条件后调度器会尝试分发执行进入RUNNING执行成功进入SUCCEEDED失败进入FAILED超过设定阈值没完成就进入TIMEOUT如果上游被跳过比如依赖的数据源当日为空下游默认进SKIPPED。为什么这么设计很多系统把“等待依赖”和“等待执行”混在一起排查问题时很难分清一个任务到底是因为资源不够在排队还是因为上游没跑完。分状态以后看板上一眼就能看出瓶颈在哪里。在实现上每个状态变更都会写入一条变更流水这样不仅知道当前状态还能回溯整个状态变化历史。排查问题的时候这个历史会特别有用你能看到任务在哪个环节停了多久。2.2 依赖解析从 YAML 配置到有向无环图调度器内部把用户提交的配置解析成一个 DAG有向无环图。节点是任务边是依赖关系。我选择 YAML 而不是 JSON 作为配置格式主要因为 YAML 支持注释业务方维护长配置时可以写清楚每个任务的用途。依赖定义采用显式声明的方式每个任务必须声明depends_on字段列出所有上游任务名称。这个方案参考了 Makefile 的思路任务本身不带时间触发条件而是完全由上游驱动。比如一个数据同步管道A 任务从 API 拉数据B 任务做清洗C 任务写入数仓配置就是dag_name: api_to_warehouse tasks: - name: extract_api_data worker_type: python command: python tasks/extract.py retry: 3 timeout: 1800 depends_on: [] - name: clean_data worker_type: python command: python tasks/clean.py retry: 2 timeout: 1200 depends_on: - extract_api_data - name: load_to_warehouse worker_type: python command: python tasks/load.py retry: 1 timeout: 3600 depends_on: - clean_data在解析阶段调度器会做一次环路检测。实现时用拓扑排序如果有环配置直接拒绝提交并给出具体的环路路径。这一步拦截了大量配置错误。我试过不检查环的版本结果任务永远卡在WAITING状态排查起来很难受。有了环检测至少能在提交阶段就把问题挡住。解析完成的 DAG 会序列化成一份不可变的调度文件后面执行引擎只读取这份文件不关心原始配置。这样做的好处是配置改了以后历史已提交的任务不受影响版本可回溯。2.3 调度算法如何避免任务饿死调度核心需要考虑一个公平性问题如果两个 DAG 都包含耗时很长的任务是否会产生一个任务把资源全部吃掉的极端情况。因为执行节点数量有限可同时运行的任务并发数有上限如果没有合理的算法后提交的任务可能一直得不到执行机会。adwawd 的调度器维护一个优先队列排序键是任务提交时间加 DAG 权重。DAG 权重是用户在配置里可选的一个参数不配置时默认权重相同提交时间早的先执行。如果配置了权重那么高权重的 DAG 里的任务会被优先调度。实测下来这个简单的两级调度策略已经能满足绝大多数场景既保证了公平又允许业务方做资源倾斜。真正的难点在于调度频率。调度器每 5 秒扫描一次待调度任务表然后批量触发可执行任务。每次扫描都会把 WAITING 的任务拉起来检查依赖这个操作在任务量大的时候容易变成瓶颈。目前通过两个索引优化解决——DAG 状态索引和任务依赖完成时间索引。在任务量万级以内调度延迟可以控制在 5 到 10 秒内这对绝大多数调度场景完全够用。2.4 执行节点任务真正运行的地方任务分配器把任务投递给空闲的执行节点。每个执行节点注册到 Redis Stream 的消费组里通过竞争方式获取任务但这不是简单的抢占式节点会先上报自己的空闲槽位数量调度器根据槽位情况决定是否投递。这个设计避免了一个节点堆积大量任务而其他节点空闲。执行节点拿到任务以后会以子进程方式运行用户的命令进程级别加超时控制。这个超时不是简单的在外部算一个时间点然后 kill而是通过 Linux 的waitpid加上可配置的SIGTERM,SIGKILL两级终止机制实现的。先允许跑一段时间不等超过 30 秒还退出不了再强杀尽量给业务方一个优雅收尾的机会。每个任务执行都会保存三份关键日志规格化标准输出、规格化错误输出、执行摘要。这三份日志对排查问题特别重要。标准输出和错误输出用来给业务方看执行细节执行摘要用来给调度器判断结果。如果命令返回非零退出码任务直接判失败如果返回零但输出中超时关键字也判失败。这就是为什么状态机需要独立设计因为“跑完”和“成功”完全是两码事。3. 实操过程从零搭建可用的 adwawd 实例纸上谈兵到这里接下来是实战部分。我默认你用的是 Linux 环境Python 版本在 3.9 以上Redis 6.0 以上机器能联网安装依赖。整个搭建过程按步骤来每一步都会说明命令的含义和可能遇到的坑。3.1 部署装依赖、初始化、启动服务先拉代码安装 Python 依赖。项目核心依赖很轻只需要pyyaml、redis、sqlalchemy和click这几个库没有其他重量级第三方包。然后初始化元数据库默认生成一个adwawd.sqlite3文件里面会建好 DAG 表、任务表、日志表这些基础表。git clone gitgithub.com:your-repo/adwawd.git cd adwawd python -m pip install -r requirements.txt python manage.py init-db接着启动三个服务调度器、执行节点、Web 控制台。调度器负责扫描待调度任务执行节点负责跑命令Web 控制台负责展示状态和提供操作入口。三个服务可以部署在同一台机器上也可以分开部署区别只在于配置里的 Redis 地址是否指向同一个实例。# 启动调度器 python scheduler.py start # 启动两个执行节点每个节点默认 4 个并发槽位 python worker.py start --name worker-1 --slots 4 python worker.py start --name worker-2 --slots 4 # 启动 Web 控制台 python web.py start --port 8080启动完以后打开浏览器访问http://localhost:8080能看到一个空的控制台页面。到这里基础设施已经就绪下一步就是提交一个实际的 DAG。实际部署中我踩过的最大的一个坑是执行节点并发数设置。早期版本默认并发数等于 CPU 核心数结果一个节点同时跑 16 个 Python 任务CPU 直接打满任务互相拖累。建议是 IO 密集型的任务并发数可以高一些比如两倍核心数CPU 密集型的任务并发数最好不要超过核心数。节点槽位设置的逻辑后面会讲。3.2 编写 DAG 配置一个真实的跨表同步场景我拿一个比较常见的场景练手每天凌晨从一个订单库同步数据到分析库中间做一个简单的数据校验。这个场景在电商公司很典型涉及三个任务抽取订单表、校验数据完整性、写入分析库。dag_name: order_sync schedule: 0 2 * * * description: 每日订单数据同步到分析库 tasks: - name: extract_order worker_type: shell command: /opt/scripts/extract_orders.sh --date {{ ds }} retry: 3 timeout: 2400 depends_on: [] resource: cpu: 1.0 memory: 1024 - name: validate_orders worker_type: shell command: /opt/scripts/validate_orders.py --date {{ ds }} retry: 0 timeout: 600 depends_on: - extract_order resource: cpu: 0.5 memory: 512 - name: load_to_warehouse worker_type: shell command: /opt/scripts/load_orders.py --date {{ ds }} retry: 2 timeout: 3600 depends_on: - extract_order - validate_orders resource: cpu: 1.0 memory: 2048注意load_to_warehouse这里同时依赖extract_order和validate_orders。按照 DAG 的定义只有当两个上游都成功以后它才会进入待调度状态。这个配置体现了 AND 语义。提交 DAG 的命令很简单控制台提交或者命令行提交都可以python manage.py submit-dag --file examples/order_sync.yaml提交成功后控制台的 DAG 列表里会出现order_sync点击进去能看到三任务形成的菱形结构节点分别为待调度状态。如果extract_order还没执行validate_orders和load_to_warehouse会保持 WAITING 状态这完全符合预期。这里有个关于占位符的细节。配置里我用了{{ ds }}这是内置的日期变量在执行时会自动替换为任务实际执行日格式为YYYY-MM-DD。调度器为每个任务生成独立的执行上下文所以即使在同一个 DAG 里有两个任务使用{{ ds }}它们拿到的日期也是同一天但如果重跑历史任务日期会变成重跑指定的日期。这个设计保证任务可以在任意时间点回放数据对修数据特别重要。3.3 运行、重试与参数调优DAG 提交后调度器会根据schedule字段的 Crontab 表达式自动触发。如果不想等控制台里也可以手动触发一次。手动触发和自动触发在内部走同一条路径都会生成一次dag_run记录这个设计让排障更简单因为不用区分任务是被谁触发的。运行后关注的第一个指标是调度延迟就是从任务变为 RUNNING 状态到实际开始执行的时间差。正常情况下这个差值在 5 秒以内。如果超过 10 秒说明执行节点的槽位不够。想要验证这个判断可以在控制台查看节点的空闲槽位数量直接增加节点或者调大 slot 数就能解决。排查问题的时候任务级别的时间线视图最有效。它按时间顺序列出该任务从提交到结束的所有关键事件包括进入 WAITING 的时间、被调度器扫描到的时间、投递到哪个节点、子进程启动时间、退出码、结束时间。有一次我发现一个任务总是比预期晚 20 分钟开始打开时间线才发现是上游任务重试了 3 次每次都等到超时阈值边缘才失败下游自然被拖慢了。遇到失败任务点击日志按钮能直接看到退出码和最后 50 行输出不需要登录服务器翻日志这对研发效率提升很明显。控制台还提供“重新执行”按钮点击后该任务以及所有下游任务会被标记为待调度从失败点继续跑不会把整个 DAG 重跑一遍。这个能力在数据管道场景里极其重要因为上游数据往往很大重跑全链路成本太高。3.4 用代码自定义执行器如果内置的 shell 和 python 执行器满足不了需求项目还留了执行器接口。执行器的核心协议是接收一个任务上下文返回执行结果。上下文里包含配置参数、执行日期、上游产出物路径。结果里包含退出码、日志输出、自定义指标。# custom_executor.py import time from adwawd.executors.base import BaseExecutor from adwawd.executors.context import TaskContext class CustomExecutor(BaseExecutor): def run(self, context: TaskContext): start_time time.time() # 这里写自定义逻辑 result self.execute_with_retry(context) elapsed time.time() - start_time return {exit_code: 0, elapsed: elapsed, message: ok}这个接口的扩展性很强。比如我后来接了一个 Spark 任务的提交逻辑只需要在自定义执行器里调spark-submit命令再把 Spark 的 application_id 写入结果里控制台展示的时候就能直接跳转到 Spark UI 看执行详情。这套“执行器即适配器”的思路让 adwawd 从纯命令行调度器扩展成了数据处理的中枢。4. 常见问题与排查技巧实录这块我整理了项目运行期遇到的高频问题每一类都给出当时排查的过程和最终的解决方案希望能帮你少踩一些坑。4.1 任务显示失败但实际跑成功了听起来很荒谬但确实发生了。有一个 Python 任务业务脚本本身可以正常完成但 adwawd 却把它标记为 FAILED。打开日志一看退出码是 0但在标准错误里输出了一段警告里面包含了Error关键字。旧版本执行器里有一个“错误输出检测”的逻辑只要扫描到关键字就认为任务失败。这个逻辑本意是抓一些退出码正常但实际报错的场景但误伤了不少正常任务。这个问题的根本结论是退出码 0 就是成功标准错误里的内容只是参考信息。后续版本加了一个配置项默认关闭错误输出检测。这个改动直接让任务失败率降了一个数量级。排查这个问题时的启示是做调度系统时“成功”这个定义越简单越好不要自作聪明地加太多启发式判断。4.2 调度延迟越来越高最终完全卡住上线一个月后的一天我收到告警说大量任务积压调度延迟从 5 秒涨到了几个小时。打开调度器的日志发现扫描一次待调度任务表的时间越来越长。一开始以为是任务数量增长导致的但看监控发现任务量只涨了一倍不至于性能恶化到这种程度。后来定位到是 Redis Stream 的消费者组堆积问题。一个执行节点因为网络问题断连重连但消费组的 ack 机制没有正确处理导致消息不断堆积。这些积压消息占用了调度器的心跳通道形成了一个负反馈循环调度器越慢任务堆积越多任务堆积越多调度器越慢。解决方法是升级了 Redis Stream 的消费逻辑重连后先重新确认所有未 ack 的消息再继续消费新消息。同时在调度器侧加了健康检查如果发现连续三次扫描耗时超过阈值自动把调度器切换成“降级模式”——优先处理存量任务的超时检测暂停新任务分发。这个设计虽然不是完美方案但在极端负载下能避免系统完全不可用。4.3 多个任务同时修改同一份数据文件这是资源控制不到位导致的。两个 DAG 都指向同一个下游文件一个任务写入数据另一个任务做清理两个任务并行执行时清理任务把写入任务的数据删了一半导致数据损坏。这个问题在调度系统里非常典型本质上是资源读写冲突。解决方案是引入资源锁机制。用户可以在配置里给任务声明resource_lock比如order_sync:write。调度器在执行前会先请求锁持有锁期间其他任务即使调度到同一节点也会等待。锁的粒度尽可能细不要锁大范围资源否则反而会引入额外的排队延迟。4.4 节点失联但任务还显示 RUNNING有一次一个执行节点因为宿主机 OOM 被操作系统杀掉了但调度器并没有立刻感知到。该节点上的任务就一直显示 RUNNING下游任务全部卡在 WAITING。排查发现调度器对执行节点的心跳超时阈值设的是 60 秒而系统 OOM 到调度器隔离节点的延迟远超过这个时间。解决办法是缩短心跳超时到 15 秒同时增加一个“节点失联自动重分配”的机制。如果失联节点上有 R 任务调度器会在 30 秒后把它们重新标记为 READY允许其他节点重新执行。为了让重分配过程安全任务执行本身必须保证幂等——重复执行不会产生副作用。如果你的任务不具备幂等性调度器做得再完善也无法避免数据问题。4.5 问题与定位速查表现象优先排查方向常见根因任务一直 WAITING检查上游任务失败状态上游还在重试或依赖配置失配调度延迟高查看节点槽位占用执行节点并发数不足心跳中断频繁检查宿主机负载、网络OOM 或网络抖动的节点被隔离重复执行后有脏数据业务任务非幂等未配置资源锁或幂等策略配置提交报错查看环路检测提示DAG 成环或依赖指向不存在的任务5. 影响范围扩展从调度器到数据工作流平台adwawd 做的时间久了我发现它能承载的已经不只是“定时跑任务”这个职能。当任务配置、执行日志、资源控制、失败重试都标准化以后调度器实际上成了一个数据工作流平台的基础设施。往下可以接不同语言写的业务脚本往上可以接分析人员的临时跑数需求中间是统一的调度控制和监控面板。后来我把几个常用的数据同步工具DataX、Sqoop封装成了标准执行器业务方只需要在 YAML 里声明数据源和目标表不用关心工具内部怎么调用。封装过程让我体会到调度系统做好以后最大的价值不在于本身的功能而在于它把底层的复杂度包装掉了让业务方专注于业务逻辑本身。5.1 与监控告警体系的对接调度器只负责执行是远远不够的真正要靠的是一个可观测体系。adwawd 把任务的元数据和执行日志统一以结构化格式落库暴露了一套简单但完整的指标单位时间调度任务数、成功数、失败数、平均调度延迟、任务积压量。这套指标可以直接接入 Prometheus配合 Grafana 做可视化。刚开始没有做监控告警时任务失败只能靠值班同学刷控制台发现。接入告警以后失败任务会在 1 分钟内通过 Webhook 推送到即时通讯群带上失败原因、日志链接、耗时数据。这个改动对团队协作的提升非常明显因为问题反馈从被动值班变成主动感知了。如果你对监控体系不熟悉建议从最基础的失败率和延迟两个指标开始后续再逐步加长尾指标。5.2 跨团队协作模式另一个有意思的方向是跨团队协作。原来部门 A 要跑数据需要找部门 B 的研发开发任务再约定触发方式。引入 adwawd 以后每个团队有自己的命名空间可以管理自己团队的任务同时团队之间可以公开或授权执行某些 DAG。这意味着部门 A 可以直接把部门 B 发布的某个 DAG 作为自己 DAG 的上游依赖当别人的任务跑完自己的任务会自动触发。这种模式下团队之间的连接从“人工通知 手动触发”变成了“系统编排 自动执行”。任务之间的数据契约由双方通过 YAML 配置定义实际上形成了一套轻量的接口协议。我在实际使用中看到这种模式能让团队协作效率提升一个量级但也对任务命名的规范性提出了要求。约定大于配置一个好的命名和注释习惯可以避免大量沟通成本。5.3 权限与资源水位管理权限设计上adwawd 做了三层角色管理员、维护者、查看者。管理员管理集群和用户授权维护者可以提交和修改 DAG控制任务重试和取消查看者只能看状态和日志。这个权限模型足够应对绝大多数内部工具的场景。资源管理方面除了前面提到的节点并发槽位DAG 维度还可以配置最大并行任务数。有的 DAG 有几十个任务但业务方希望同时最多跑 5 个避免对数据库造成太大压力。这个配置在参与实时数据同步的场景里特别重要因为目标库的写入能力通常有限盲目并行会把数据库压垮。6. 延伸思考与实操心得把一个调度器从零到一落地再到稳定服务几万次任务执行中间的过程让我对任务调度这件事有了更切身的理解。最后分享几个我个人的实操心得希望有点参考价值。关于技术选型我见过很多团队一上来就选重型框架觉得功能越多越好。但调度系统这种基础组件恰恰应该遵守最小可用原则——先满足核心场景留好扩展点不要奢求一步到位。adwawd 到今天也只有九个核心配置项大部分用户十分钟就能上手这种“心里有数”的感觉让使用者敢于改配置、试新任务系统也就更容易推广起来。关于状态设计把任务状态切细看起来会增加实现复杂度但它让系统具备了可解释性。很多分布式系统最大的问题不是因为逻辑复杂而是因为出问题时无法说清当前处于什么阶段。状态机是基础状态变更历史是宝藏。每次排查完问题我都会回看状态流转记录很多时候瓶颈原因就藏在某两个状态之间的时间差里。关于稳定性调度系统直接决定上层业务任务的执行所以宁可保守也不能激进。新增功能时先在小范围灰度跑一段时间确认没副作用再全量开放。我自己因为在没有充分验证的情况下发布了一个调度算法的改动导致高峰期出现大面积任务堆积那次事故让我记住了一个铁律调度器的改动必须做流量压测和回放验证并且要准备兜底方案。关于后续扩展目前我在研究的方向有两个一个是把 DAG 的配置文件做成 Git 友好模式这样任务版本变更可以直接走代码评审流程另一个是把执行节点做成容器化方式每次任务运行时动态拉起一个轻量容器来做资源隔离这样能彻底解决不同任务环境依赖冲突的问题代价是部署复杂度会高不少。如果你有类似的调度场景我建议先从任务编排和可观测性入手这两个点带来的收益最直接也最容易让团队接受一套新系统。最后说一点所谓“调度”本质上是在回答一个非常朴素的问题什么时候、在什么条件下、去做哪件事。把这个道理想明白了无论用什么框架或语言实现都能设计出靠谱的系统。adwawd 只是这个朴素问题的一个落地方案我更希望分享的是这套思考方式——先想清楚状态、依赖、可控性再动手写代码。