Airflow 2.7.1:两个同调度DAG无法同时运行的问题求助
Airflow多DAG无法并行运行的配置方案
问题背景
有两个每5分钟调度一次的DAG(dagA和dagB),当前无法同时运行——同一时间只有一个DAG的单个任务在执行,另一个DAG的任务甚至不会入队,必须等前者任务完成才启动。已确认全局parallelism配置为32,两个DAG的参数、任务结构一致,相关信息如下:
调度与执行表现
调度规则:两个DAG均为 */5 * * * * dagA 包含任务:task1、task2、task3 dagB 包含任务:task1、task2 执行顺序示例: dagA.task1 (运行中) → dagB无任务入队 dagA.task1 (执行完成) → dagB.task2入队、dagB.task1开始运行
DAG代码示例
@dag( dag_id='dagA', schedule_interval='*/5 * * * *', # 每5分钟调度一次 start_date=datetime(2023, 9, 1, tz="UTC"), catchup=False, max_active_runs=1, default_args=default_args, ) # --- 分割线 --- @dag( dag_id='dagB', schedule_interval='*/5 * * * *', # 每5分钟调度一次 start_date=datetime(2023, 9, 1, tz="UTC"), catchup=False, max_active_runs=1, default_args=default_args, ) def process_gdelt_etl(): @task( retries=1, retry_delay=timedelta(minutes=4) ) # 两个DAG的任务结构一致 def test(): pass
解决配置调整
1. 调整单个DAG的最大并行任务数
检查全局配置max_active_tasks_per_dag(Airflow 2.x也可用DAG级参数max_active_tasks),这个参数限制单个DAG同时运行的任务数量,默认可能为1,导致单个DAG同一时间只能跑一个任务,还会占用执行资源阻塞其他DAG。
- 全局配置修改(修改
airflow.cfg):
max_active_tasks_per_dag = 3 # 可根据每个DAG的任务数量设置,比如dagA有3个任务就设为3
- DAG单独配置(在DAG定义中添加):
@dag( dag_id='dagA', # 其他现有参数... max_active_tasks=3, # 覆盖全局,指定该DAG允许同时运行的任务数 )
2. 确认dag_concurrency配置
这个参数是单个DAG的并发任务上限(部分Airflow版本中是max_active_tasks的别名),确保全局或DAG级别的值大于1,避免单个DAG的任务串行执行占用所有资源。
3. 切换支持并行的执行器
如果你用的是默认的SequentialExecutor,它是单线程执行器,同一时间只能跑一个任务,完全不支持并行。必须切换到以下执行器:
LocalExecutor:适合单机部署,支持多进程并行执行CeleryExecutor:适合分布式部署,支持大规模任务并行
修改airflow.cfg中的执行器配置:
executor = LocalExecutor
修改后必须重启Airflow的webserver和scheduler服务才能生效。
4. 检查任务队列配置
如果你的default_args中指定了queue参数,要确保对应队列的worker数量足够。如果没有特殊需求,保持默认队列即可,避免因队列资源不足导致任务无法并行。
5. 确认max_active_runs配置
你当前已经设置max_active_runs=1,这个参数控制单个DAG同时运行的实例数,当前配置没问题,无需调整,只要不设为0或负数即可。
验证方法
- 修改配置后重启Airflow的webserver和scheduler
- 手动触发两个DAG,观察任务是否同时进入运行状态
- 通过Airflow UI的Graph View或Task Instance List,确认跨DAG的多个任务同时处于
running状态
内容的提问来源于stack exchange,提问作者nickanor
相关产品推荐
相关产品推荐

