如何基于输入参数对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
相关产品推荐
相关产品推荐

