如何在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会在其启动时触发检查。
- 编写自定义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消费该队列:
- 操作步骤:
- 在
airflow.cfg中配置任务路由,把目标DAG的所有任务路由到mutex_queue:[celery] task_routes = {"dag_a.*": {"queue": "mutex_queue"}, "dag_b.*": {"queue": "mutex_queue"}} - 启动Celery Worker时指定队列和并发数:
airflow celery worker -q mutex_queue -c 1
- 在
- 特点:Celery会保证该队列同一时间只有1个任务在执行,间接实现DAG互斥;缺点是DAG内的任务会串行执行,适合任务量较小的场景。
内容的提问来源于stack exchange,提问作者Dark Matter
相关产品推荐
相关产品推荐

