Airflow技术问题:如何在多个Worker上运行任务
如何在多个Airflow Worker上运行任务
嘿,你已经通过Celery Executor搭建好Airflow,并且在DAG里给任务指定了不同队列——这已经走对了关键一步!要让任务在多个Worker上运行,核心就是让不同Worker监听指定队列,配合你DAG里的队列配置就能实现,具体步骤如下:
1. 启动对应队列的Celery Worker
Celery Worker是通过监听特定队列来获取任务的,你需要为每个目标队列(或组合队列)启动独立的Worker进程:
- 启动监听
test_queue的Worker:airflow celery worker -q test_queue - 启动监听
local_queue的Worker:airflow celery worker -q local_queue - 如果想让一个Worker同时处理多个队列的任务,用逗号分隔队列名称即可:
airflow celery worker -q test_queue,local_queue
2. 确认DAG中的任务队列配置
你当前的DAG代码已经做对了关键配置:
t1绑定了test_queue,会被自动分发到监听该队列的Worker执行t2和t3绑定了local_queue,会被分发到监听该队列的Worker执行
这个设定完全没问题,Airflow会根据每个任务的queue参数,自动把任务投递到对应的Celery队列中。
3. 可选:配置全局默认队列
如果有部分任务没有显式指定queue参数,Airflow会使用airflow.cfg里的default_queue配置(默认值为default)。你可以启动监听默认队列的Worker来处理这类任务:
airflow celery worker -q default
4. 验证任务执行效果
启动所有需要的Worker后,触发你的dummy_for_testing DAG,就能通过以下方式验证任务是否在多Worker上运行:
- 查看Worker的运行日志,能看到任务被分配到对应Worker的记录
- 检查
/tmp/dag_output.log文件,如果你的Worker部署在不同机器上,hostname的输出会不一样,这就说明任务确实在不同Worker上执行了
另外补充一点:如果你的Worker部署在不同机器上,只要所有Worker都能连接到同一个Celery消息中间件(比如Redis或RabbitMQ),并且Airflow的核心配置(数据库连接、airflow.cfg参数)保持一致,就能正常协同工作,接收并执行任务。
内容的提问来源于stack exchange,提问作者Mukul Jain
相关产品推荐
相关产品推荐

