You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.08.26 19:09:23