运行机制与调优指南:从单实例启动到高可用调度)
Apache Airflow 调度器Scheduler运行机制与调优指南从单实例启动到高可用调度【免费下载链接】airflowApache Airflow - A platform to programmatically author, schedule, and monitor workflows项目地址: https://gitcode.com/GitHub_Trending/ai/airflow本文基于 airflow-core/docs/administration-and-deployment/scheduler.rst 整理扩充并结合 airflow-core 源码中调度器的真实实现scheduler_job_runner.py与其配置定义config.yml进行源码级印证。读者读完后可以掌握如何启动并理解 Airflow Scheduler 的工作循环、如何部署多个 Scheduler 实现高可用HA、影响调度吞吐的瓶颈在哪、以及官方推荐用于调优的每一组[scheduler]配置项及其默认值与底层影响。Scheduler 是什么职责与基本运行方式Airflow Scheduler 是整个平台的调度大脑它监控所有 DAG 与任务一旦任务的依赖满足就触发对应的 Task Instance 执行。官方文档的概括是The Airflow scheduler monitors all tasks and DAGs, then triggers the task instances once their dependencies are complete.在其背后Scheduler 会派生一个子进程subprocess持续监控并同步指定 DAG 目录中的所有 DAG。默认情况下调度器每分钟收集一次 DAG 解析结果并检查当前是否有处于活跃状态的任务可以被触发执行。在生产环境中Scheduler 被设计为一个需要常驻运行的持久化服务。启动调度器启动一个 Scheduler 非常简单只需要执行airflow scheduler该命令读取的配置全部来自airflow.cfg在 Airflow 3.x 中配置的权威定义位于 airflow-core/src/airflow/config_templates/config.ymlairflow.cfg由此模板生成部署时也可以完全用环境变量替代配置文件。调度器启动成功后你的 DAG 就会开始被执行。调度器使用配置中所选的 :doc:Executor如 LocalExecutor、CeleryExecutor、KubernetesExecutor 等来真正运行那些已就绪的任务。也就是说Scheduler 只负责决策谁该跑而在哪跑、怎么跑由 Executor 负责。关于首次 DagRun 与调度时机的两个关键认知第一个 Dag Run 何时产生首次 Dag Run 是基于 DAG 中任务的最小start_date创建的此后的 Dag Run 则按照 DAG 的 :doc:timetable时间表规则持续创建。调度是滞后于区间结束的对于使用 cron 或 timedelta 作为schedule的 DAGScheduler 不会在区间开始时触发任务而是要等到该区间覆盖的周期结束之后才触发。例如一个schedule为daily的任务会在当天结束2019-11-21T23:59之后才开始运行对应数据区间的那个 DagRun。这种设计是有意为之它确保该周期所需的数据在 DAG 执行前已经完整可用。副作用是——在 Airflow UI 上看起来任务像是晚了一天在跑。再强调一遍核心语义Scheduler 是在start_date之后的一个 schedule 周期末尾运行你的作业即数据区间结束后触发。更完整的调度语义可参考官方文档 airflow-core/docs/core-concepts/dag-run其中详细解释了 data interval、logical date 与 catchup 的关系。高吞吐取向与任务优先级的真实语义一个容易引起误解的点是优先级。官方文档明确指出Scheduler 被设计为追求高吞吐high throughput每个调度循环中Scheduler 会检查 Pool 中当前有多少空闲槽位并最多按该数量调度任务实例。这意味着任务优先级priority_weight只有在等待被调度的任务数超过了队列槽位数时才会生效。因此可能出现低优先级任务与高优先级任务处于同一批次时先被调度的情况。如果你的业务对优先级有强诉求需要理解这一批次化调度模型而不是期望严格的全局优先级排序。调度循环背后的源码结构理解调优参数之前最好先看 Scheduler 在外壳之下到底做了什么。从当前仓库源码看调度器的主体实现在 airflow-core/src/airflow/jobs/scheduler_job_runner.py 中的SchedulerJobRunner类第 315 行起。主循环_run_scheduler_loop_run_scheduler_loop第 1792 行起是实际的主循环其 docstring 将其概括为四步通过DagFileProcessorAgent收割HarvestDAG 解析结果查找并排队可执行的任务先在数据库中变更任务实例状态再把任务交给 Executor 排队Executor 心跳异步执行已排队任务并同步正在运行任务的状态检查过期的 Deadline如超时的任务如有则移交处理。循环中还通过EventScheduler注册了一批周期性任务它们各自对应一个配置项可以在 config.yml 中查到默认值周期性动作注册间隔的配置项默认值秒孤儿任务检查与收养scheduler.orphaned_tasks_check_interval300.0触发器超时检查scheduler.trigger_timeout_check_interval15Pool 使用量指标上报scheduler.pool_metrics_interval5.0Task Instance 状态指标上报scheduler.ti_metrics_interval30.0运行中 DagRun 指标上报scheduler.dagrun_metrics_interval30.0心跳超时任务清理scheduler.task_instance_heartbeat_timeout_detection_interval10.0卡在 queued 状态的任务处理scheduler.task_queued_timeout_check_interval依据task_queued_timeout陈旧 DAG / 孤儿资产清理scheduler.parsing_cleanup_interval60每次循环迭代都会用stats.timer(scheduler.scheduler_loop_duration)统计耗时并在会话结束后expunge_all()释放 ORM 对象——源码注释里甚至提醒我们刚才可能查看了 500 个 ORM 对象可见其内存管理是刻意为之的。决策函数_do_scheduling单次调度决策集中在_do_scheduling第 1992 行起其行为与官方 HA 文档描述一一对应创建必要的 DagRun通过DagModel的next_dagrun_create_after列判定哪些 DAG 需要新 DagRun。由于创建 DagRun 相对耗时默认每轮最多处理 10 个 DAGscheduler.max_dagruns_to_create_per_loop——设得过高会导致一个 Scheduler 忙于建 Run 而没空调度任务。取出最久未被审视的 N 个运行中 DagRun 进行推进N 默认为 20scheduler.max_dagruns_per_loop_to_schedule。这些行是带行锁选出的见下文 HA 数据库要求因此同时只有一个 Scheduler 能处理它们。调大该值对小 DAG 更友好但对超过 500 个任务的大 DAG 反而可能变慢。进入临界区Critical Section排队任务加锁 Pool 表行把任务送入 Executor对应_critical_section_enqueue_task_instances。随后是回调发送、session.expunge_all()再判断total_free_executor_slots如果所有 Executor 都已满则直接跳过临界区否则进入临界区排队任务。若在临界区遇到OperationalError且判定为锁不可用说明另一个 Scheduler 正持有锁本实例会rollback并直接结束本轮指标scheduler.critical_section_busy会增加而不是阻塞等待。运行多个 Scheduler高可用HA部署从 Airflow 2.0.0 起当前仓库为 3.x完全支持Airflow 支持同时运行多个 Scheduler——既是为了性能扩展也是为了故障韧性。设计思路复用元数据库而非引入共识协议HA Scheduler 的核心设计是充分利用已有的元数据库Metadata Database。这样做的初衷是运维简单性每个组件本来就要与数据库通信因此不再引入 Scheduler 之间的直接通信或共识算法如 Raft、Paxos也不需要额外的共识组件如 ZooKeeper、Consul从而把运维面降到最低。在 HA 模式下Scheduler 使用序列化 DAGserialized DAG表示来做调度决策。其调度循环的粗略轮廓是检查是否有 DAG 需要创建新的 DagRun如有则创建审视一批 DagRun找出可调度的 TaskInstance 或已完成的 DagRun选出可调度的 TaskInstance在遵守 Pool 限制与其他并发限制的前提下将其入队enqueue执行。对数据库的要求为什么需要行级锁保持高吞吐的代价是调度循环中有一段计算需要在内存中完成如果每个 TaskInstance 都往返数据库做判断速度会慢到不可接受因此必须保证同一时刻只有一个 Scheduler 处于该临界区——否则并发与 Pool 限制无法被正确执行。实现手段是数据库行级锁SELECT ... FOR UPDATE。临界区正是任务从 scheduled 状态被入队到 Executor、同时守住各类并发与 Pool 限制的地方获取临界区的方式是对slot_poolPool 表的每一行申请行级写锁——语义上等价于SELECT * FROM slot_pool FOR UPDATE NOWAIT实际 SQL 略有差异。数据库兼容矩阵PostgreSQL 12 / MySQL 8.0开箱即用可以随意运行任意数量的 Scheduler 副本无需额外配置。MariaDB直到 10.6.0 才实现SKIP LOCKED/NOWAIT子句。缺少这些能力时运行多 Scheduler不被支持且有社区报告的死锁错误10.6.0 及以后在多 Scheduler 下可能工作正常但未经测试。Microsoft SQL Server未经 HA 测试。从源码看临界区的真实行为在 scheduler_job_runner.py 第 1202 行起的_critical_section_enqueue_task_instances中临界区被描述为三步按优先级挑选 TaskInstance约束是状态符合预期、且不超出max_active_runs与 Pool 限制原子地变更上述 TaskInstance 的状态将 TaskInstance 入队到 Executor。关于锁竞争源码给出了更精确的补充支持 NOWAIT 的数据库上一个被阻塞的 Scheduler 会跳过临界区继续做其它事创建新 DagRun、把任务从 None 推进到 SCHEDULED 等不支持 NOWAIT 的数据库如 MariaDB、MySQL 5.x上其它 Scheduler 会等待锁释放后再继续。这也再次解释了上面 MariaDB 的兼容性警告。此外该函数在计算本轮最多可调度多少任务时会同时参考max_tis_per_query与全局core.parallelism若max_tis_per_query为 0则max_tis parallelism - 当前已占用槽位否则max_tis min(max_tis_per_query, parallelism - 当前已占用槽位)。core.parallelism代表单个 Scheduler 最多可同时运行的任务数上限在引入多 Executor 后Scheduler 负责在所有 Executor 之间总量不超过该上限。这就是一个参数控制整条调度流水线节流的底层逻辑。Scheduler 性能调优指南什么在影响 Scheduler 的性能调优前需要先盘点影响面官方将其归纳为三类部署形态可用内存、CPU、网络吞吐。DAG 结构与逻辑DAG 数量、复杂度任务数、依赖数。Scheduler 配置Scheduler 实例数、单循环处理的 TaskInstance 数、单循环创建/调度的新 DagRun 数、执行清理与孤儿任务检测/收养的频率。社区也公开过深挖调度器内部的演讲资料如 Airflow Summit 2021 的 Deep Dive into the Airflow Scheduler可作为理解under-the-hood行为的补充材料。调优的方法论官方并不推荐任何特定监控工具而是强调一套通用的性能优化流程用你惯用的监控手段持续观测系统——本文不展开具体指标与工具重点在于明确该监控哪些资源确定当前最想优化的性能维度提高吞吐降低延迟观测瓶颈出现在哪CPU、内存、I/O 通常是最主要的限制因素根据期望与观测决定下一步改动然后回到第 1 步——性能改进是一个迭代过程。Airflow 提供了大量旋钮但选择旋哪些、往哪个方向拧取决于你的部署、DAG 结构、硬件与预期。管理部署的一部分工作就是先想清楚你要为哪个指标做优化。可能限制 Scheduler 性能的资源点数据库连接与数据库使用Airflow 以吃数据库连接著称——DAG 越多、并行处理越多打开的数据库连接越多。MySQL 的线程模型下通常不是问题但 Postgres 是进程模型连接多会很吃力。业界普遍建议即使是中等规模的基于 Postgres 的 Airflow 部署也最好用 PGBouncer 作为数据库代理。仓库中的 chart/Airflow Helm Chart开箱即用地支持 PGBouncer。CPUAirflow Scheduler 在多个实例下几乎线性扩展。如果瓶颈是 CPU增加 Scheduler 实例数即可。内存观察内存时要注意看的是哪类内存——通常应该关注**工作内存working memory**而非总内存使用量命名依部署而异。可以采取的性能改进措施提高资源利用率当 CPU/内存/I/O/网络存在闲置容量时增加 Scheduler 数量或缩短高频动作的间隔可能以更高资源占用为代价换取性能提升。增加硬件能力CPU 打满往往是系统能力不足例如单机 CPU 全占时可在一台新机器上再加一个 Scheduler。多数情况下添加第 2、第 3 个 Scheduler 后调度能力近似线性增长除非共享数据库等已成为瓶颈。实验不同的调度器可调参数tunables调优往往是用一个性能面交换另一个性能面的艺术——例如让 Scheduler 每轮处理更多 DagRun 提升了小 DAG 的吞吐却可能损害大 DAG 的吞吐。Scheduler 配置项详解官方 Tunables以下是官方文档推荐的、用于控制 Scheduler 行为与性能的核心配置项均位于[scheduler]section更完整的非性能类参数见 airflow-core/docs/configurations-ref.rst。默认值与类型来自当前仓库的配置权威定义 config.yml对应行区间约 第 2650–2821 行。max_dagruns_to_create_per_loop默认值10类型 integer版本 2.0.0。作用改变每个 Scheduler 在创建 DagRun 时锁定的 DAG 数量。调优提示如果你的 DAG 极大单 DAG 任务数在 1 万以上并且跑多个 Scheduler你通常不希望一个 Scheduler 独揽全部创建工作此时可适当调低该值。max_dagruns_per_loop_to_schedule默认值20类型 integer版本 2.0.0。作用调度与排队任务时一个 Scheduler 最多审视并锁定多少个 DagRun。调优提示增大该值可提高小 DAG 的吞吐但很可能拖慢大 DAG例如单 DAG 超过 500 个任务的吞吐多 Scheduler 场景下设得过高还可能导致某个 Scheduler 抢走所有 DagRun让其它 Scheduler 无事可做。use_row_level_locking默认值True类型 boolean版本 2.0.0。作用是否在相关查询中下发SELECT ... FOR UPDATE。警告如果设为False就不要同时运行超过一个 Scheduler——否则并发限制将无法被正确执行。pool_metrics_interval默认值5.0秒类型 float版本 2.0.0。作用Pool 使用量统计发送到 StatsD 的频率在启用statsd_on时生效。调优提示该统计查询相对昂贵建议设置成与 StatsD 聚合周期一致。ti_metrics_interval默认值30.0秒类型 float版本 3.0.0。作用Task Instancescheduled、queued、running、deferred 状态统计发送到 StatsD 的频率在启用statsd_on时生效。调优提示同样是相对昂贵的查询建议与 StatsD 聚合周期一致。另外从 3.1.0 起还有同类参数dagrun_metrics_interval默认30.0控制运行中 DagRun 指标的发送频率。orphaned_tasks_check_interval默认值300.0秒类型 float版本 2.0.0。作用Scheduler 检查孤儿任务orphaned tasks与已死 SchedulerJob 的频率。含义它决定了一个已死 Scheduler 被发现、其监管的任务被另一个 Scheduler 接管的速度。任务本身会继续运行所以一段时间没检测出来并无大碍。当一个 SchedulerJob 被判定为死亡依据scheduler.scheduler_health_check_threshold该进程之前启动的 running/queued 任务将被**收养adopted**并由当前 Scheduler 继续监管。该周期在源码_run_scheduler_loop中通过EventScheduler.call_regular_interval注册启动时也会立即执行一次检查见 scheduler_job_runner.py。max_tis_per_query默认值16类型 integer。作用调度主循环中查询的批大小即每轮调度循环评估多少个 TaskInstance。警告不应大于core.parallelism。设得太高会因查询谓词复杂导致 SQL 性能下降或锁竞争过度还可能触达数据库单条 SQL 的最大长度限制。设为0时使用core.parallelism的值这正是临界区源码中的max_tis计算逻辑。scheduler_idle_sleep_time默认值1秒类型 float版本 2.2.0。作用控制 Scheduler 在本轮无事可做时、两轮循环之间的休眠时长。如果本轮实际调度了任务会立刻开始下一轮迭代不进入休眠。备注官方文档坦言该参数命名不佳历史原因未来会弃用当前名称并更名。与 HA/健康相关的邻近配置源码可查、便于组合使用在 config.yml 的[scheduler]section 中还有一组与调度器生命周期相关的参数可与上述 tunables 组合使用scheduler_health_check_threshold默认30秒。最后一次 Scheduler 心跳距今超过该值时Scheduler 被判为不健康用于/health端点与airflow jobs checkCLI。enable_health_check默认False。设为True时启动 Scheduler 会额外拉起一个微型 Web 服务子进程用于健康检查其监听端口由scheduler_health_check_server_port默认8974决定、绑定地址由scheduler_health_check_server_host默认0.0.0.0决定。num_runs默认-1表示每个 DAG 文件尝试调度多少次-1 为无限。only_idle3.2.0 引入默认False可让计数只在调度器空闲本轮没有排队/完成任务时递增处理过任务后计数会重置。catchup_by_default默认False。控制 Scheduler 是否默认执行补跑catchup可在 DAG 定义中按 DAG 覆盖catchup。task_instance_heartbeat_timeout默认300与task_instance_heartbeat_timeout_detection_interval默认10.0本地任务作业LocalTaskJob心跳超时多久判死、以及 Scheduler 以多高频率扫描心跳超时的任务实例。需要指出的是[scheduler]下还有大量非性能相关参数如use_job_schedule、trigger_timeout_check_interval、task_queued_timeout、parsing_cleanup_interval等它们共同组成了调度循环的全部定时行为。调优时建议以 config.yml 为唯一权威来源核对参数名与默认值避免新旧版本参数名不一致导致的误配置。将各部分拼起来一条任务从 DAG 文件到 Executor 的旅程把文档与源码对照一次完整的调度可以描述为DagFileProcessor 子进程解析 DAG 目录目录监听与解析节奏由dag_processor与解析相关配置控制产生解析结果主循环每分钟左右收集解析结果写入 serialized DAG 表示_do_scheduling依据next_dagrun_create_after批量创建 DagRun受max_dagruns_to_create_per_loop限制选出最久未处理的一批运行中 DagRun受max_dagruns_per_loop_to_schedule限制推进任务状态与 DagRun 终态若 Executor 有空闲槽位进入临界区锁定 Pool 行受use_row_level_locking控制按优先级从高到低挑选max_tis_per_query个任务并原子改状态为 queued任务被分发到对应 Executor 排队、执行、回报状态各类周期任务孤儿任务收养、心跳超时清理、指标上报等按各自 interval 在循环中穿插执行。理解这条链路后再去调整上表中的任何旋钮都会比盲调更有依据。运维建议小结生产环境务必让 Scheduler 常驻配合进程守护或容器平台的重启策略参考仓库内作业生命周期说明 airflow-core/src/airflow/jobs/JOB_LIFECYCLE.md 可以了解 SchedulerJob 在数据库中记录的状态流转。HA 场景请使用 PostgreSQL 12 或 MySQL 8.0保持use_row_level_locking True避免 MariaDB10.6.0 以下与 SQL Server。性能调优是迭代过程先监控CPU/内存/I-O/数据库连接再确定优化目标最后小步调整max_dagruns_to_create_per_loop、max_dagruns_per_loop_to_schedule、max_tis_per_query等参数并回测。需要扩容吞吐时优先考虑横向增加 Scheduler 实例——它几乎线性扩展。数据库连接焦虑优先交给 PGBouncerHelm Chart 开箱支持而不是粗暴调大连接池上限。【免费下载链接】airflowApache Airflow - A platform to programmatically author, schedule, and monitor workflows项目地址: https://gitcode.com/GitHub_Trending/ai/airflow创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考