如何让紧急DagRun跳过Celery队列优先执行?
实现Airflow DagRun级别的紧急优先执行(Celery执行器)
Airflow本身没有直接为整个DagRun设置全局优先级的原生功能,但针对Celery执行器,可以通过以下几种变通方案实现“紧急DagRun插队”的需求:
1. 为紧急DagRun分配专属高优先级Celery队列
这是最可靠的方案,通过隔离队列实现优先级调度:
- 第一步:配置高优先级队列
在Airflow的airflow.cfg中添加新队列:
同时在Celery的配置文件(如[celery] task_queue_names = default,high_priorityceleryconfig.py)中设置队列优先级权重,确保high_priority队列被优先处理:task_routes = { '*': {'queue': 'default'}, } # 配置队列优先级,数值越大优先级越高 queue_priorities = { 'high_priority': 10, 'default': 1, } - 第二步:让DAG支持动态指定队列
修改DAG代码,允许通过DagRun的配置动态给所有任务分配队列:from airflow import DAG from airflow.operators.python import PythonOperator from datetime import datetime def get_task_queue(**context): # 从DagRun的conf中读取队列,默认用default return context['dag_run'].conf.get('queue', 'default') with DAG( dag_id='your_dag_id', start_date=datetime(2023, 1, 1), schedule=None, ) as dag: task1 = PythonOperator( task_id='task1', python_callable=lambda: print("Running task"), queue=get_task_queue, provide_context=True ) task2 = PythonOperator( task_id='task2', python_callable=lambda: print("Running task 2"), queue=get_task_queue, provide_context=True ) task1 >> task2 - 第三步:触发紧急DagRun
使用CLI或API触发时,通过conf指定高优先级队列:
这样该DagRun的所有任务都会进入airflow dags trigger -c '{"queue": "high_priority"}' your_dag_idhigh_priority队列,Celery Worker会优先处理这个队列的任务,实现跳过默认队列等待任务的效果。
2. 动态提升紧急DagRun任务的优先级权重
如果不想新增队列,可以通过修改任务实例的priority_weight让Celery优先调度:
- 触发紧急DagRun后,批量更新任务优先级
编写脚本或在Airflow中添加一个辅助任务,将目标DagRun的所有任务实例优先级设为最高:from airflow.models import DagRun, TaskInstance from airflow.utils.session import create_session def elevate_dagrun_priority(dag_id, run_id, priority=1000): with create_session() as session: dag_run = session.query(DagRun).filter( DagRun.dag_id == dag_id, DagRun.run_id == run_id ).first() if not dag_run: print("DagRun not found") return # 获取该DagRun下所有未执行的任务实例 tis = session.query(TaskInstance).filter( TaskInstance.dag_run_id == dag_run.id, TaskInstance.state.in_(['queued', 'scheduled']) ).all() for ti in tis: ti.priority_weight = priority # 设为远高于其他任务的值(默认是1) session.add(ti) session.commit() print(f"Updated {len(tis)} tasks to priority {priority}") # 调用示例 elevate_dagrun_priority("your_dag_id", "emergency_run_xxxx") - 配套Celery配置
确保Celery启用优先级支持,在celeryconfig.py中添加:task_priority = 10 # 允许的最大优先级值 worker_prefetch_multiplier = 1 # 禁用预取,避免Worker提前拿到低优先级任务 task_acks_late = True
注意事项
- 上述方案均针对Celery执行器,不适用于其他执行器(如KubernetesExecutor)
- 调整数据库或Celery任务状态时,务必在低峰期操作,避免影响正常调度
- 如果需要频繁使用紧急DagRun,建议优先采用“专属高优先级队列”的方案,更易于维护和监控
内容的提问来源于stack exchange,提问作者OhadBasan
相关产品推荐
相关产品推荐

