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

如何在Airflow中限制一组DAG的最大活跃运行数为1?

解决DAG组互斥运行(含Backfill场景)的其他方案

针对你遇到的Backfill场景下ExternalTaskSensor失效的问题,除了pool、优先级权重和自定义Operator检查外,还有以下几种实用方案:

1. Airflow全局变量+ShortCircuitOperator互斥

借助Airflow内置的Variable实现全局锁逻辑,每个DAG启动前先检查并抢占锁,执行完成后释放:

  • 核心逻辑:
    • 用ShortCircuitOperator在DAG启动时检查active_dag变量,若为空则写入当前DAG ID并继续执行;若变量存在且不是当前DAG,则中断任务。
    • 用PythonOperator在DAG所有任务结束(无论成功失败)后,清空active_dag变量。
    • 开启serialize_writes=True保证变量修改的原子性,避免并发竞争;同时可以加个定时清理DAG,定期检查锁中记录的DAG是否真的在运行,防止意外终止导致锁一直占用。
  • 代码示例:
from airflow.models import Variable
from airflow.operators.python import ShortCircuitOperator, PythonOperator
from airflow.utils.session import create_session

def check_and_acquire_lock():
    with create_session() as session:
        current_running = Variable.get("active_dag", default_var=None, session=session)
        if not current_running:
            Variable.set("active_dag", "{{ dag.dag_id }}", serialize_writes=True, session=session)
            return True
        # 允许当前DAG的重跑/Backfill任务继续执行
        return current_running == "{{ dag.dag_id }}"

def release_lock(**context):
    with create_session() as session:
        current_running = Variable.get("active_dag", default_var=None, session=session)
        if current_running == context["dag"].dag_id:
            Variable.delete("active_dag", session=session)

# 在你的DAG中添加以下任务
lock_check = ShortCircuitOperator(
    task_id="check_running_dag_lock",
    python_callable=check_and_acquire_lock,
    dag=dag
)

lock_release = PythonOperator(
    task_id="release_dag_lock",
    python_callable=release_lock,
    trigger_rule="all_done",  # 无论任务成功/失败都释放锁
    dag=dag
)

# 设置依赖:lock_check >> 你的核心任务 >> lock_release

2. 分布式锁(Redis为例)

如果是多Scheduler节点的集群环境,用Redis分布式锁比Airflow Variable更可靠,能避免单点竞争问题:

  • 核心逻辑:
    • DAG启动时尝试获取Redis锁,设置合理的过期时间(比DAG最长运行时间长),防止意外终止导致锁无法释放。
    • Backfill任务同样会触发锁检查,天然支持历史任务的互斥。
  • 代码示例:
import redis
from airflow.operators.python import PythonOperator

def acquire_redis_lock(**context):
    r = redis.Redis(host="your-redis-host", port=6379, db=0)
    lock_key = "dag_group_mutex_lock"
    # nx=True表示只有锁不存在时才设置,ex=86400是锁过期时间(24小时)
    lock_acquired = r.set(lock_key, context["dag"].dag_id, ex=86400, nx=True)
    if not lock_acquired:
        owner_dag = r.get(lock_key).decode() if r.get(lock_key) else "unknown"
        raise RuntimeError(f"互斥锁已被DAG {owner_dag}占用,当前任务终止")

def release_redis_lock(**context):
    r = redis.Redis(host="your-redis-host", port=6379, db=0)
    lock_key = "dag_group_mutex_lock"
    current_owner = r.get(lock_key).decode() if r.get(lock_key) else None
    if current_owner == context["dag"].dag_id:
        r.delete(lock_key)

# DAG内任务配置
acquire_lock_task = PythonOperator(
    task_id="acquire_redis_mutex",
    python_callable=acquire_redis_lock,
    dag=dag
)

release_lock_task = PythonOperator(
    task_id="release_redis_mutex",
    python_callable=release_redis_lock,
    trigger_rule="all_done",
    dag=dag
)

# 依赖设置:acquire_lock_task >> 核心任务 >> release_lock_task

3. 自定义DagRunListener

利用Airflow的DagRunListener钩子,在DAG运行启动时自动检查同组其他DAG的状态,实现全局拦截:

  • 核心逻辑:
    • 编写自定义Listener,继承DagRunListener,重写on_dag_run_start方法,查询同组内是否有正在运行/排队的DagRun,若有则标记当前DagRun为失败。
    • 该方式对Backfill任务完全生效,因为Backfill会生成新的DagRun,Listener会在其启动时触发检查。
  • 代码示例:
from airflow.listeners.listener import DagRunListener
from airflow.models import DagRun
from airflow.utils.session import create_session

class MutexDagGroupListener(DagRunListener):
    # 定义需要互斥的DAG组ID列表
    TARGET_DAG_GROUP = ["dag_a", "dag_b", "dag_c"]

    def on_dag_run_start(self, dag_run: DagRun, msg: str) -> None:
        if dag_run.dag_id not in self.TARGET_DAG_GROUP:
            return
        with create_session() as session:
            # 查询同组内其他处于运行/排队状态的DagRun
            running_count = session.query(DagRun).filter(
                DagRun.dag_id.in_(self.TARGET_DAG_GROUP),
                DagRun.dag_id != dag_run.dag_id,
                DagRun.state.in_(["running", "queued"])
            ).count()
            if running_count > 0:
                dag_run.state = "failed"
                dag_run.set_state("failed", session=session)
                session.commit()

# 在airflow.cfg中注册Listener:
# listener_classes = your.package.path.MutexDagGroupListener

4. CeleryExecutor专属:单并发队列

如果你的Airflow使用CeleryExecutor,可以给互斥DAG组指定一个专属队列,并且只启动1个worker消费该队列:

  • 操作步骤:
    1. 在airflow.cfg中配置任务路由,把目标DAG的所有任务路由到mutex_queue:
      [celery]
      task_routes = {"dag_a.*": {"queue": "mutex_queue"}, "dag_b.*": {"queue": "mutex_queue"}}
      
    2. 启动Celery Worker时指定队列和并发数:
      airflow celery worker -q mutex_queue -c 1
      
  • 特点:Celery会保证该队列同一时间只有1个任务在执行,间接实现DAG互斥;缺点是DAG内的任务会串行执行,适合任务量较小的场景。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.01 04:46:48