Airflow能否标记DagRun为失败或终止?如何编程实现该操作
Airflow 标记DagRun为失败/终止运行中DagRun的实现方法
Airflow原生支持将DagRun标记为失败状态、终止运行中DagRun的操作,可通过Python API、REST API、UI操作三种方式实现,编程场景可直接用下面的方法:
1. 内置Python API 编程实现
- DAG任务内部终止当前DagRun
在任务逻辑里满足终止条件时,直接抛出AirflowFailException即可,Airflow会自动将当前DagRun标记为失败状态,无需额外操作:
from airflow.exceptions import AirflowFailException def task_logic(**context): # 自定义终止判断逻辑 if meet_terminate_condition: raise AirflowFailException("触发终止条件,当前DagRun标记为失败")
如果需要直接控制DagRun状态、不等待任务调度链路走完,可以直接操作DagRun模型对象,同时清理残留的未完成任务:
from airflow.models.dagrun import DagRun from airflow.utils.state import State def stop_target_dagrun(dag_id: str, run_id: str): # 查询匹配的目标DagRun target_dagrun = DagRun.find(dag_id=dag_id, run_id=run_id)[0] # 将DagRun状态置为失败 target_dagrun.set_state(State.FAILED) # 同步将该DagRun下所有未完成的任务标记为失败,避免残留任务继续占用资源 unfinished_tasks = target_dagrun.get_task_instances( state=[State.RUNNING, State.QUEUED, State.SCHEDULED, State.UP_FOR_RETRY] ) for ti in unfinished_tasks: ti.set_state(State.FAILED)
注意:如果是在Airflow服务进程外(比如独立的Python脚本)调用上述代码,需要先加载Airflow配置、初始化元数据库连接,否则会出现连接报错。
- 批量终止DagRun
如果需要批量终止符合条件的DagRun,只需要给DagRun.find()传入对应的过滤参数(比如执行时间范围、DagRun状态、DAG标签等),遍历查询结果执行上述状态更新逻辑即可。
2. REST API 远程调用实现
如果不想直接连接Airflow元数据库,或者需要跨服务调用,可以用Airflow 2.0+自带的原生REST API实现:
- 调用接口:
PATCH /api/v1/dags/{dag_id}/dagRuns/{dag_run_id} - 请求体参数:
{"state": "failed"} - 调用前需要提前创建拥有DagRun编辑权限的账号,通过Basic Auth或JWT Token做鉴权,接口会自动处理DagRun状态更新、关联任务的终止清理,不会出现状态不一致的问题。
3. UI手动操作(非编程场景)
临时处理单个DagRun时,可以直接在Airflow Web UI的DagRun列表页找到目标运行记录,点击操作栏的「标记为失败」按钮,系统会自动终止该DagRun下所有运行、排队中的任务,将DagRun状态更新为失败。
注意:不要直接手动修改Airflow元数据库中dag_run表的state字段,跳过Airflow内置的状态更新逻辑会导致DagRun和下属任务实例状态不一致,出现残留任务占用资源、后续调度异常的问题。
内容的提问来源于stack exchange,提问作者Jiew Meng
相关产品推荐
相关产品推荐

