AWS托管Airflow中max_active_runs无法限制DAG并发运行问题
问题根因
你遇到的配置不生效问题来自两个层面:
- 参数作用混淆:
concurrency参数控制的是单DAG下所有运行实例允许同时执行的任务总数,并非DAG运行实例的并发数,控制DAG运行并发的参数确实是max_active_runs。 - 版本BUG:Airflow 2.0.2存在已知缺陷,手动触发的DAG运行不会校验
max_active_runs限制,该问题在2.2及以上版本才被官方修复,这是你配置不生效的核心原因。
你提供的DAG运行截图如下:
可行解决方案
方案1:升级Airflow版本
将AWS托管Airflow(MWAA)环境升级到2.2及以上版本,升级后你当前的max_active_runs=1配置即可生效,同时触发的第二个DAG运行会自动进入排队状态。
方案2:新增前置并发检查任务(无需升级版本)
如果暂时无法升级版本,可以在DAG的第一个节点添加运行状态校验逻辑,检测到已有运行中实例时直接跳过当前运行:
from airflow.models import DagRun from airflow.operators.python import PythonOperator from airflow.exceptions import AirflowSkipException def check_existing_runs(**context): current_dag_id = context["dag"].dag_id current_run_id = context["run_id"] # 查询当前DAG所有运行中实例 running_runs = DagRun.find(dag_id=current_dag_id, state="running") # 排除当前运行实例本身,如果存在其他运行中实例则跳过 if any(run.run_id != current_run_id for run in running_runs): raise AirflowSkipException("存在正在运行的同DAG实例,本次运行终止") # 注册为DAG第一个任务 check_task = PythonOperator( task_id="check_concurrency", python_callable=check_existing_runs, provide_context=True, dag=dag ) # 后续任务全部依赖check_task # check_task >> 你的其他任务
方案3:绑定专属资源池
给该DAG单独创建大小为1的资源池,强制限制所有任务的并发能力:
- 进入Airflow UI → Admin → Pools,新建名为
single_run_pool的资源池,Slots设置为1 - 在DAG的
default_args中添加配置pool: "single_run_pool"即可生效
内容的提问来源于stack exchange,提问作者Birat Bade Shrestha
相关产品推荐
相关产品推荐


