任务失败时能否通过on_failure_callback触发Airflow DAG?
如何通过on_failure_callback触发另一个DAG
完全可以通过on_failure_callback实现任务失败时触发另一个DAG,核心思路是在回调函数中调用Airflow内置的DAG触发逻辑,替代TriggerDagRunOperator的触发方式。
实现步骤(Airflow 2.x)
- 导入所需模块
- 定义失败回调函数,在其中调用
trigger_dagrun触发目标DAG - 在任务中配置
on_failure_callback参数
代码示例
from airflow import DAG from airflow.operators.python import PythonOperator from airflow.api.common.trigger_dagrun import trigger_dagrun from datetime import datetime # 定义失败时触发目标DAG的回调函数 def trigger_failure_handler(context): # 触发指定DAG,并传递失败上下文信息 trigger_dagrun( dag_id="target_failure_handling_dag", # 要触发的目标DAG ID conf={ "failed_dag_id": context["dag"].dag_id, "failed_task_id": context["task_instance"].task_id, "execution_date": str(context["execution_date"]) }, # 传递失败相关元数据 execution_date=context["execution_date"], replace_microseconds=False ) # 主DAG定义 with DAG( dag_id="primary_workflow_dag", start_date=datetime(2024, 1, 1), schedule_interval="@daily", catchup=False ) as dag: # 模拟会失败的任务 def simulate_failure(): raise RuntimeError("任务执行出错,触发失败处理流程") failing_task = PythonOperator( task_id="critical_task", python_callable=simulate_failure, on_failure_callback=trigger_failure_handler # 绑定失败回调 )
关键说明
- 确保目标DAG(
target_failure_handling_dag)已创建并处于启用状态 context参数包含了任务实例、DAG信息、执行日期等上下文数据,可通过conf参数传递给被触发的DAG,用于后续的失败排查或处理- 若使用Airflow 1.x,可通过
airflow.models.DagRun手动创建DAG运行记录来实现触发,核心逻辑类似
内容的提问来源于stack exchange,提问作者Dylan García
相关产品推荐
相关产品推荐

