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

如何让紧急DagRun跳过Celery队列优先执行?

实现Airflow DagRun级别的紧急优先执行(Celery执行器)

Airflow本身没有直接为整个DagRun设置全局优先级的原生功能,但针对Celery执行器,可以通过以下几种变通方案实现“紧急DagRun插队”的需求:

1. 为紧急DagRun分配专属高优先级Celery队列

这是最可靠的方案,通过隔离队列实现优先级调度:

  • 第一步:配置高优先级队列
    在Airflow的airflow.cfg中添加新队列:
    [celery]
    task_queue_names = default,high_priority
    
    同时在Celery的配置文件(如celeryconfig.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指定高优先级队列:
    airflow dags trigger -c '{"queue": "high_priority"}' your_dag_id
    
    这样该DagRun的所有任务都会进入high_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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.11 14:32:20