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

如何终止由TriggerDagRunOperator启动的子Dag?

如何终止TriggerDagRunOperator启动的子DAG?

TriggerDagRunOperator本身没有内置的终止子DAG的方法,因为子DAG是独立的Airflow运行实例。要实现父DAG失败时终止子DAG,可以通过以下几种方案实现:


方案1:父DAG失败回调 + Airflow核心API直接操作实例

通过父DAG的失败回调函数,利用Airflow的数据库模型找到对应执行日期的子DAG实例,将其标记为失败并终止所有运行中的任务。

代码示例

from airflow.models import DagRun, TaskInstance
from airflow.utils.state import State
from airflow.settings import Session
from airflow import DAG
from airflow.operators.trigger_dagrun import TriggerDagRunOperator
from datetime import datetime

def terminate_child_dag_on_parent_failure(context):
    # 获取父DAG的执行日期(与子DAG的执行日期一致)
    parent_exec_date = context["logical_date"]
    child_dag_id = "child_dag"

    session = Session()
    try:
        # 定位目标子DAG运行实例
        child_dag_run = session.query(DagRun).filter(
            DagRun.dag_id == child_dag_id,
            DagRun.execution_date == parent_exec_date
        ).first()

        if child_dag_run and child_dag_run.state not in [State.SUCCESS, State.FAILED]:
            # 标记子DAG为失败状态
            child_dag_run.state = State.FAILED
            session.commit()

            # 终止子DAG中所有正在运行的任务
            running_tasks = session.query(TaskInstance).filter(
                TaskInstance.dag_id == child_dag_id,
                TaskInstance.execution_date == parent_exec_date,
                TaskInstance.state == State.RUNNING
            ).all()

            for task in running_tasks:
                task.state = State.FAILED
                session.commit()
    finally:
        session.close()

# 父DAG配置
default_args = {
    "owner": "airflow",
    "start_date": datetime(2024, 1, 1),
    "on_failure_callback": terminate_child_dag_on_parent_failure  # 绑定失败回调
}

with DAG(
    dag_id="parent_dag",
    default_args=default_args,
    schedule_interval="@daily"
) as dag:
    my_trigger = TriggerDagRunOperator(
        task_id="my_trigger",
        trigger_dag_id="child_dag",
        wait_for_completion=True,
        execution_date="{{ logical_date }}"
    )

    # 父DAG后续任务示例
    # some_other_task = ...

注意事项

  • 确保Airflow Worker进程有数据库访问权限(对应Airflow元数据库的读写权限);
  • 子DAG的执行日期必须与父DAG的logical_date严格匹配(你的代码中已经通过execution_date="{{ logical_date }}"实现,所以能准确定位)。

方案2:失败回调中调用Airflow CLI命令

如果不想直接操作数据库模型,可以在回调中调用Airflow官方CLI命令来标记子DAG为失败,CLI会自动处理终止运行任务的逻辑。

代码示例

import subprocess
from airflow.utils.state import State

def terminate_child_dag_via_cli(context):
    parent_exec_date = context["logical_date"].strftime("%Y-%m-%dT%H:%M:%S%z")
    child_dag_id = "child_dag"

    # 调用airflow dags fail命令终止子DAG
    try:
        subprocess.run(
            ["airflow", "dags", "fail", child_dag_id, "-e", parent_exec_date],
            check=True,
            capture_output=True
        )
    except subprocess.CalledProcessError as e:
        # 可选:处理命令执行失败的情况
        print(f"Failed to terminate child DAG: {e.stderr.decode()}")

注意事项

  • 运行Worker的系统用户必须有权限执行airflow命令;
  • 需确保Worker环境的Airflow配置正确(如AIRFLOW_HOME环境变量已设置)。

方案3:自定义可终止的TriggerDagRunOperator

将终止逻辑封装成自定义Operator,继承原生TriggerDagRunOperator,方便复用。

代码示例

from airflow.operators.trigger_dagrun import TriggerDagRunOperator
from airflow.models import DagRun, TaskInstance
from airflow.utils.state import State
from airflow.settings import Session

class TriggerDagRunWithTerminateOperator(TriggerDagRunOperator):
    def terminate_child_dag(self, execution_date):
        session = Session()
        try:
            child_dag_run = session.query(DagRun).filter(
                DagRun.dag_id == self.trigger_dag_id,
                DagRun.execution_date == execution_date
            ).first()

            if child_dag_run and child_dag_run.state not in [State.SUCCESS, State.FAILED]:
                child_dag_run.state = State.FAILED
                session.commit()

                # 终止运行中的任务
                running_tasks = session.query(TaskInstance).filter(
                    TaskInstance.dag_id == self.trigger_dag_id,
                    TaskInstance.execution_date == execution_date,
                    TaskInstance.state == State.RUNNING
                ).all()

                for task in running_tasks:
                    task.state = State.FAILED
                    session.commit()
        finally:
            session.close()

# 父DAG中使用自定义Operator
def parent_failure_callback(context):
    trigger_task = context["task_instance"].task
    if isinstance(trigger_task, TriggerDagRunWithTerminateOperator):
        trigger_task.terminate_child_dag(context["logical_date"])

default_args = {
    "on_failure_callback": parent_failure_callback
}

with DAG(dag_id="parent_dag", default_args=default_args) as dag:
    my_trigger = TriggerDagRunWithTerminateOperator(
        task_id="my_trigger",
        trigger_dag_id="child_dag",
        wait_for_completion=True,
        execution_date="{{ logical_date }}"
    )

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.28 15:08:16