Airflow 2.6.1中DAG未按priority_weight与weight_rule预期调度求助
Airflow任务优先级未按预期生效问题排查
环境信息
- Airflow版本:2.6.1
- Python版本:3.7
- 执行器:Celery Executor(1个worker)
- 部署方式:Docker Compose
问题描述
已将DAG的concurrency设为1,weight_rule配置为'absolute',按预期任务应仅依据priority_weight执行优先级排序。但实际运行时,设置了priority_weight=2的critical_group并未优先执行所有任务,与预期不符。
相关代码
# 设置weight_rule为absolute,concurrency=1 default_args = {"weight_rule": "absolute"} with DAG( dag_id="priority_test", concurrency=1, start_date=pendulum.datetime(2023, 7, 10, tz="Asia/Seoul"), default_args=default_args, ) as dag: extract_op = BashOperator(task_id="extract_op", bash_command="sleep 5") # critical_group -> priority_weight=2 with TaskGroup( group_id="critical_group", tooltip="critical_grouping", default_args={"retries": 3, "priority_weight": 2}, ) as critical_group: critical_extract = BashOperator( task_id="critical_extract", bash_command="sleep 5" ) with TaskGroup( group_id="critical_abc", tooltip="critical_abc_grouping" ) as critical_section_1: cr_extract_a = BashOperator(task_id="cr_extract_a", bash_command="sleep 5") cr_extract_b = BashOperator(task_id="cr_extract_b", bash_command="sleep 5") cr_extract_c = BashOperator(task_id="cr_extract_c", bash_command="sleep 5") cr_extract_a >> cr_extract_b >> cr_extract_c with TaskGroup( group_id="critical_123", tooltip="critical_123_grouping" ) as critical_section_2: cr_extract_1 = BashOperator(task_id="cr_extract_1", bash_command="sleep 5") cr_extract_2 = BashOperator(task_id="cr_extract_2", bash_command="sleep 5") cr_extract_3 = BashOperator(task_id="cr_extract_3", bash_command="sleep 5") cr_extract_1 >> cr_extract_2 >> cr_extract_3 critical_extract >> [critical_section_1, critical_section_2] # non_critical_group -> priority_weight=1 with TaskGroup( group_id="non_critical_group", tooltip="non_critical_grouping", default_args={"retries": 3, "priority_weight": 1}, ) as non_critical_group: non_critical_extract = BashOperator( task_id="non_critical_extract", bash_command="sleep 5", ) with TaskGroup( group_id="non_critical_abc", tooltip="non_critical_abc_grouping" ) as non_critical_section_1: non_cr_extract_a = BashOperator( task_id="non_cr_extract_a", bash_command="sleep 5" ) non_cr_extract_b = BashOperator( task_id="non_cr_extract_b", bash_command="sleep 5" ) non_cr_extract_c = BashOperator( task_id="non_cr_extract_c", bash_command="sleep 5" ) non_cr_extract_a >> non_cr_extract_b >> non_cr_extract_c with TaskGroup( group_id="non_critical_123", tooltip="non_critical_123_grouping" ) as non_critical_section_2: non_cr_extract_1 = BashOperator( task_id="non_cr_extract_1", bash_command="sleep 5" ) non_cr_extract_2 = BashOperator( task_id="non_cr_extract_2", bash_command="sleep 5" ) non_cr_extract_3 = BashOperator( task_id="non_cr_extract_3", bash_command="sleep 5" ) non_cr_extract_1 >> non_cr_extract_2 >> non_cr_extract_3 non_critical_extract >> [non_critical_section_1, non_critical_section_2] extract_op >> [critical_group, non_critical_group]
DAG图

问题根因
- TaskGroup的default_args不传递给子任务:你在
critical_group的default_args中设置的priority_weight=2,仅作用于TaskGroup本身(虚拟任务节点),不会自动传递给其下的子任务和子TaskGroup,导致实际子任务优先级未生效。 - Celery Executor配置缺失:若Celery未开启优先级队列支持,即使任务设置了优先级,调度器也不会按权重排序,但此问题核心是优先级未正确传递到子任务。
修复方案
方案1:显式给子TaskGroup设置优先级
要让critical_group下所有任务拥有高优先级,需在每个子TaskGroup的default_args中显式设置priority_weight=2,子任务会继承该值:
# 修改critical_abc子TaskGroup with TaskGroup( group_id="critical_abc", tooltip="critical_abc_grouping", default_args={"priority_weight": 2} ) as critical_section_1: cr_extract_a = BashOperator(task_id="cr_extract_a", bash_command="sleep 5") cr_extract_b = BashOperator(task_id="cr_extract_b", bash_command="sleep 5") cr_extract_c = BashOperator(task_id="cr_extract_c", bash_command="sleep 5") cr_extract_a >> cr_extract_b >> cr_extract_c # 修改critical_123子TaskGroup with TaskGroup( group_id="critical_123", tooltip="critical_123_grouping", default_args={"priority_weight": 2} ) as critical_section_2: cr_extract_1 = BashOperator(task_id="cr_extract_1", bash_command="sleep 5") cr_extract_2 = BashOperator(task_id="cr_extract_2", bash_command="sleep 5") cr_extract_3 = BashOperator(task_id="cr_extract_3", bash_command="sleep 5") cr_extract_1 >> cr_extract_2 >> cr_extract_3
方案2:利用继承简化配置
在DAG的default_args中设置默认优先级为1,再在critical_group的default_args中覆盖为2,子TaskGroup会自动继承父级优先级:
default_args = {"weight_rule": "absolute", "priority_weight": 1} with DAG( dag_id="priority_test", concurrency=1, start_date=pendulum.datetime(2023, 7, 10, tz="Asia/Seoul"), default_args=default_args, ) as dag: with TaskGroup( group_id="critical_group", tooltip="critical_grouping", default_args={"retries": 3, "priority_weight": 2}, ) as critical_group: critical_extract = BashOperator(task_id="critical_extract", bash_command="sleep 5") # 子TaskGroup自动继承priority_weight=2 with TaskGroup( group_id="critical_abc", tooltip="critical_abc_grouping" ) as critical_section_1: cr_extract_a = BashOperator(task_id="cr_extract_a", bash_command="sleep 5") # 其他子任务无需重复设置
方案3:开启Celery优先级支持
打开Airflow的airflow.cfg,在celery部分添加/修改配置:
[celery] task_queue_max_priority = 10 # 数值需大于设置的最大priority_weight
重启Celery Worker和Airflow调度器,确保配置生效,Celery才会按任务权重排序执行。
验证方式
修改后触发DAG运行,在Airflow UI的任务实例页面,查看每个任务的Priority Weight列,确认critical_group下任务优先级为2,同时观察执行顺序是否符合高优先级任务先完成的预期。
内容的提问来源于stack exchange,提问作者형한결
相关产品推荐
相关产品推荐

