Airflow:如何将流程任务分发至多个Worker节点执行?
解决Airflow任务仅单Worker运行的问题
首先修正你代码中的拼写错误:python_callabale 应为 python_callable,这个错误会导致任务无法正常执行,先修复后再进行以下配置调整。
要将任务分发到多个Worker节点,可从以下几个方面配置:
确认使用CeleryExecutor
只有CeleryExecutor支持多Worker集群部署,如果你当前用的是LocalExecutor或SequentialExecutor,需要切换到CeleryExecutor:- 在
airflow.cfg中设置executor = CeleryExecutor - 配置消息队列(Redis或RabbitMQ),填写
broker_url和result_backend - 在所有Worker节点启动Celery Worker服务
- 在
调整任务调度策略
Airflow默认的任务分配策略可能导致任务集中到单个Worker,可修改配置切换为更均衡的策略:
在airflow.cfg中设置task_execution_strategy = round_robin(轮询分配)或least_loaded(分配到负载最低的Worker),调度器会按该策略将任务分发到不同Worker。使用任务队列分流
给任务指定不同的队列,让不同Worker监听对应队列,实现任务分流:- 在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 - 启动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
- Worker1:
- 在DAG代码中给PythonOperator添加
调整并发参数
检查并调整以下配置确保任务能并行分发:- 在
airflow.cfg中设置worker_concurrency(单个Worker的并发任务数),避免单个Worker承接过多任务 - 调整DAG的
max_active_tasks参数,允许更多任务同时处于调度状态
- 在
内容的提问来源于stack exchange,提问作者drake10k
相关产品推荐
相关产品推荐

