简介:这份PDF完整收录了快手大数据任务调度系统Kwaiflow的设计与实践分享,适合大数据平台工程师、数据架构师以及从事调度系统研发的读者。内容从调度系统分类切入,梳理了从Airflow到Kwaiflow 3.0的演进路径,重点解析双层实体调度模型、Scheduler主备高可用方案、分布式秒级调度关键技术,以及容器化执行在减少调度延迟方面的实践细节。资源包含1个PDF文件,压缩包大小6.1MB,已有282人学习浏览。作者结合快手数据工厂真实场景,讲清了数十万任务、百万依赖交织下如何实现高容量、高可用调度,对想构建大规模工作流调度系统或优化现有任务的团队具有直接参考价值。
1. 快手大数据任务调度系统:从 Airflow 的 P99 数分钟到 Kwaiflow 的秒级调度
快手大数据任务调度系统的演进路径在业内很有代表性:2016 到 2019 年用 Airflow,任务数千个,没有外部平台接入;2019 年换 Kwaiflow 1.0 解决百万级规模问题;到 2021 年的 Kwaiflow 3.0,任务数涨到数十万,接入平台数十个。Airflow 单集群不超过 1 万 DAG、调度延迟 P99 高达数分钟、Scheduler 无 HA、年均故障约 8 个,这四条痛点基本说清了为什么要自研。下面按调度模型、低延迟设计、高可用和落地避坑四个方向拆解,适合正在选型或自研调度系统的后端、数据平台工程师。最值得借鉴的不是某段代码,而是“双层实体模型 + 事件驱动 Actor + 分级调度”这套组合设计思路。
2. 双层实体模型与系统架构:Task/DAG 怎么组织,Scheduler 怎么分工
Kwaiflow 的调度模型是双层实体模型,这是理解整个系统的钥匙。对比 Airflow 那种偏单层的 DAG 模型,Kwaiflow 把“执行模板”和“工作流集合”拆开:Task 是执行模板,承载某类代码的执行;DAG 是一系列 Task 的集合,挂着调度定时、依赖关系等属性。更关键的是 DAG 之间可以互相依赖,形成 DAG 级别的上下游链路,而不只是 Task 与 Task 之间的边。这个看似微小的差异,在快手这种多业务方、多任务类型、多调度周期混合的场景里,直接决定了系统能不能撑住几十万任务和百万条依赖关系。
2.1 双层模型为什么比单层模型能打
先看快手数据工厂里的真实任务类型:MySQL2Hive、Kafka2Hive、Bash、Hive SQL、指标生产、ABTest 任务、机器学习训练、Hive2Druid、Hive2Ch、Hive2Redis,既有亚秒级的 Bash 脚本,也有跑几十分钟的 Hive 任务,还有按小时、按天不同节奏触发的指标生产链路。如果所有任务平铺在一个 DAG 里,任务一多就变成一张大饼,权限隔离、定时配置、复用抽取都很难做。
双层模型的表达力来自“DAG 作为整体参与依赖”。比如指标生产链路是一个 DAG,机器学习训练任务依赖指标产出的结果,那就在 DAG 级别建立依赖,而不是让训练任务去感知指标 DAG 内部几十个 Task 的细节。数据开发只需要声明“我等 DAG B 完成”,DAG B 内部怎么编排由对应的业务方自己管。单层模型表达这种关系通常要把上游所有 Task 都列成依赖,维护成本随 Task 数量线性上升。
我一般会把两层模型的差异归纳成表格来看:
| 对比项 | 双层实体模型(Task + DAG) | 单层实体模型(仅 Task) |
|---|---|---|
| 依赖粒度 | Task 间依赖 + DAG 间依赖 | 只有 Task 级依赖 |
| 表达力 | 强,DAG 可作为整体被依赖,适合多业务方协作 | 弱,跨业务方协作要展开到 Task 级 |
| 定时属性挂载 | 挂在 DAG 上,DAG 内部 Task 可复用 | 每个 DAG 都要单独维护 |
| 权限与控制粒度 | 可按 DAG 隔离业务方权限 | 权限只能做在 Task 或目录上 |
| 系统复杂度 | 高,需要处理双层状态机 | 低,上手快 |
实现一个 Task 模板时,我习惯把执行类型、资源规格、重试策略、超时时间这些抽成字段。类似下面的示意配置:
task_template: task_id: "hive_sync_orders" type: "HiveSQL" owner: "data_dev_team" schedule: "dag_level" resources: cpu: 2 memory: "4G" retry: times: 3 interval: "5m" timeout: "60m" script_path: "hdfs://warehouse/scripts/sync_orders.hql"上面的配置里,type 决定了这个 Task 用哪类执行器,HiveSQL 会走 Hive 的执行通道;resources 是资源规格,快手实际环境里有 0.5C1G、1C1G、16C32G 这类档位,不同规格对应不同容器配额;retry 是任务失败后的重试策略,重试次数和间隔要按任务类型区分,Bash 任务我通常不重试,Hive 任务看情况给 1 到 3 次;timeout 防止任务卡死占用 Worker。注意这里的 schedule 字段声明为 dag_level,意思是这个 Task 自己不挂定时,完全由所属 DAG 的调度属性控制。
2.2 系统架构:Scheduler、Queue Service、Worker 三层怎么分工
Kwaiflow 整体架构可以分成接入层、核心模块和其他模块来看。接入层主要是 API Server,统一接收所有外部请求,比如创建 DAG、发布任务、查询实例状态;核心模块是 Scheduler、Queue Service、Worker;其他模块包括 Log Server 日志服务、Alarm Server 报警服务、Event Emitter 事件发送、Instance Lineage 实例血缘。
Scheduler 是大脑,负责实例生成与调度。它内部有 Routine Scheduler 的 Active 和 Standby 两套角色,Active 承担实际调度,Standby 热备;还拆了 DAG Loader、TimeDetector、Dep.Detector、Resource Detector 等不同职责的调度器组件。Queue Service 是实例分发层,不是简单的一个队列,而是分通道的:P0 Channel、Px Channel、K8s Channel、Biz Channel,不同优先级和不同类型的任务走不同通道,避免大任务把通道堵死。Worker 负责实例执行,内部再分 Worker Group,每个组里有 Local Executor 和 Remote Executor 两种执行方式。
一条任务从发布到执行,大致链路是这样的:API Server 接收创建请求,Scheduler 里的 DAG Loader 加载 DAG 定义,TimeDetector 根据定时配置生成实例,Dep.Detector 检测上游依赖是否满足,满足后把实例交给 Queue Service 分发到对应的 Worker Group,Worker 上的 Executor 拉取用户代码并执行,执行结果再回报给 Scheduler 更新状态。整个过程如果走轮询,几十万任务每轮都要扫一遍状态,延迟和数据库压力都扛不住,所以快手在 Scheduler 内部用了 Actor 模型做全事件触发,这块后面细说。
2.3 接入层和其他模块的边界
接入层里除了 API Server,还有 Log Server、Alarm Server、Event Emitter、Instance Lineage。Log Server 负责统一收任务运行日志,避免用户直接登录 Worker 查日志;Alarm Server 根据监控规则产生报警;Event Emitter 给外部系统发标准事件,比如实例成功、失败、重试,方便上层数据平台订阅;Instance Lineage 记录实例血缘,这个对数据质量追踪特别有用,任务跑挂了能快速定位影响面。
把这些模块拆出核心链路之外,是快手这套设计里很务实的一点。Scheduler 只管调度,日志、报警、血缘都走旁路模块,核心链路的压力就小很多。实际拆项目时我特别留意了这个边界:如果日志和报警都塞进 Scheduler,任务高峰时日志 IO 会直接把调度线程拖死,这是很多自研调度系统翻车的常见原因。
3. 低调度延迟的关键设计:Actor 无轮询链路与百万级定时器
调度延迟的定义是理论起调时刻到实际开始运行用户代码的时间差。注意它不包含用户代码本身的运行时间,只看“到点该跑,到底多久真的跑起来”。Kwaiflow 把这段延迟拆成了四个主要来源:状态转变、运行环境准备、定时器触发、数据库访问。状态转变可以做毫秒级,运行环境准备在本地执行时是毫秒到亚秒、容器化执行时是秒级,定时器触发在传统实现里是亚秒到分钟,数据库访问是毫秒级。目标的秒级调度延迟,意味着这四个环节每一个都不能拖后腿。
| 延迟来源 | 常规量级 | Kwaiflow 的做法 |
|---|---|---|
| 状态转变 | 毫秒 | 全链路事件触发,无轮询 |
| 运行环境准备 | 毫秒~亚秒(本地)/ 秒级(容器预热) | 预加载、镜像预热 |
| 定时器触发 | 亚秒~分钟 | 定时器索引、读写分离、分库分表 |
| 数据库访问 | 毫秒 | 读写分离、分库分表 |
3.1 Actor 事件触发链路:从 EntryActor 到 Result Handler
Scheduler 内部把整个调度流程拆成 19 个 Actor 步骤,全部事件触发,没有轮询。大体链路是这样:EntryActor 接收 API Server 发来的请求,DagLoaderActors 加载 DAG 定义,Inst. Generator Actors 生成实例,Dep. Detector Actor 检测上游依赖状态,Dep. Decider Actor 决策依赖是否满足,Resource Detector Actor 检测资源,Resource Decider Actor 决策资源是否够,Retry Actor 处理重试,Inst. Receiver Actor 把实例交给渲染环节,Render Actor 渲染参数(比如业务日期、动态变量),Prehook Actor 跑执行前钩子,Execute Actor 真正执行,Posthook Actor 跑执行后钩子,Processor Actor 汇总结果,Result Handler Actor 处理结果并写回状态。
这套 Actor 设计的好处是三个:全流程事件触发,从任务发布到执行落地没有一次轮询,状态一变立刻有人响应;异步高并发,Actor 之间通过消息传递,单次操作不会阻塞后续任务;简单易用,不需要手动加锁或者管理线程池,Actor 框架把并发控制包掉了。我拆过不少调度系统,很多系统延迟高不是业务逻辑复杂,而是每隔几秒扫一次数据库判断有没有到期任务,任务一多扫描本身就是瓶颈。Kwaiflow 用事件驱动把“状态转变”压到毫秒级,等于把调度延迟的地板直接拉低了。
3.2 百万级高吞吐定时器的实现思路
几十万任务挂在系统里,每个 DAG 有自己的定时规则,到点就该触发实例生成。如果用一个大的延迟队列或者每分钟扫一次数据库,延迟和数据库压力都不可控。Kwaiflow 的做法是定时器索引、读写分离、分库分表三件套。
定时器索引的意思是,不直接扫任务表,而是维护一张专门用于扫描的索引表,只记录任务 ID、下次触发时间、触发周期。扫描窗口只覆盖近几分钟的数据,索引命中后精准触发,避免全表扫描。读写分离是把实例生成这类写操作和依赖查询这类读操作拆到不同的库,读库可以做只读副本。分库分表针对的是任务量再往上走的阶段,按 DAG ID 哈希分片,把单库的压力摊开。
我做容量评估时习惯先算一笔账:假设 50 万任务,每分钟需要判断到期的任务约占 1%,也就是 5000 个,单库单表扫 5000 行不是问题;但如果每次都全表扫 50 万行,延迟就会从毫秒级退化到秒级。所以定时器索引的价值不在于省那几千行,而在于把“全部任务都背上”变成“只扫即将到期的少量任务”。
3.3 运行环境准备:本地执行与容器化执行怎么选
状态转变做到毫秒级之后,环境准备就成了主要矛盾。Kwaiflow 把执行分成了两种方式:本地执行和容器化远程执行。本地执行直接跑在 Worker 的进程里,启动快、Overhead 小,适合 Hive、Sensor 这类需要频繁调度但对隔离要求不高的任务;容器化执行把用户脚本丢进容器跑,资源和环境隔离更好,适合 Bash 这类自定义脚本,但镜像拉取和容器启动是分钟级的开销。
秒级延迟目标下,分钟级的镜像拉取显然不合格,所以快手做了镜像预热。通用镜像提前加载到 Worker 节点上,任务执行时直接从本地启动容器,把分钟级压到秒级。我自己的经验是,预热策略不能只在集群初始化时做一遍,还要监听 Worker 节点的上下线事件,新节点一注册就强制触发热镜像拉取,否则高峰期扩容的节点会集体卡在拉镜像上,那画面相当酸爽。
4. 高可用与分级调度:Exactly Once、主备切换和资源争抢
任务调度系统的高可用比普通微服务难做,因为调度的是别人的任务,漏跑会导致下游数据缺,重跑会导致重复数据写库。快手的思路是两条腿走路:一条腿是组件层面做到“尽量 Exactly Once”,另一条腿是业务层面用分级调度保障高优链路。组件故障、组件失联这类场景,在分布式系统里永远存在,设计目标不是杜绝故障,而是故障时任务不漏不重。
4.1 尽量 Exactly Once:故障转移、消息 Ack 和状态机
任务实例的可靠执行面临两个典型场景:组件故障,比如 Scheduler crash;组件失联,比如组件进程还在运行但网络不通。Kwaiflow 对应用了两套手段。组件故障时做 Failover 自动故障转移,配合鲁棒通信协议,Active Scheduler 挂了 Standby 能接管;组件失联时靠消息 Ack 和定期轮检,Ack 保证消息至少送达一次,轮检用来发现失联的组件并把任务重新分发。
避免遗漏执行和避免重复执行要同时解决。避免遗漏靠消息 Ack 加定期轮检,消息没确认就继续投递;避免重复靠状态机和重试清理。状态机记录每个实例的完整状态流转,比如 READY、RUNNING、SUCCESS、FAILED,只有特定状态才能执行特定操作;重试清理会把标记为 RUNNING 但实际已经结束的实例清掉,防止重复消费。
需要说清楚的是,分布式系统里的 Exactly Once 通常是用“at least once + 幂等/清理”逼近的,Kwaiflow 也不例外。实例 ID 是幂等键,结果写入前先查状态,重复的消息直接丢弃。我在实际工作中验证这类系统时,不会只看正常流程,而是专门做故障注入:kill 掉 Scheduler,看任务会不会漏;网络抖动,看会不会重。只有这两种注入都过了,才敢说高可用基本靠谱。
4.2 主备切换:Routine Scheduler 的 Active 与 Standby
Scheduler 的高可用通过 Active/Standby 实现,平时 Active 承担调度,Standby 热备。切换时最关键的不是谁能当主,而是状态怎么交接。如果 Active 有一批已经调度但还没确认执行结果的消息,Standby 接管后要把这批消息重新处理;但如果这批消息里的任务已经跑起来了,直接重放就会导致重复执行。
所以主备切换必须配合“状态 + 幂等”一起做。状态是每个实例当前到哪一步了,幂等是同一个实例 ID 只接受一次最终结果。我拆 Kwaiflow 的 Failover 设计时,特别注意到了一个细节:它强调鲁棒通信协议,意思是消息发送方和接收方之间对“什么算成功”有明确的共识,比如先写状态再回 Ack,还是先回 Ack 再写状态,顺序不同,故障时的表现完全不同。这个顺序问题,是很多自研系统主备切换后数据不一致的根本原因。
4.3 分级调度:资源紧张时保高优链路
快手有大促、大型活动这类场景,数据量突然翻倍,计算资源不能让所有任务同时按时产出。分级调度的思路是:把任务和 Worker Group 都分成 P0 到 P3 等级,高优任务优先调度优先执行,低优任务限流,同时允许人工管控。
具体做法是依赖规则管理器。Rule Manager 维护 Allow Rules 和 Block Rules,下发到 Scheduler;Resource Manager 感知底层资源情况,把 Yarn、K8s 的资源状态和反压信息上报给 Scheduler。调度时先看优先级再看资源,P0 通道的任务来了优先分配资源,P3 任务在资源紧张时主动限流。人工管控是最后一道保险,可以手动放行某些任务,也可以阻断某些任务,适合数据出问题时的紧急止血。
| Worker Group 等级 | 典型任务 | 资源紧张时的行为 |
|---|---|---|
| P0 | 核心指标生产、在线服务依赖 | 优先调度、优先执行 |
| P1 | 常规小时级数仓任务 | 正常调度,不主动抢占 |
| P2 | 天级任务、ABTest 数据 | 可被低优先级任务影响 |
| P3 | 探索性分析、模型实验 | 资源紧张时限流 |
这套分级设计里,我认为最值得抄的是“规则下发”而不是“写死在代码里”。规则管理器独立于 Scheduler,意味着运维人员不需要改代码就能调整调度策略。大促前把活动相关任务的优先级调高,活动结束后把规则删掉,全部通过配置完成。
5. 常见问题与避坑:Kwaiflow 落地时最容易被忽略的四个细节
自研调度系统也好,深度定制开源调度系统也好,很多问题不是架构设计阶段暴露的,而是任务量上来之后才翻车。下面这几条是我拆这类系统时认为最常见的坑,按“现象 → 原因 → 解决”来写。
5.1 主备切换后任务重复执行
现象:凌晨发生了 Scheduler 主备切换,第二天发现部分小时级任务同一个实例跑了两遍,下游表出现重复数据。
原因:Failover 时 Standby 把一批状态为“已发送未确认”的消息重新回放,但消息里的任务实际上已经执行完了,状态还没写回;重复回放导致同一个实例 ID 被两个进程同时消费。
解决:状态机必须做持久化,且每个 Actor 步骤完成后先写状态再回 Ack;重放消息前先做一次孤儿实例清理,把标记为 RUNNING 但已经失联的实例重置;执行结果写入做成幂等,同一个实例 ID 只接受第一次写入。我在自研系统里会把这三件事绑定成一条开关,缺一个都不允许 Failover 上线。
5.2 镜像预热像玄学:预热了还是慢
现象:镜像列表里明明已经预热的镜像,高峰期新扩容的 Worker 节点执行任务还是要等几十秒,甚至超过一分钟。
原因:新扩的 Worker 节点是后注册的,预热只覆盖了存量节点,没覆盖新节点;还有一种情况是预热用的镜像 tag 和任务执行时用的 tag 不一致,tag 指向的镜像内容已经变了,等于没预热。
解决:镜像预热必须和节点上下线联动,新节点注册后强制触发一次热镜像拉取,拉完才接收任务;执行端统一改用镜像 digest 而不是 tag,避免 tag 漂移导致预热失效。从那以后我每次设计镜像预热方案,都会把“节点生命周期事件”和“镜像不可变”两条写进验收标准。
5.3 分级调度规则翻车:宽泛的 Block 规则把高优任务拦了
现象:大促期间 P0 任务被 Block 规则拦截,P3 任务反而正常执行,报警群里炸锅。
原因:Allow 和 Block 规则没有做优先级排序,一条写得很宽泛的 Block 规则先匹配到所有任务,高优任务也被拦了;或者有人工放行规则,但规则没设置有效期,活动结束后还在生效。
解决:规则下发前先做影响面模拟,看看会匹配到哪些任务,再实际下发;Block 规则必须比 Allow 规则更具体,冲突时 Allow 优先;所有人工放行规则强制带过期时间,到期自动失效。我一般会在规则管理器里加一个“模拟执行”按钮,上线规则先跑一遍模拟,匹配量级确认没问题再正式生效。
5.4 分库分表后依赖判断从毫秒退化到秒级
现象:任务量到几十万之后,依赖判断延迟明显上升,数据库连接数经常被打满。
原因:依赖边分散到多个分片后,判断一个 DAG 的上游是否完成需要跨分片查多个表,JOIN 和多次查询把数据库拖垮了。
解决:把全量依赖关系从数据库搬到内存,启动时加载一次依赖图,后续依赖判断全部查内存;数据库只做变更日志,依赖关系有变动时异步更新内存图;定时器索引单独分库,读写分离,扫描窗口只扫即将到期的任务。这套方案在任务量进一步上涨时还能继续扛,内存依赖图可以做多副本。
6. 进阶验证与演进:监控预案和千万级规模怎么评估
系统上线只是开始,真正考验调度系统的是长期运维。快手的监控预案体系分了三个层次:使用层,看业务方的任务延迟、失败率;服务层,看 Scheduler、Queue Service、Worker 的 CPU、内存、GC、队列积压;依赖层,看 MySQL、K8s、Yarn 这些底层组件。三个层次分别监控,才能快速定位问题到底出在调度逻辑还是资源供给。
分层之外还有分级监控,不同优先级对应不同报警和值班方式。P0 报警直接电话,P1 报警要求 5 分钟内响应,P2 报警进日报汇总。这个设计看起来简单,但非常管用:没有分级,所有报警一起响,值班人员反而分不清主次。我在自己维护的平台里也是这套逻辑,宁可少报警,不能把重要报警淹没在大量噪音里。
故障预案要分系统故障和数据异常两类来准备。系统故障处理预案针对 Scheduler 切换、队列积压、Worker 失联这类问题,核心是先恢复调度能力再排查根因;数据异常处理预案针对被调度任务的数据质量问题,数据出现异常时先用阻断恢复工具阻断错误任务,避免坏数据继续往下游传播,再恢复正确数据。演练也一样,系统故障演练是主动把某个模块搞挂,验证故障处理流程;数据故障演练是模拟数据质量问题,验证阻断和恢复工具。这两类演练我都会固定做,否则预案就是纸上谈兵。
验证调度系统能力,我习惯压两个数:调度延迟 P99 和可用性。压调度延迟不能只看平均,直接用任务发布时间压出分位数:
# 以 200 QPS 发布 5 万个测试实例,统计起调时间延迟 for i in $(seq 1 50000); do curl -s -X POST "http://scheduler-api/publish" \ -H "Content-Type: application/json" \ -d '{"dag_id":"perf_test","execute_time":"now"}' > /dev/null & if (( i % 200 == 0 )); then wait; fi done # 再从日志里统计从 publish 到 worker 开始执行的时间差,取 P50/P99 awk -F',' '{print $3-$2}' scheduler.log | sort -n | \ awk '{a[NR]=$1} END{print "P50="a[int(NR*0.5)], "P99="a[int(NR*0.99)]}'命令里的 execute_time 设为 now 表示立即触发,绕开定时器环节,专门压调度链路本身;wait 每 200 个请求做一次,避免同时建立太多连接把压测机自己打挂。统计脚本取的是日志里的起调时间和实际开始执行时间,差值就是调度延迟。压测时我会顺带观察 Scheduler 的 CPU 和 MySQL 的慢查询数,如果延迟没上去但 CPU 先打满,瓶颈就在 Scheduler;如果两个都上去了,就要看数据库了。
容量评估方面,我的经验是先单 Scheduler 压出单机瓶颈,再按“百万级任务、秒级调度延迟”的指标反推需要多少节点。如果单机每秒能处理 500 次调度决策,百万级任务按峰值每分钟触发 2 万次折算,需要不到 1 个 Scheduler 就能扛住计算量,但为了 HA 至少保持 Active/Standby 各一个。真正要扩容的往往是 Worker 和数据库,而不是 Scheduler。那次大促切换踩坑之后,我每次做调度系统方案都会强制走一遍主备切换演练加重试清理检查,预防性验证到位再谈上线。这种细节平时看着不起眼,关键时刻就是后悔药。希望帮到你。
本文还有配套的精品资源,点击获取