Airflow中如何实现DAG运行的抢占式触发?
Airflow实现DAG运行抢占的方案
Airflow本身没有内置的“自动抢占”配置,但可以通过以下几种方式实现你需要的「终止旧运行实例、让新实例立即启动」的需求:
方法一:在触发新实例前主动终止旧实例
在管理DAG(dag1)中,先通过PythonOperator调用Airflow内部API,查询并终止目标DAG(dag2)的所有活跃运行实例,再触发新的dag2实例。
示例代码:
from airflow import DAG from airflow.operators.python import PythonOperator from airflow.operators.trigger_dagrun import TriggerDagRunOperator from airflow.api.client.local_client import Client from airflow.utils.state import State from datetime import datetime def terminate_running_dag2(): # 初始化本地客户端(适用于Airflow 2.x) client = Client(None, None) # 获取dag2所有处于RUNNING状态的实例 running_runs = client.get_dag_runs( dag_id="dag2", states=[State.RUNNING] ) # 逐个终止实例(可根据需求选择标记为FAILED或SUCCESS) for run in running_runs: client.set_dag_run_state( dag_run_id=run.dag_run_id, state=State.FAILED, message="被新触发的实例抢占终止" ) with DAG( dag_id="dag1", schedule=None, # 按需触发 start_date=datetime(2024, 1, 1), catchup=False ) as dag: terminate_task = PythonOperator( task_id="terminate_dag2_running_instances", python_callable=terminate_running_dag2 ) trigger_task = TriggerDagRunOperator( task_id="trigger_dag2", trigger_dag_id="dag2", wait_for_completion=False ) terminate_task >> trigger_task
方法二:结合DAG配置与事件触发优化
- 给dag2设置
max_active_runs=1,确保同一时间只有一个实例处于活跃状态; - 在dag1触发新实例前,先终止旧实例(同方法一),这样新实例会立即从排队状态转为运行状态,无需等待旧实例自然结束。
注意事项
- 终止运行中的任务可能导致数据中间状态不一致,需确保dag2的任务具备幂等性(即重复执行不会产生错误结果),或有完善的回滚/清理机制;
- 若使用Airflow的远程执行环境(如CeleryExecutor),需确保
local_client的权限足够操作DAG运行状态; - 对于Airflow 1.x版本,API调用方式略有不同,需使用
airflow.api.common.experimental下的方法(如get_dag_runs、set_dag_run_state)。
内容的提问来源于stack exchange,提问作者FrustratedWithFormsDesigner
相关产品推荐
相关产品推荐

