运行中的Celery5应用如何修改并发数、管理任务队列?
FastAPI集成Celery 5动态调整Worker并发数方案
背景
在FastAPI应用中使用Celery 5管理异步任务,初始通过如下命令启动了单Worker、并发数为8的Celery服务:celery worker --app=app.worker.celery --concurrency=8 --loglevel=info --logfile=logs/celery.log
需求为直接从FastAPI应用中动态调整Celery的并发数,需要验证方案可行性并找到最优实现。
初期踩坑方案
最初未找到直接修改运行中Worker并发数的方法,尝试通过代码直接启动新Worker,实现代码如下:
from celery import current_app as app cmd = ["--app=app.worker.celery", "--concurrency=8", "--loglevel=info", "--logfile=logs/celery.log", "--without-gossip" , "--detach", "-E"] app.worker_main(cmd)
该方案存在缺陷:即使传入--detach后台运行参数,调用时仍会阻塞FastAPI的请求,无法在生产环境使用。
最终最优实现
参考Flower 1.0.1的实现逻辑,使用Celery官方提供的远程控制API即可直接动态调整运行中Worker的并发数,无需重启Worker,也不会阻塞请求,实现代码如下:
from celery import current_app as app response = app.control.pool_grow( n=4, reply=True, destination=[worker_name])
API补充说明
pool_grow方法用于给指定Worker扩容并发数,参数n为新增的并发进程/协程/线程数- 如需缩容并发数,可调用
pool_shrink方法:app.control.pool_shrink(n=2, reply=True, destination=[worker_name]),参数n为要减少的并发数 destination参数为要调整的Worker名称列表,如省略该参数则会调整所有在线Worker的并发数reply=True会返回对应Worker的执行结果,可用于校验调整操作是否成功
内容的提问来源于stack exchange,提问作者Rafael
相关产品推荐
相关产品推荐

