news 2026/9/11 17:06:54

Apache Airflow DAG 执行指标全景:从 schedule_delay 到 task duration 的源码级解读

作者头像

张小明

前端开发工程师

1.2k 24
文章封面图
Apache Airflow DAG 执行指标全景:从 schedule_delay 到 task duration 的源码级解读

Apache Airflow DAG 执行指标全景:从 schedule_delay 到 task duration 的源码级解读

【免费下载链接】airflowApache Airflow - A platform to programmatically author, schedule, and monitor workflows项目地址: https://gitcode.com/GitHub_Trending/ai/airflow

Airflow 通过 StatsD/Datadog/OTel 等后端持续采集调度与执行过程的关键时序指标,帮助运维与开发人员量化"DAG 从被调度到执行完成"全链路各环节的延迟与耗时。本文以 METRICS.md 中那张甘特图为骨架,结合 dagrun.py 与 scheduler_job_runner.py 的源码实现,逐一还原每条指标的真实含义、计算方式与采集时机,读完后你将能准确读懂 Airflow 指标面板上的每一根时间线。

一、指标总览:一条 DAG 生命周期上的时间锚点

METRICS.md 用一张 Mermaid 甘特图直观展示了"一次 DAG 运行"中事件与指标在时间轴上的对应关系。以下为该图的核心时间线(事件均为里程碑式时间点,指标则是这些时间点之间的差值或区间):

