ARTICLE DETAIL

资讯详情

深耕郑州网站建设与运营推广的一线实战洞察。

Apache Airflow 指标标签一致性修复:为 `dagrun.duration.failed` 补充 `run_type` 标签

Apache Airflow 指标标签一致性修复:为 `dagrun.duration.failed` 补充 `run_type` 标签 Apache Airflow 指标标签一致性修复为dagrun.duration.failed补充run_type标签【免费下载链接】airflowApache Airflow - A platform to programmatically author, schedule, and monitor workflows项目地址: https://gitcode.com/GitHub_Trending/ai/airflow导读当 Dag Run 因超时被强制标记为失败时Airflow Scheduler 会向 StatsD 发出dagrun.duration.failed计时指标。本次修复为该指标补齐了run_type标签使其与正常完成success/failed路径上由DagRun._emit_duration_stats_for_finished_state发出的同一系列指标保持一致。读完本文你将掌握该指标的两条完整发射链路、stats_tags的标签构成规则、run_type的全部取值枚举以及如何用仓库内的单元测试验证这一行为。1. 问题背景两条发射路径的标签不一致在 Apache Airflow 中dagrun.duration.*系列指标用于度量一次 Dag Run 从start_date到end_date的持续时长timing 类型。该系列指标名以 Dag Run 的最终状态为后缀如dagrun.duration.success、dagrun.duration.failed并携带dag_id、run_type、team_name等标签tags用于多维度聚合。在本次修复之前该系列指标存在一条发射路径不一致正常完成路径Dag Run 以 success/failed 等终结状态结束时由 DagRun._emit_duration_stats_for_finished_state 统一发射指标。该方法先做防御性检查状态为 RUNNING、start_date/end_date缺失时直接返回再计算duration self.end_date - self.start_date以fdagrun.duration.{self.state}为指标名、{**self.stats_tags, dag_id: self.dag_id}为标签调用stats.timing。由于stats_tags属性天然包含run_type因此正常路径下该指标总是带有run_type标签。超时强制失败路径当 Dag Run 超过 DAG 配置的dagrun_timeout仍未能完成时Scheduler 会将其状态强制置为FAILED并单独发射dagrun.duration.failed指标。问题就出在这里——超时路径的发射代码在 scheduler_job_runner.py 中独立实现此前构造标签时遗漏了run_type导致同名的dagrun.duration.failed指标在不同触发场景下标签集合不一致。2. 修复内容超时路径复用stats_tags本修复对应 newsfragment 文件 67765.bugfix.rst的核心改动是在超时路径发射dagrun.duration.failed时将dag_run.stats_tags展开进标签字典if dag_run.end_date and dag_run.start_date: duration dag_run.end_date - dag_run.start_date stats.timing( dagrun.duration.failed, duration, tagsprune_dict( { **dag_run.stats_tags, team_name: self._get_team_names_for_dag_ids([dag_run.dag_id], session).get( dag_run.dag_id ) if self._multi_team else None, } ), )关键点拆解**dag_run.stats_tags将 Dag Run 的标准统计标签含dag_id、run_type、team_name展开合并这是本次修复的核心——run_type由此进入超时场景的指标标签prune_dict包裹过滤掉值为假None、空串等的标签键避免向 StatsD 发送空值标签条件补写team_name当 Scheduler 启用了多团队_multi_team时从数据库查询并补充团队名未启用时该键值为None随后被prune_dict剔除。注意这里team_name的计算结果会覆盖stats_tags中已有的同名键因为写在展开之后。修复后超时路径与正常完成路径在标签构成上完全对齐二者都包含dag_id与run_type只是在team_name的取值细节上保留了各自的上下文差异。2.1 正常完成路径的对照实现为了理解一致性的具体含义对照查看正常完成路径 DagRun._emit_duration_stats_for_finished_statedef _emit_duration_stats_for_finished_state(self): if self.state DagRunState.RUNNING: return if self.start_date is None: self.log.warning(Failed to record duration of %s: start_date is not set., self) return if self.end_date is None: self.log.warning(Failed to record duration of %s: end_date is not set., self) return duration self.end_date - self.start_date stats.timing( fdagrun.duration.{self.state}, dtduration, tags{**self.stats_tags, dag_id: self.dag_id}, )可以看到正常路径直接以fdagrun.duration.{self.state}作为指标名并且总是携带stats_tags内含run_type。修复前的超时路径缺失的正是这一关键标签。3. 标签的底层来源DagRun.stats_tags属性run_type标签之所以能进入两条路径根源在于DagRun模型统一提供的 stats_tags 属性property def stats_tags(self) - dict[str, str]: # prune_dict strips falsy values, so merge dag tags after it runs so standalone # tags (empty value) are preserved for DogStatsD emission. base prune_dict( { dag_id: self.dag_id, # bare value so it serializes as e.g. scheduled, not dagruntype.scheduled run_type: getattr(self.run_type, value, self.run_type), team_name: getattr(self, _team_name, None), } ) dag_tags self.dag_tags_for_stats() # Built-in keys win on collision; dag tags fill in everything else. return {**dag_tags, **base}实现细节值得注意run_type的序列化方式getattr(self.run_type, value, self.run_type)确保标签值取枚举的裸值如scheduled、manual而不是dagruntype.scheduled这样的枚举字符串形式内置键优先级{**dag_tags, **base}保证dag_id、run_type、team_name这三个内置键在与用户自定义 DAG 标签冲突时胜出prune_dict的双重应用先过滤内置标签中的假值再合并 DAG 标签这样值为空串的独立 DAG 标签也能保留下来传给 DogStatsD。3.1run_type的完整取值run_type标签的所有可能值由 DagRunType 枚举定义class DagRunType(str, enum.Enum): Class with DagRun types. BACKFILL_JOB backfill SCHEDULED scheduled MANUAL manual OPERATOR_TRIGGERED operator_triggered ASSET_TRIGGERED asset_triggered ASSET_MATERIALIZATION asset_materialization这意味着在修复之后dagrun.duration.failed指标可以按run_type维度区分由调度触发的超时scheduled、手动触发后超时manual、资产Asset事件触发后超时asset_triggered等从而在监控看板上分别观察不同来源 Dag Run 的超时耗时分布。4. 超时判定与指标发射的完整链路要理解该指标在什么时机被发射需要梳理 Scheduler 的超时处理逻辑。在 SchedulerJobRunner._schedule_dag_run 中超时判定的核心条件是if ( dag_run.start_date and dag.dagrun_timeout and dag_run.start_date timezone.utcnow() - dag.dagrun_timeout ): dag_run.set_state(DagRunState.FAILED)链路顺序如下超时判定dag_run.start_date存在、DAG 配置了dagrun_timeout、且start_date早于当前时间 - dagrun_timeout三者同时满足即判定超时状态强制置为 FAILED调用dag_run.set_state(DagRunState.FAILED)未完成任务置为 SKIPPED对仍在运行的任务实例统一置为SKIPPED并 flush回调与通知构造timed_out原因的失败回调callback_to_execute随后调用notify_dagrun_state_changed(msgtimed_out)触发监听器发射超时耗时指标在end_date与start_date均存在时发射dagrun.duration.failed本次修复点。dagrun_timeout是 DAG 级别的配置定义于 DAG 构造参数中仅在start_date存在时才会参与超时判定。此外超时后若 Dag Run 处于终结状态且run_type属于SCHEDULED/MANUAL/ASSET_TRIGGERED还会触发_set_exceeds_max_active_runs逻辑为后续调度预留容量。5. 测试验证仓库中的行为断言本次修复配套了专门的单元测试位于 test_scheduler_job.py 的 test_dagrun_timeout_duration_metric_has_run_type。该测试完整复现了超时场景并断言标签内容mock.patch(airflow._shared.observability.metrics.stats._get_backend) def test_dagrun_timeout_duration_metric_has_run_type(self, mock_get_backend, dag_maker): ... mock_stats mock.MagicMock(specStatsLogger) mock_get_backend.return_value mock_stats session settings.Session() with dag_maker( dag_idtest_dagrun_timeout_duration_metric, dagrun_timeoutdatetime.timedelta(seconds60), sessionsession, ): EmptyOperator(task_iddummy) dr dag_maker.create_dagrun(start_datetimezone.utcnow() - datetime.timedelta(days1)) scheduler_job Job() self.job_runner SchedulerJobRunner(jobscheduler_job) self.job_runner._schedule_dag_run(dr, session) session.flush() session.refresh(dr) assert dr.state State.FAILED mock_stats.timing.assert_any_call( dagrun.duration.failed, mock.ANY, tags{dag_id: dr.dag_id, run_type: dr.run_type}, )测试的关键构造手法dagrun_timeoutdatetime.timedelta(seconds60)为测试 DAG 配置 60 秒超时create_dagrun(start_datetimezone.utcnow() - datetime.timedelta(days1))将 Dag Run 的start_date设为 1 天前使已运行时长远超 60 秒超时阈值确保触发超时判定断言mock_stats.timing.assert_any_call(dagrun.duration.failed, mock.ANY, tags{dag_id: dr.dag_id, run_type: dr.run_type})直接验证了指标名、时长值与标签集合其中tags精确匹配dag_id与run_type两个键——这正是本次修复的验收标准。测试同时断言dr.state State.FAILED确认指标发射前提超时后状态为 failed成立。6. 对监控实践的指导意义该修复对使用 StatsD/Datadog 等时序监控体系的 Airflow 运维者具有直接价值标签维度统一修复前同一指标名dagrun.duration.failed在不同场景下标签集合不同会导致时序数据库中的 series 数量分裂、聚合查询结果失真。修复后无论 Dag Run 因何失败正常失败或超时失败dagrun.duration.failed都携带一致的dag_idrun_type标签组合可按触发来源拆分告警借助run_type标签可以针对scheduled、manual、asset_triggered等不同来源分别设置超时耗时告警阈值例如仅对调度触发的超时提高关注度与正常完成指标对齐dagrun.duration.success、dagrun.duration.failed等指标现在共享同一套stats_tags逻辑监控面板的标签筛选条件可以统一编写无需针对超时场景做特殊处理。结语本次修复虽然只是一处标签字段的补齐却消除了dagrun.duration.failed指标在两条发射路径上的语义分裂超时路径通过展开DagRun.stats_tags获得run_type标签与_emit_duration_stats_for_finished_state的正常路径保持一致并配有精确断言标签集合的单元测试兜底。对于所有基于 Airflow 指标做 SLA 监控与超时分析的用户而言这保证了数据口径的完整与可对比性。相关实现与测试均可在当前仓库中直接查阅DagRun.stats_tags、超时发射代码、DagRunType 枚举 以及 验收测试。【免费下载链接】airflowApache Airflow - A platform to programmatically author, schedule, and monitor workflows项目地址: https://gitcode.com/GitHub_Trending/ai/airflow创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
返回列表