Airflow 2.2.0跨不同pool设置priority_weight无法优先调度单DAG run问题问询
根因说明
- Airflow的
priority_weight优先级是按Pool单独维护调度队列的,不同池内的任务优先级不会跨池比较,所以你给任务分配不同池之后,原先的优先级配置自然失效 - 你当前代码中所有任务的
priority_weight都用了默认值1,weight_rule仅用来计算上游累计权重,基础权重一致的前提下也无法体现优先级差异
最优解决方案
你的需求是同DAG的多次运行串行、前一次所有任务执行完成后再启动下一次,直接配置DAG级活跃运行数限制即可,完全不受任务分属不同池的影响:
在DAG初始化后添加配置即可:
dag = generate_dag('test_prio') # 新增这一行即可 dag.max_active_runs = 1
该参数会限制同一个DAG最多同时存在1个活跃DAG Run,只有当前运行的所有任务全部结束(无论成功/失败),下一个被触发的DAG Run才会开始调度执行。
特殊场景适配方案
如果业务需要同时保留多个活跃DAG Run,仅要求更早触发的Run优先执行完,可按以下规则调整:
- 给每个任务的
priority_weight赋值为对应DAG Run执行时间的时间戳整数,越早触发的Run权重越高,保证旧Run的所有任务不管在哪个池,优先级都高于新Run的任务 - 所有关联的Pool需要预留足够的slots,避免单池占满导致高优先级任务无法调度
内容的提问来源于stack exchange,提问作者Atef Arfaoui
相关产品推荐
相关产品推荐

