Airflow多DAG任务调度异常:任务卡在Scheduled状态无法并行
你的问题核心是两个DAG无法同时并行执行,本质是Airflow的资源调度限制导致dynamic_Task_Concurrency_two无法获取足够的任务槽。结合你的配置(parallelism=36、max_active_tasks_per_dag=12),全局总槽数足够支撑两个DAG各12个任务(24<36),所以问题大概率出在以下几个细节配置上:
可能的原因及对应解决步骤
1. 默认Pool槽数不足
Airflow默认的default_pool槽数为12,若两个DAG都使用默认Pool,dynamic_Task_Concurrency会占满全部12个槽,导致dynamic_Task_Concurrency_two无法启动任务。
解决方法:
- 调整默认Pool槽数:进入Airflow UI的
Admin -> Pools,找到default_pool,将槽数修改为24或更高(满足两个DAG的并发需求)。 - 为DAG分配独立Pool:创建两个自定义Pool,分别分配给两个DAG:
- 在UI中创建
dag1_pool(槽数12)和dag2_pool(槽数5,因为dynamic_Task_Concurrency_two只有5个任务)。 - 在DAG定义中指定Pool:
# 给dynamic_Task_Concurrency指定Pool dag1 = DAG( 'dynamic_Task_Concurrency', default_args=default_args, pool='dag1_pool', concurrency=12 ) # 给dynamic_Task_Concurrency_two指定Pool dag2 = DAG( 'dynamic_Task_Concurrency_two', default_args=default_args, pool='dag2_pool', concurrency=5 )
- 在UI中创建
2. dynamic_Task_Concurrency的活跃运行实例过多
若dynamic_Task_Concurrency的max_active_runs(单个DAG允许同时运行的实例数)设置过大,比如3,那么它的3个实例会占用3*12=36个全局槽,直接占满parallelism的上限,导致另一个DAG无槽可用。
解决方法:
- 在DAG定义中限制
max_active_runs为1:dag = DAG( 'dynamic_Task_Concurrency', default_args=default_args, max_active_runs=1, concurrency=12 ) - 或者修改全局配置
airflow.cfg中的max_active_runs_per_dag=1,然后重启Airflow服务。
3. dynamic_Task_Concurrency_two的DAG级并发限制
如果该DAG自身定义了concurrency参数且值过小(比如0或1),会覆盖全局的max_active_tasks_per_dag配置,导致任务无法并行启动。
解决方法:
在DAG定义中显式设置concurrency为至少5(或全局的12):
dag = DAG( 'dynamic_Task_Concurrency_two', default_args=default_args, concurrency=12 )
4. 未重启Airflow服务使配置生效
修改airflow.cfg后,必须重启Webserver、Scheduler和Worker服务,新配置才能生效。
解决方法:
执行重启命令(根据你的部署方式调整):
# 若使用systemd sudo systemctl restart airflow-webserver airflow-scheduler airflow-worker # 若使用docker-compose docker-compose restart webserver scheduler worker
5. 跨DAG依赖导致等待
检查dynamic_Task_Concurrency_two的任务是否包含ExternalTaskSensor等依赖dynamic_Task_Concurrency的组件,若有则会等待依赖任务完成才启动。
解决方法:
检查DAG任务定义,移除不必要的跨DAG依赖逻辑。
内容的提问来源于stack exchange,提问作者lifetolearn22

