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

Airflow:如何将流程任务分发至多个Worker节点执行?

解决Airflow任务仅单Worker运行的问题

首先修正你代码中的拼写错误:python_callabale 应为 python_callable,这个错误会导致任务无法正常执行,先修复后再进行以下配置调整。

要将任务分发到多个Worker节点,可从以下几个方面配置:

  • 确认使用CeleryExecutor
    只有CeleryExecutor支持多Worker集群部署,如果你当前用的是LocalExecutor或SequentialExecutor,需要切换到CeleryExecutor:

    1. 在airflow.cfg中设置executor = CeleryExecutor
    2. 配置消息队列(Redis或RabbitMQ),填写broker_url和result_backend
    3. 在所有Worker节点启动Celery Worker服务
  • 调整任务调度策略
    Airflow默认的任务分配策略可能导致任务集中到单个Worker,可修改配置切换为更均衡的策略:
    在airflow.cfg中设置task_execution_strategy = round_robin(轮询分配)或least_loaded(分配到负载最低的Worker),调度器会按该策略将任务分发到不同Worker。

  • 使用任务队列分流
    给任务指定不同的队列,让不同Worker监听对应队列,实现任务分流:

    1. 在DAG代码中给PythonOperator添加queue参数,循环分配不同队列:
      queues = ['worker_queue_1', 'worker_queue_2', 'worker_queue_3']
      for idx, value in enumerate(values_list):
          execute_func_1 = PythonOperator(
              task_id=f"{str(value)}_pyth_op_1",
              python_callable=function_1,
              op_kwargs={'value': value},
              queue=queues[idx % 3]
          )
          execute_func_2 = PythonOperator(
              task_id=f"{str(value)}_pyth_op_2",
              python_callable=function_2,
              op_kwargs={'value': value},
              queue=queues[idx % 3]
          )
          execute_func_1 >> execute_func_2
      
    2. 启动Worker时指定监听的队列,比如:
      • Worker1:airflow celery worker -q worker_queue_1
      • Worker2:airflow celery worker -q worker_queue_2
      • Worker3:airflow celery worker -q worker_queue_3
  • 调整并发参数
    检查并调整以下配置确保任务能并行分发:

    • 在airflow.cfg中设置worker_concurrency(单个Worker的并发任务数),避免单个Worker承接过多任务
    • 调整DAG的max_active_tasks参数,允许更多任务同时处于调度状态

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.28 19:42:37