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

如何基于输入参数对Airflow Dagrun进行优先级排序?

Airflow DagRun级优先级实现方案(基于level参数)

针对你需要根据level参数为DagRun设置调度优先级的需求,以下是两种可行的实现方案,适配不同Airflow版本:

方法一:任务级动态绑定DagRun参数(兼容所有Airflow版本)

利用Airflow的模板渲染功能,将DagRun传入的level参数动态映射为任务的priority_weight。由于你的任务是串行执行,只要第一个任务的优先级足够高,整个DagRun就会被优先调度。

代码示例

from airflow import DAG
from airflow.operators.python import PythonOperator
from datetime import datetime

def long_running_task(**context):
    # 替换为你的长时任务逻辑
    import time
    time.sleep(3600)

default_args = {
    'owner': 'airflow',
    'start_date': datetime(2023, 1, 1),
}

with DAG(
    dag_id='DagX',
    default_args=default_args,
    concurrency=16,
    schedule_interval=None,  # 手动触发
    catchup=False
) as dag:
    # 串行任务1:动态设置优先级权重
    task1 = PythonOperator(
        task_id='long_task_1',
        python_callable=long_running_task,
        priority_weight="{{ dag_run.conf.get('level', 1) * 100 }}",
        provide_context=True
    )

    # 串行任务2:复用相同的优先级权重
    task2 = PythonOperator(
        task_id='long_task_2',
        python_callable=long_running_task,
        priority_weight="{{ dag_run.conf.get('level', 1) * 100 }}",
        provide_context=True
    )

    task1 >> task2

逻辑说明

  • priority_weight通过模板语法从DagRun的conf中读取level值,乘以100放大权重差异(避免优先级数值过于接近导致调度逻辑失效)
  • 串行任务共享相同的优先级权重,确保整个DagRun的调度优先级一致
  • Airflow调度器会优先执行priority_weight更高的任务,当Dag并发数耗尽时,队列中的高level任务会被优先调度

方法二:直接设置DagRun的priority_weight(Airflow 2.2+)

Airflow 2.2及以上版本支持DagRun级别的priority_weight属性,调度器会直接根据该值排序待执行的DagRun,无需为每个任务单独配置。

方案1:触发时手动映射level到优先级

通过Airflow API触发Dag时,直接将level参数转换为priority_weight:

from airflow.api.client.local_client import Client

client = Client(None, None)
# 触发level=5的DagRun,设置最高优先级
client.trigger_dag(
    dag_id='DagX',
    conf={'level': 5},
    priority_weight=5 * 100  # level值放大100倍作为权重
)

方案2:自动绑定level与DagRun优先级

通过Dag的dagrun_create_callback回调函数,在DagRun创建时自动根据level参数设置优先级:

from airflow import DAG
from airflow.operators.python import PythonOperator
from airflow.models import DagRun
from datetime import datetime

def long_running_task(**context):
    # 替换为你的长时任务逻辑
    import time
    time.sleep(3600)

def set_dagrun_priority(context):
    """DagRun创建时自动设置优先级"""
    dag_run = context['dag_run']
    level = dag_run.conf.get('level', 1)
    dag_run.priority_weight = level * 100
    dag_run.session.commit()

default_args = {
    'owner': 'airflow',
    'start_date': datetime(2023, 1, 1),
}

with DAG(
    dag_id='DagX',
    default_args=default_args,
    concurrency=16,
    schedule_interval=None,
    catchup=False,
    dagrun_create_callback=set_dagrun_priority
) as dag:
    task1 = PythonOperator(
        task_id='long_task_1',
        python_callable=long_running_task
    )

    task2 = PythonOperator(
        task_id='long_task_2',
        python_callable=long_running_task
    )

    task1 >> task2

逻辑说明

  • dagrun_create_callback在DagRun被创建时触发,自动将level值映射为priority_weight
  • 调度器会优先处理priority_weight更高的DagRun,整个DagRun的所有任务都会获得优先级提升
  • 无需为每个任务配置优先级,实现真正的DagRun级优先级控制

注意事项

  • 权重放大倍数可根据实际需求调整,确保不同level之间的优先级差异足够明显
  • 若使用方法一,需确保所有串行任务的priority_weight保持一致,避免中间任务优先级过低导致DagRun停滞
  • Airflow调度器的max_tis_per_query配置可能影响队列排序效率,建议保持默认值或适当调大

内容的提问来源于stack exchange,提问作者Atur

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.11 18:15:02