Airflow 1.10:如何通过DagRun对象清除调度运行并触发失败任务?
在Airflow 1.10中通过DagRun对象清除失败运行实例
在Airflow 1.10版本中,完全可以用你获取到的DagRun对象来清除对应的失败运行实例并触发重新调度。
具体实现代码
假设你已经通过DagRun.find()方法拿到了目标失败实例,直接调用对象的clear()方法即可完成操作:
from airflow.models import DagRun # 获取指定DAG的所有失败运行实例 failed_dag_runs = DagRun.find(dag_id="your_target_dag_id", state="failed") # 遍历并逐个清除 for dag_run in failed_dag_runs: # 清除该运行实例,重置状态以触发重新调度 dag_run.clear( start_date=dag_run.execution_date, end_date=dag_run.execution_date, dag_id=dag_run.dag_id, include_subdags=False, reset_dag_runs=True )
关键参数说明
reset_dag_runs=True:强制重置DagRun的状态,允许Airflow重新调度这个运行实例start_date/end_date:指定为当前dag_run的执行日期,确保只清除这个特定的实例,不会影响其他运行记录include_subdags:如果你的DAG包含子DAG,可根据需求设置是否一并清除子DAG的运行实例
注意事项
- 执行此操作需要具备Airflow的管理员级权限
- 代码需在Airflow的运行环境中执行(比如Airflow部署的Python环境、通过
airflow python命令运行脚本)
内容的提问来源于stack exchange,提问作者Aman Singh
相关产品推荐
相关产品推荐

