如何终止由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
相关产品推荐
相关产品推荐

