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/airflowAirflow 通过 StatsD/Datadog/OTel 等后端持续采集调度与执行过程的关键时序指标帮助运维与开发人员量化DAG 从被调度到执行完成全链路各环节的延迟与耗时。本文以 METRICS.md 中那张甘特图为骨架结合 dagrun.py 与 scheduler_job_runner.py 的源码实现逐一还原每条指标的真实含义、计算方式与采集时机读完后你将能准确读懂 Airflow 指标面板上的每一根时间线。一、指标总览一条 DAG 生命周期上的时间锚点METRICS.md 用一张 Mermaid 甘特图直观展示了一次 DAG 运行中事件与指标在时间轴上的对应关系。以下为该图的核心时间线事件均为里程碑式时间点指标则是这些时间点之间的差值或区间事件 / 指标时间位置 / 计算区间DAG Scheduleddag_schedDAG 被调度器排定计划的时间点DAG Startsdag_startDAG Run 真正开始执行First task Scheduledtask1_sched首个任务被调度Task N ScheduledtaskN_sched第 N 个任务被调度Task N starts runningtaskN_start第 N 个任务开始运行Task N donetaskN_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 UIdag_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, tagsprune_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( fdagrun.{dag.dag_id}.first_task_scheduling_delay, true_delay, tagsself.stats_tags ) stats.timing(dagrun.first_task_scheduling_delay, true_delay, tagsself.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, tagsself.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.pydef _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( fdagrun.duration.{self.state}, dtduration, 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.pyproperty 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 DogStatsDdatadog_logger.py使用statsd_host/statsd_prefix作为 namespacemetrics.statsd_on→ 标准 StatsDstatsd_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 sizefirst_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),仅供参考