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

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配置与事件触发优化

  1. 给dag2设置max_active_runs=1,确保同一时间只有一个实例处于活跃状态;
  2. 在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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.26 10:20:10