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

任务失败时能否通过on_failure_callback触发Airflow DAG?

如何通过on_failure_callback触发另一个DAG

完全可以通过on_failure_callback实现任务失败时触发另一个DAG,核心思路是在回调函数中调用Airflow内置的DAG触发逻辑,替代TriggerDagRunOperator的触发方式。

实现步骤(Airflow 2.x)

  1. 导入所需模块
  2. 定义失败回调函数,在其中调用trigger_dagrun触发目标DAG
  3. 在任务中配置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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.15 21:09:49