如何限制Celery应用全局活动任务数量为2
配置失效原因
你之前尝试的三种concurrency=2配置全部无效,核心原因有两个:
concurrency是Celery Worker的启动参数,不是Celery应用初始化、任务装饰器的直接传参,写在Celery(__name__, concurrency=2)、@celery.task(concurrency=2)里的参数不会被Worker识别- 写在全局配置里时,正确的配置键名是
worker_concurrency,不是concurrency,之前的配置键名写错,Celery会直接忽略该配置
从你贴的运行日志看,12个任务几乎同时启动、耗时都在5-6s区间,说明Worker实际用了默认并发值(等于机器CPU核心数),完全没有做并发限制,才会导致高内存任务同时跑满占满资源。
可行配置方案
方案1:启动命令直接指定参数(最稳妥无兼容问题)
启动Worker时直接追加并发参数,不用改业务代码,适合临时调试场景:
# 替换成你自己的Celery应用模块路径 celery -A your_celery_app_path worker --concurrency=2 --loglevel=info
补充:如果你的任务是IO密集型,启动时用了gevent/eventlet协程池,需要额外加--prefetch-multiplier=1参数,禁止Worker提前预取超过并发数的任务到本地缓存,严格保证同时运行的任务数不超过2。
方案2:写入应用全局配置(固化配置不用每次启动输参数)
修正配置键名,把并发限制写到Celery配置里,启动Worker时不需要额外加参数:
from celery import Celery celery = Celery(__name__) # 注意配置键名必须是worker_concurrency,不能简写为concurrency celery.conf.update( worker_concurrency=2, # 协程池/严格并发限制场景加下面这行,进程池场景可省略 worker_prefetch_multiplier=1, ) @celery.task(name="task_runner", bind=True) def generic_task(self, function_type): task_id = self.request.id if function_type == 1: high_ram_usage_1() # 内存占用极高的函数 if function_type == 2: high_ram_usage_2() # 另一个高内存占用函数 # 其余业务逻辑
注意:
@celery.task()装饰器本身不支持concurrency参数,全局并发限制是Worker级别的配置,无法通过单任务装饰器设置。
配置生效验证
启动Worker后看启动日志的开头输出,找到并发数打印行,显示如下内容就代表配置生效:
concurrency: 2 (prefork)
配置生效后,新投递的任务最多同时只有2个处于运行状态,其余任务会在Broker队列中排队,等待运行中的任务执行完成释放资源后才会被调度执行,不会再出现十几个高内存任务同时运行占满内存的问题。
内容的提问来源于stack exchange,提问作者vahvero
相关产品推荐
相关产品推荐

