多彩编程 多彩编程MZPH · CODE BLOG
ARTICLE DETAIL

文章详情

深耕前端与后端开发技术的一线实战笔记与踩坑复盘。

从Airflow到Kwaiflow:快手大数据任务调度系统秒级调度架构详解

从Airflow到Kwaiflow:快手大数据任务调度系统秒级调度架构详解 简介这份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 GroupWorker 上的 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 HandlerScheduler 内部把整个调度流程拆成 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 与 StandbyScheduler 的高可用通过 Active/Standby 实现平时 Active 承担调度Standby 热备。切换时最关键的不是谁能当主而是状态怎么交接。如果 Active 有一批已经调度但还没确认执行结果的消息Standby 接管后要把这批消息重新处理但如果这批消息里的任务已经跑起来了直接重放就会导致重复执行。所以主备切换必须配合“状态 幂等”一起做。状态是每个实例当前到哪一步了幂等是同一个实例 ID 只接受一次最终结果。我拆 Kwaiflow 的 Failover 设计时特别注意到了一个细节它强调鲁棒通信协议意思是消息发送方和接收方之间对“什么算成功”有明确的共识比如先写状态再回 Ack还是先回 Ack 再写状态顺序不同故障时的表现完全不同。这个顺序问题是很多自研系统主备切换后数据不一致的根本原因。4.3 分级调度资源紧张时保高优链路快手有大促、大型活动这类场景数据量突然翻倍计算资源不能让所有任务同时按时产出。分级调度的思路是把任务和 Worker Group 都分成 P0 到 P3 等级高优任务优先调度优先执行低优任务限流同时允许人工管控。具体做法是依赖规则管理器。Rule Manager 维护 Allow Rules 和 Block Rules下发到 SchedulerResource 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 P50a[int(NR*0.5)], P99a[int(NR*0.99)]}命令里的 execute_time 设为 now 表示立即触发绕开定时器环节专门压调度链路本身wait 每 200 个请求做一次避免同时建立太多连接把压测机自己打挂。统计脚本取的是日志里的起调时间和实际开始执行时间差值就是调度延迟。压测时我会顺带观察 Scheduler 的 CPU 和 MySQL 的慢查询数如果延迟没上去但 CPU 先打满瓶颈就在 Scheduler如果两个都上去了就要看数据库了。容量评估方面我的经验是先单 Scheduler 压出单机瓶颈再按“百万级任务、秒级调度延迟”的指标反推需要多少节点。如果单机每秒能处理 500 次调度决策百万级任务按峰值每分钟触发 2 万次折算需要不到 1 个 Scheduler 就能扛住计算量但为了 HA 至少保持 Active/Standby 各一个。真正要扩容的往往是 Worker 和数据库而不是 Scheduler。那次大促切换踩坑之后我每次做调度系统方案都会强制走一遍主备切换演练加重试清理检查预防性验证到位再谈上线。这种细节平时看着不起眼关键时刻就是后悔药。希望帮到你。本文还有配套的精品资源点击获取
返回列表