You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

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图

DAG图

问题根因

  1. TaskGroup的default_args不传递给子任务:你在critical_group的default_args中设置的priority_weight=2,仅作用于TaskGroup本身(虚拟任务节点),不会自动传递给其下的子任务和子TaskGroup,导致实际子任务优先级未生效。
  2. 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,提问作者형한결

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.07.16 22:21:59