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

如何避免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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.25 05:15:40