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_sched) | DAG 被调度器排定计划的时间点 |
DAG Starts(dag_start) | DAG 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_delay | dag_sched之后 1 小时(= 实际开始 - 计划时间) |
first_task_scheduling_delay | dag_sched之后 2 小时(= 首个任务开始 - 计划时间) |
duration.success/failure.dag_id | dag_start之后 5 小时(= DAG 总执行时长) |
task_id.duration | taskN_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 == SCHEDULED且clear_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_date到end_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.success、dagrun.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_id、run_type标签。对应的单元测试 test_scheduler_job.py 明确断言:超时失败场景下上报的dagrun.duration.failed必须携带run_type标签,与正常完成路径保持一致。
五、task_id.duration与 landing time:任务级耗时
- 任务执行时长:甘特图中的
task_id.duration对应任务实例从start_date到end_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_id、run_type(序列化为scheduled而非dagruntype.scheduled)、team_name,并可叠加 DAG 自定义标签;内置 key 优先于 DAG 标签,避免冲突覆盖。schedule_delay在启用多团队(_multi_team)时还会额外附带team_name标签。
上报后端选择由 stats_utils.py 决定,优先级依次为:
metrics.statsd_datadog_enabled→ Datadog DogStatsD(datadog_logger.py,使用statsd_host/statsd_prefix作为 namespace);metrics.statsd_on→ 标准 StatsD(statsd_logger.py,支持 UDP / Unix socket、statsd_prefix、InfluxDB tags、metrics_allow_list/metrics_block_list过滤、stat_name_handler自定义命名);metrics.otel_on→ OpenTelemetry metrics;- 均未开启 →
NoStatsLogger(静默丢弃)。
metrics.*配置项由旧版scheduler.statsd_on、scheduler.statsd_host、scheduler.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_after、queued_at、start_date、end_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),仅供参考