事件 / 指标时间位置 / 计算区间
DAG Scheduled(dag_schedDAG 被调度器排定计划的时间点
DAG Starts(dag_startDAG Run 真正开始执行
First task Scheduled(task1_sched首个任务被调度
Task N Scheduled(taskN_sched第 N 个任务被调度
Task N starts running(taskN_start第 N 个任务开始运行
Task N done(taskN_done第 N 个任务完成
Last task ends, DAG execution ends最后一个任务结束,DAG 执行结束
dagrun.schedule_delaydag_sched之后 1 小时(= 实际开始 - 计划时间)
first_task_scheduling_delaydag_sched之后 2 小时(= 首个任务开始 - 计划时间)
duration.success/failure.dag_iddag_start之后 5 小时(= DAG 总执行时长)
task_id.durationtaskN_start之后 1 小时(= 任务执行时长)
task N landing time(仅 Airflow UI)dag_sched之后 5 小时(= 从调度到任务完成)

甘特图中里程碑被设置为 2 分钟的占位时长,指标条目的长度(1h、2h、5h 等)只是为了在图中错落排布,便于阅读时人工比对"某指标对应哪两个事件之间的区间",并非真实数值。

原图中各指标的具体定义如下:

二、dagrun.schedule_delay:调度器把 Run 置为 RUNNING 的延迟

指标含义:一次调度触发的 DAG Run 从"应开始时间(run_after)"到"调度器实际将其置为 RUNNING 状态"之间的延迟,反映调度器的积压与调度响应速度。

源码实现位于 scheduler_job_runner.py:调度器在处理"已入队、待置为 running"的 DAG Run 时,将其start_date设为当前时间,并与dag_run.run_after(即时间表计算的期望开始时间)做差:

def _update_state(dag: SerializedDAG, dag_run: DagRun): dag_run.state = DagRunState.RUNNING dag_run.start_date = timezone.utcnow() if ( dag.timetable.periodic and dag_run.run_type != DagRunType.MANUAL and dag_run.triggered_by != DagRunTriggeredByType.ASSET and dag_run.clear_number < 1 ): expected_start_date = dag_run.run_after schedule_delay = dag_run.start_date - expected_start_date stats.timing( "dagrun.schedule_delay", schedule_delay, tags=prune_dict({...}), )

需要注意的采集前置条件(源码中明确体现):

  • dag.timetable.periodic必须为真——即 DAG 使用周期性时间表(如 Cron),一次性/事件驱动的 DAG 不采集;
  • run_type必须是SCHEDULED,手动触发(MANUAL)与资产(Asset)触发的 Run 不采集;
  • clear_number < 1,即该 Run 没有被 clear 过(重跑会污染该指标);
  • dag_run.start_date由调度器在此刻写入timezone.utcnow(),因此该指标本质上是"排队等待被调度器接管"的时间。

三、first_task_scheduling_delay:从计划时间到首个任务真正启动

指标含义:DAG 中第一个开始运行的任务的start_date减去 DAG Run 的期望开始时间run_after,度量"计划调度 → 实际开始干活"的完整延迟(包含调度器延迟、任务排队、Executor 领取、worker 启动等)。

源码实现位于 dagrun.py,由_emit_true_scheduling_delay_stats_for_finished_state在 DAG Run 进入终态(success/failure)时计算并上报:

first_start_date = min(ti.start_date for ti in finished_tis if ti.start_date) true_delay = first_start_date - self.run_after if true_delay.total_seconds() > 0: stats.timing( f"dagrun.{dag.dag_id}.first_task_scheduling_delay", true_delay, tags=self.stats_tags ) stats.timing("dagrun.first_task_scheduling_delay", true_delay, tags=self.stats_tags) if self.queued_at is not None: start_delay = first_start_date - self.queued_at if start_delay.total_seconds() > 0: stats.timing("dagrun.first_task_start_delay", start_delay, tags=self.stats_tags)

该指标同时上报两个版本

  • dagrun.first_task_scheduling_delay:全局聚合名,用于跨 DAG 对比;
  • dagrun.{dag_id}.first_task_scheduling_delay:按 DAG 维度,便于定位特定 DAG 的调度健康度。

此外,若queued_at存在,还会额外上报dagrun.first_task_start_delay(首个任务从入队到启动的延迟)。

采集前置条件(源码 docstring 与守卫条件明确):

  • 仅当run_type == SCHEDULEDclear_number == 0(即调度器触发、未被清理重跑);
  • 必须存在已结束任务(finished_tis非空);
  • DAG 时间表必须是周期性的(dag.timetable.periodic),否则没有可参照的下一次调度来计算延迟;
  • true_delay.total_seconds() > 0才会上报,负值(时钟回拨等异常)被丢弃。

关于离群值:源码注释特别提醒——当首个任务被 clear 后,会取第二个任务的start_date作为最小值,从而产生离群点;这类离群值应在 StatsD/Datadog 侧通过 dashboard 工具过滤,而非在 Airflow 内部处理。

四、dagrun.duration.*:一次 DAG Run 的总耗时

指标含义:DAG Run 从start_dateend_date的执行总时长,按最终状态区分指标名。

源码实现位于 dagrun.py:

def _emit_duration_stats_for_finished_state(self): if self.state == DagRunState.RUNNING: return if self.start_date is None or self.end_date is None: return duration = self.end_date - self.start_date stats.timing( f"dagrun.duration.{self.state}", dt=duration, tags={**self.stats_tags, "dag_id": self.dag_id}, )

指标名中的状态取值即DagRunState(如dagrun.duration.successdagrun.duration.failed)。其触发点位于 dagrun.py 的update_state中:当 DAG Run 状态被更新为终态时,同时调用_emit_true_scheduling_delay_stats_for_finished_state_emit_duration_stats_for_finished_state

另一条独立的失败时长上报路径在 scheduler_job_runner.py:当 DAG Run 因dagrun_timeout超时被判定失败时,调度器直接计算end_date - start_date并上报dagrun.duration.failed,且携带dag_idrun_type标签。对应的单元测试 test_scheduler_job.py 明确断言:超时失败场景下上报的dagrun.duration.failed必须携带run_type标签,与正常完成路径保持一致。

五、task_id.duration与 landing time:任务级耗时

  • 任务执行时长:甘特图中的task_id.duration对应任务实例从start_dateend_date的耗时。任务实例模型直接以ti.duration属性暴露该值(见 scheduler_job_runner.py 与执行 API 路由 task_instances.py),并在 Airflow UI 中以"Duration"列展示。该指标关注的是任务真正运行(running)阶段的耗时,不含排队时间。
  • landing time:甘特图中特别标注"task N landing time (only in airflow UI)"——即任务从被调度(dag_sched)到完成(taskN_done)的总跨度,仅存在于 Airflow UI(任务实例表格中的 Landing Time 列),并非以 StatsD 指标形式对外上报。它度量的是任务"从排定计划到最终落地"的全链路时间,涵盖调度、排队、运行全过程。

六、指标标签与上报后端:tags 如何携带上下文

所有 DAG 级指标统一携带stats_tags,其构造见 dagrun.py:

@property def stats_tags(self) -> dict[str, str]: base = prune_dict( { "dag_id": self.dag_id, "run_type": getattr(self.run_type, "value", self.run_type), # 如 "scheduled" "team_name": getattr(self, "_team_name", None), } ) dag_tags = self.dag_tags_for_stats() return {**dag_tags, **base}

即默认携带dag_idrun_type(序列化为scheduled而非dagruntype.scheduled)、team_name,并可叠加 DAG 自定义标签;内置 key 优先于 DAG 标签,避免冲突覆盖。schedule_delay在启用多团队(_multi_team)时还会额外附带team_name标签。

上报后端选择由 stats_utils.py 决定,优先级依次为:

  1. metrics.statsd_datadog_enabled→ Datadog DogStatsD(datadog_logger.py,使用statsd_host/statsd_prefix作为 namespace);
  2. metrics.statsd_on→ 标准 StatsD(statsd_logger.py,支持 UDP / Unix socket、statsd_prefix、InfluxDB tags、metrics_allow_list/metrics_block_list过滤、stat_name_handler自定义命名);
  3. metrics.otel_on→ OpenTelemetry metrics;
  4. 均未开启 →NoStatsLogger(静默丢弃)。

metrics.*配置项由旧版scheduler.statsd_onscheduler.statsd_hostscheduler.statsd_prefix迁移而来(见 config_command.py 中的重命名映射)。

七、典型观测场景与排查指引

将上述指标组合使用,可以快速定位调度链路瓶颈:

  • dagrun.schedule_delay持续偏高:调度器本身积压或 DAG Run 排队未被及时接管,重点排查 scheduler 负载、max_active_runs限制、Executor 排队(queue size);
  • first_task_scheduling_delay高而schedule_delay正常:瓶颈在任务调度/分发环节(任务依赖不满足、池(pool)被占满、Executor slot 不足、worker 拉取延迟);
  • dagrun.duration.success/failed异常拉长:结合task_id.duration定位是哪个任务拖慢了整体,再区分是任务自身执行慢还是排队慢;
  • dagrun.first_task_start_delay(入队→启动)偏大:关注 Executor 与 worker 之间的任务领取与资源分配效率;
  • 离群值处理:任务被 clear/重跑会导致first_task_scheduling_delay出现离群点,观测时应结合 run_type/clear 信息在 dashboard 侧过滤,详见 dagrun.py 的注释说明。

对应的单元测试提供了指标语义的权威佐证:test_dagrun.py 验证了调度延迟统计在dag_run.update_state()时被正确触发、且同时上报全局与按 dag_id 两个版本;test_scheduler_job.py 验证了dagrun.schedule_delay携带 dag_id 标签的上报路径。

八、小结

Airflow 的 DAG 执行指标虽然命名直观,但每个指标都隐含了严格的采集条件(周期性时间表、SCHEDULED run 类型、未被 clear 等)与特定的计算锚点(run_afterqueued_atstart_dateend_date)。理解 METRICS.md 中那张甘特图背后的源码逻辑,是搭建可靠的调度健康监控面板、快速定位调度延迟问题的基础。实际接入时,只需开启metrics.statsd_on(或 Datadog/OTel)并配置上报端点,上述指标便会随调度与执行流程自动上报,配合 tags 即可完成多维度的下钻分析。

【免费下载链接】airflowApache Airflow - A platform to programmatically author, schedule, and monitor workflows项目地址: https://gitcode.com/GitHub_Trending/ai/airflow

创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考

版权声明: 本文来自互联网用户投稿,该文观点仅代表作者本人,不代表本站立场。本站仅提供信息存储空间服务,不拥有所有权,不承担相关法律责任。如若内容造成侵权/违法违规/事实不符,请联系邮箱:809451989@qq.com进行投诉反馈,一经查实,立即删除!
网站建设 2026/9/11 17:05:41

Shan-Chen LBM两相流C++实现:从伪势力到VTK可视化

简介&#xff1a;本资源是一份面向计算流体力学初学者与C编程学习者的两相流数值模拟实践代码&#xff0c;聚焦Lattice Boltzmann Method&#xff08;LBM&#xff09;与Shan-Chen多相模型的工程实现。它解决了二维两相流中界面演化、表面张力建模等关键问题&#xff0c;适用于高…

作者头像 李华
网站建设 2026/9/11 17:03:33

燃气营销管理系统:5大核心功能与3类应用场景实战拆解

燃气行业的竞争格局正在发生深刻变化。随着管网规模持续扩张、终端用户数量不断增长&#xff0c;传统的客户台账登记、抄表收费和业务办理模式&#xff0c;已经难以支撑精细化管理需求。尤其是在市场化改革推进的背景下&#xff0c;燃气企业既要保障安全供气的底线&#xff0c;…

作者头像 李华
网站建设 2026/9/11 17:02:25

量化投资:月末交易策略回测与优化实践

1. 月末策略标的回测研究概述月末交易策略是量化投资领域一个经典的研究方向。每到月末&#xff0c;市场往往会出现特定的资金流动模式&#xff0c;这为策略开发提供了天然的逻辑基础。我在过去三年持续跟踪这个策略时发现&#xff0c;单纯依靠传统的月末效应已经很难获得稳定收…

作者头像 李华