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
相关产品推荐
相关产品推荐

