如何避免Airflow中两个DAG同时运行并解决死锁问题?
解决两个Airflow DAG互斥运行的死锁问题
死锁原因分析
你之前的实现逻辑是让每个DAG等待对方最近一次执行完成,但如果两个DAG刚好同时启动,它们各自锁定的「最近一次执行」就是对方当前正在运行的实例,结果就是互相等待、谁也不让谁,直接形成死锁。
可行解决方案
方案1:给其中一个DAG设置「优先级让步」逻辑
让运行频率更低的DAG(比如DAG B)在检测到冲突时主动放弃本次调度,而非死等。这种方式简单直接,能从根源避免互相僵持。
修改DAG B的传感器逻辑,用PythonSensor替代原有的ExternalTaskSensor:
from airflow.sensors.python import PythonSensor from airflow.models import DagRun, settings from airflow.utils.state import State def check_dag_a_not_running(): session = settings.Session() # 检查DAG A是否有正在运行的实例 running_dag_a = session.query(DagRun)\ .filter(DagRun.dag_id == 'DAG_A', DagRun.state == State.RUNNING)\ .first() session.close() # 返回True表示可以继续,False则继续等待/触发跳过 return not bool(running_dag_a) sensor = PythonSensor( task_id='wait_for_dag_a_to_finish', python_callable=check_dag_a_not_running, mode='reschedule', poke_interval=60, # 每分钟检查一次状态 dag=dag )
DAG A的传感器可以保留原有逻辑(等待DAG B完成),因为DAG B运行频率低,即便同时启动,DAG A会先获取运行权限,DAG B只会等待不会反向阻塞。
方案2:使用全局共享锁
借助Airflow元数据库或Redis实现全局互斥锁,确保同一时间只有一个DAG能执行。这种方式更严谨,适合对互斥要求极高的场景。
示例:基于数据库行级锁实现互斥操作
from airflow.models import BaseOperator from airflow.utils.decorators import apply_defaults from airflow.settings import Session from sqlalchemy import text class MutexLockOperator(BaseOperator): @apply_defaults def __init__(self, lock_name, **kwargs): super().__init__(**kwargs) self.lock_name = lock_name def execute(self, context): session = Session() try: # 尝试获取锁,0表示立即返回不等待 lock_result = session.execute( text("SELECT GET_LOCK(:lock_name, 0)"), {'lock_name': self.lock_name} ).scalar() if lock_result != 1: raise Exception(f"无法获取互斥锁,已有其他DAG在运行") finally: session.close() class MutexUnlockOperator(BaseOperator): @apply_defaults def __init__(self, lock_name, **kwargs): super().__init__(**kwargs) self.lock_name = lock_name def execute(self, context): session = Session() try: session.execute( text("SELECT RELEASE_LOCK(:lock_name)"), {'lock_name': self.lock_name} ) session.commit() finally: session.close()
然后在两个DAG中添加锁流程:
- DAG A:
MutexLockOperator→ 原有任务 →MutexUnlockOperator - DAG B:
MutexLockOperator→ 原有任务 →MutexUnlockOperator
先启动的DAG会获取锁,后启动的会因锁被占用而失败,你可以结合失败回调让DAG自动重新调度,直到锁被释放。
方案3:修改ExternalTaskSensor的等待逻辑
不要等待对方的最近一次执行,改为等待对方所有正在运行的实例完成,避免锁定到当前运行的冲突实例。
修改DAG A的传感器逻辑:
from airflow.sensors.external_task import ExternalTaskSensor from airflow.models import DagRun, settings from airflow.utils.state import State def get_running_dag_b_exec_dates(context): session = settings.Session() # 获取DAG B所有正在运行实例的执行时间 running_exec_dates = session.query(DagRun.execution_date)\ .filter(DagRun.dag_id == 'DAG_B', DagRun.state == State.RUNNING)\ .all() session.close() # 返回需要等待的执行时间列表,无运行实例则直接通过 return [date[0] for date in running_exec_dates] if running_exec_dates else None sensor = ExternalTaskSensor( task_id='wait_for_dag_b_running_instances', external_dag_id='DAG_B', external_task_id=None, dag=dag, check_existence=True, execution_date_fn=get_running_dag_b_exec_dates, allowed_states=["success", "failed"], mode='reschedule' )
DAG B做同样修改即可,这样传感器只会等待真正在跑的实例完成,不会陷入死锁。
推荐方案
如果DAG B运行频率更低(每小时),优先选方案1,实现简单且能有效避免死锁;如果对互斥严谨性要求高,方案2的全局锁更可靠。
内容的提问来源于stack exchange,提问作者stacked
相关产品推荐
相关产品推荐

