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

Airflow技术问询:任务失败时如何重启整个DAG而非单个任务

Airflow 实现任意任务失败时重启整个DAG

可行,你可以通过失败回调函数 + Airflow API客户端的方式实现这个需求,当任意任务失败时触发整个DAG从头开始执行,具体实现如下:

核心思路

关闭单个任务的重试机制,给所有任务绑定一个失败回调函数,当任务失败时,通过Airflow的本地API客户端清除当前DAG运行的所有任务状态,并触发新的DAG运行,从而实现从task1开始重新执行。

具体代码实现

1. 导入依赖模块

from airflow import DAG
from airflow.operators.python import PythonOperator
from airflow.utils.dates import days_ago
from datetime import timedelta
from airflow.api.client.local_client import Client

2. 定义失败回调函数

这个函数会在任务失败时被调用,负责清除当前运行的任务状态并触发新的DAG运行:

def trigger_full_dag_retry(context):
    # 获取当前DAG的ID和执行时间
    dag_id = context["dag"].dag_id
    execution_date = context["execution_date"]

    # 初始化Airflow本地客户端
    client = Client(None, None)

    # 清除当前DAG运行的所有任务实例状态,确保新运行能从头执行
    client.clear_task_instances(
        dag_id=dag_id,
        execution_date=execution_date,
        only_failed=False,
        only_running=False,
        include_upstream=False,
        include_downstream=False,
        reset_dag_runs=True
    )

    # 触发当前DAG的新运行
    client.trigger_dag(dag_id=dag_id)

3. 配置DAG的default_args

关闭单个任务的重试,同时将回调函数加入默认参数:

default_args = {
    "owner": "testing",
    "retries": 0,  # 关闭单个任务重试,避免和DAG重启冲突
    "retry_delay": timedelta(minutes=1),
    "on_failure_callback": trigger_full_dag_retry  # 绑定失败回调
}

4. 定义DAG和任务依赖

with DAG(
    dag_id="full_retry_test_dag",
    default_args=default_args,
    schedule_interval="@daily",
    start_date=days_ago(1),
    catchup=False
) as dag:
    task1 = PythonOperator(
        task_id="task1",
        python_callable=lambda: print("Executing task1")
    )

    task2 = PythonOperator(
        task_id="task2",
        python_callable=lambda: print("Executing task2")
    )

    task3 = PythonOperator(
        task_id="task3",
        python_callable=lambda: 1/0  # 模拟任务失败场景
    )

    task4 = PythonOperator(
        task_id="task4",
        python_callable=lambda: print("Executing task4")
    )

    # 设置任务依赖
    task1 >> task2 >> task3 >> task4

注意事项

  • 必须将单个任务的retries设为0,否则任务会先执行自身重试,再触发DAG重启,不符合你"直接重启整个DAG"的需求。
  • 运行Airflow的进程用户需要具备触发DAG和清除任务实例的权限,否则回调函数会执行失败。
  • 清除任务实例状态是为了避免新触发的DAG运行跳过已完成的任务,确保从task1开始从头执行。

内容的提问来源于stack exchange,提问作者awakener1986

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.16 19:30:15