Apache Airflow 实战:利用 on_failure_callback 与 REST API 实现 Dag 级重试(Dag-level Retry)
【免费下载链接】airflowApache Airflow - A platform to programmatically author, schedule, and monitor workflows项目地址: https://gitcode.com/GitHub_Trending/ai/airflow
导读
Apache Airflow 内置的重试机制是任务级(task-level)的:每个任务拥有独立的retries与retry_delay,失败时只有该任务本身会被重新执行。但在真实生产环境中,有些工作流的最小重试单位是整个 Dag run——例如多步骤流水线的中间状态难以部分清理时,从头重跑比重跑单个失败任务更安全。本文基于 airflow-core/docs/howto/dag-level-retry-via-callback.rst,讲解一种"接近 Dag 级重试"的成熟配方:把on_failure_callback与 Airflow 公开 REST API 的 Clear 端点组合起来,在任务失败后清除失败的 Dag run,让调度器重新执行。读完本文,你将掌握该配方的适用场景、完整可运行的代码、重试次数限制策略,以及必须警惕的各类边界情况。
为什么需要 Dag 级重试:任务级重试的局限
Airflow 的设计哲学是"任务应当被设计成任务级幂等(idempotent at the task level)",这样内置的按任务重试机制就足够用了。这是推荐的起点,绝大多数工作流不需要超出它的东西。
但存在这样一类场景:即便任务已经做到幂等,你仍然希望一个更粗粒度的重试单元——本质上是 Dag 级重试。典型例子是:
- 工作单元天然是一个 Dag run,而不是单个任务;
- 多步骤流水线中间状态难以部分清理,从头重跑比重跑单个失败任务更安全。
此时可以组合使用 Airflow 的两大既有原语:
on_failure_callback——任务失败后的回调;- 公开 REST API 的 Clear 端点——
POST /api/v2/dags/{dag_id}/dagRuns/{dag_run_id}/clear。
当某个任务失败时,回调调用 Clear 端点清除失败的 Dag run,调度器随后会重新运行它,从而在效果上实现了"整个 Dag run 重试"。
重要定位:这是一个构建在既有原语之上的 recipe(配方),不是新特性。它有明确的观点倾向和取舍,使用前务必通读文末 Caveats 一节。
社区正在讨论一种更通用的构造——有时被称为transactional task group——作为 Airflow dev list 上的潜在新增特性。如果你的场景受益于此,请参与社区讨论,而不要把本配方作为长期依赖方案。
何时该用、何时不该用
适合使用的情况
- 工作单元天然是 Dag run,而非单个任务;
- 重跑任务是幂等的——再次运行不会产生重复副作用、重复扣费或不一致的外部状态;
- 有限次数的整体重跑在时间与成本上是可以接受的。
应避免的情况
- 任务有非幂等副作用(发送邮件、扣款、向不去重的 API 提交数据),且尚未做成幂等;
- 重试需要对下游 Dag 或资产透明——清除一个 Dag run 产生的是全新的一次 attempt(尝试),而不是对原始运行的"重试";
- 简单的按任务
retries设置已经覆盖了你实际遇到的失败模式。
工作原理:从失败到重新调度的完整链路
Dag 级重试的完整流程如下:
- 某个任务在耗尽自身任务级重试后失败;
- 任务的
on_failure_callback在dag processor中执行; - 回调判断是否允许再次进行 Dag 级尝试(使用你自己维护的计数器,见下文"限制重试次数");
- 若允许,回调调用 Airflow REST API 清除失败的 Dag run;
- 调度器拾取被清除的任务实例并重新执行。
这里用到的 REST 端点是POST /api/v2/dags/{dag_id}/dagRuns/{dag_run_id}/clear。该端点在 OpenAPI 规范 v2-rest-api-generated.yaml 中定义,路由实现位于 dag_run.py。
从源码看,clear_dag_run端点接收DAGRunClearBody请求体,最终调用perform_clear_dag_run服务(见 dag_run.py),后者在内部执行dag.clear(...)并把任务实例重置,使其可被重新调度。请求体支持以下字段(定义于 datamodels/dag_run.py):
| 字段 | 默认值 | 含义 |
|---|---|---|
dry_run | true | 若为true只返回将被清除的任务实例列表,不真正清除;配方中必须设为false才会实际生效 |
only_failed | false | 只清除失败(FAILED/UPSTREAM_FAILED)的任务实例 |
only_new | false | 只排队最新 Dag 版本中新增的任务,而不清除既有任务(与only_failed互斥) |
run_on_latest_version | null | (实验性)清除后在最新 bundle 版本上运行;未指定时依次回退到 DAG 级rerun_with_latest_version参数、[core] rerun_with_latest_version配置项,最终为False |
note | null | 附加备注,最大长度 1000 |
前置条件
使用该配方前,需要满足:
- dag processor 对 Airflow REST API 的认证访问,且拥有对目标 Dag 清除 Dag run 的权限。从源码可见,
clear_dag_run路由依赖requires_access_dag(method="PUT", access_entity=DagAccessEntity.RUN)(见 dag_run.py),即需要对该 Dag 的 RUN 实体具备 PUT 权限; - dag processor 到 API server 的网络可达性。
注意:
on_failure_callback只在失败经由正常任务执行发生时触发。通过 UI、CLI 手工进行的状态变更不会触发它。这一点与 Airflow 回调的整体行为一致——关于回调触发机制的详细说明见 Callbacks 文档。
基础示例:清除失败的 Dag run
下面这个示例会在 Dag 中任意任务耗尽任务级重试并失败后,清除整个 Dag run。任务上设置retries=0,使得回调在第一次失败时就触发,让整个 Dag run接管重试责任。
from urllib.parse import quote import requests from airflow.sdk import DAG, task AIRFLOW_API_BASE = "https://airflow.example.com/api/v2" # your deployment's API base def clear_dag_run_on_failure(context): dag_run = context["dag_run"] dag_id_path = quote(dag_run.dag_id, safe="") run_id_path = quote(dag_run.run_id, safe="") response = requests.post( f"{AIRFLOW_API_BASE}/dags/{dag_id_path}/dagRuns/{run_id_path}/clear", headers={"Authorization": "Bearer <token>"}, json={"dry_run": False}, timeout=30, ) response.raise_for_status() with DAG( dag_id="example_dag_level_retry", default_args={ "retries": 0, # let the Dag-level retry take over "on_failure_callback": clear_dag_run_on_failure, }, ): @task def step_one(): ... @task def step_two(value): ... step_two(step_one())代码要点说明:
retries=0:让 Dag 级重试接管。如果保留任务级重试,回调要等到任务级重试耗尽后才触发;- URL 编码:
dag_id与run_id用quote(..., safe="")做严格 URL 编码,避免特殊字符破坏路径; dry_run: False:真正执行清除。请求体同时支持only_failed、only_new、run_on_latest_version、note等字段,可按需扩展;timeout=30:为网络请求设置超时,避免回调长期挂起。
警告:这个最简版本在失败时无条件清除。不要在生产环境原样运行——如果某个任务持续失败,它会无限循环。部署前必须加上尝试次数上限,见下一节。
限制重试次数:基于 Variable 的计数器
一个朴素的、总是清除的回调会在任务持续失败时无限循环。注意dag_run.conf在 run 创建时被设置,清除不会改变它,因此不能用作重试计数器——你必须自己跟踪尝试次数。一个简单做法是使用作用域限定到该 run 的 Variable:
from airflow.sdk import Variable def _attempts_key(dag_id: str, run_id: str) -> str: return f"dag_run_attempts::{dag_id}::{run_id}" def _read_attempts(dag_id: str, run_id: str) -> int: return int(Variable.get(_attempts_key(dag_id, run_id), default=0)) def _bump_attempts(dag_id: str, run_id: str) -> int: new_value = _read_attempts(dag_id, run_id) + 1 Variable.set(_attempts_key(dag_id, run_id), str(new_value)) return new_value在决定是否清除之前调用_read_attempts读取当前次数;在调用 API之前调用_bump_attempts递增。当 Dag run 结束时——无论成功还是放弃——都要重置计数器,避免 Variable 不断累积:
def cleanup_attempts(context): dag_run = context["dag_run"] Variable.delete(_attempts_key(dag_run.dag_id, dag_run.run_id))将cleanup_attempts挂到 Dag 级的on_success_callback,并在on_failure_callback的"放弃重试"分支中、返回之前调用它。同时建议在达到上限时发出告警,而不是静默放弃。
并发的非原子性
注意:这种基于 Variable 的计数器在并发回调下不是原子的。如果你的 Dag 存在并行分支且可能同时失败,两个回调可能读到同一个值并都决定清除。实践中第二次清除是空操作(第一次已经重置了 run),但你可能会看到重复的 API 调用。如果需要更严格的保证,请使用带 compare-and-swap 语义的后端。
为自动化修复(Remediation)Dag 加护栏
Dag 级重试模式同样适用于自动化修复 Dag——这类 Dag 会对外部系统执行操作,例如清除失败的工作、重放消息、重启外部作业。这些操作应当有护栏,防止一个持久性故障导致相同的恢复动作在多个 Dag run 中反复执行。
Pools(池)可以限制同时运行的恢复任务数量,但无法限制同一恢复动作被尝试的频率。对于自动化修复工作流,建议为每个目标添加:
- 一个小的冷却时间标记(cooldown marker);
- 一个最大自动操作次数,超过后转入人工审查。
对于较小的护栏值——如冷却时间戳或尝试计数器——Airflow Variable 可能足够。但不要把 Variable 当作高吞吐状态存储,或当作强一致性的锁机制。如果工作流需要原子更新、严格的并发控制或大量按目标记录的状态,请使用为此设计的外部存储。
Caveats:需要警惕的边界情况
使用该配方前,请逐一核对以下取舍:
UP_FOR_RETRY等待间隔:失败的任务会在UP_FOR_RETRY状态等待retry_delay(默认 5 分钟)后,才转为FAILED并触发下一次清除。因此每次 Dag 级重试至少被延迟这个间隔;- 幂等是你的责任:每个被清除的任务都会再次运行。失败尝试中产生过的任何外部副作用——写入的行、发送的消息、上传的文件——都会再次产生,除非你的任务幂等(例如使用 upsert、去重键或条件写入);
- 循环风险:计数器或失败任务中的 bug 可能导致无限的"清除—重试"循环。务必限制尝试次数,并在达到上限时告警;
- 成本:一次 Dag 级重试会从零开始重跑所有被清除的任务。请确保时间与资源成本可接受,并把尝试上限设低;
- 正在运行的任务:如果你的 Dag 有并行分支,一个任务失败时同一 Dag run 中的其他任务可能仍在运行。请设计 Dag 及其副作用,确保在其它工作进行中清除 run 是安全的;
- Token 生命周期:大多数认证管理器签发的 token 会过期(例如 simple auth manager 的 JWT 默认 24 小时)。存储在 Variable 中的静态 token 过期后会静默返回 401;循环防护能阻止错误行为,但 Dag 级重试会停止工作。对于长期运行的部署,请定期刷新凭据,或使用更长生命周期的服务凭据。
小结
本文给出的配方用约 20 行回调代码加一个 REST 调用,把"任务级失败"升级为"整个 Dag run 重试",适用于工作单元天然是 Dag run、任务幂等且有界重跑可接受的场景。核心链路是:任务失败 → dag processor 中的on_failure_callback→POST /api/v2/dags/{dag_id}/dagRuns/{dag_run_id}/clear→ 调度器重新调度。落地时务必带上基于 Variable 的尝试计数与清理逻辑、对自动化修复场景的护栏,并认真对待UP_FOR_RETRY延迟、幂等性、循环风险、成本、并发运行与 token 过期这六类边界情况。更通用的transactional task group特性正在社区讨论中,长期方案建议跟随社区进展。
【免费下载链接】airflowApache Airflow - A platform to programmatically author, schedule, and monitor workflows项目地址: https://gitcode.com/GitHub_Trending/ai/airflow
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考