如何实现Celery资源密集型任务的串行执行?(已试队列与chain无效)
要让资源密集型任务串行执行,有两种实用方案,根据你的场景选择:
方案1:给heavy队列分配单进程worker
这是最简单的方案,无需修改任务代码,只需调整worker的启动参数:
操作步骤:
- 确保只启动一个worker进程来监听
heavy队列,并且设置该worker的并发数为1(同一时间只能处理一个任务)。启动命令如下:
celery -A your_app worker -Q heavy --concurrency=1 --loglevel=info
(把your_app替换成你的Celery应用实例所在的模块名)
原理:
这个worker只会用一个进程处理heavy队列的任务,所有提交到该队列的任务都会按顺序排队,前一个任务执行完毕后才会启动下一个。即使你提交任务的时间分散,也能保证串行执行。
方案2:在任务内部加分布式锁
如果不想限制worker数量(比如worker还需要处理其他队列的并行任务),可以在任务内部加锁,确保同一时间只有一个heavy任务实例在运行。这里以Redis锁为例:
操作步骤:
- 先安装Redis依赖:
pip install redis
- 修改任务代码,添加锁逻辑:
import redis from celery import shared_task from celery.exceptions import Retry # 初始化Redis客户端(根据你的Redis配置调整参数) redis_client = redis.Redis(host='localhost', port=6379, db=0, decode_responses=True) # 锁的唯一标识 LOCK_KEY = "heavy_task_running_lock" @shared_task(name="tasks.heavy", bind=True, max_retries=None) def some_task(self, arg): # 尝试获取锁,不阻塞,立即返回结果 # timeout设置为任务最长可能执行的时间(比如3600秒,根据实际情况调整) lock = redis_client.lock(LOCK_KEY, timeout=3600) lock_acquired = lock.acquire(blocking=False) if not lock_acquired: # 获取锁失败,10秒后重试(可调整重试间隔) raise self.retry(countdown=10) try: # 这里写你的资源密集型任务逻辑 print(f"执行任务,参数:{arg}") # ... 任务代码 ... finally: # 无论任务成功还是失败,都释放锁 lock.release()
原理:
当一个任务获取到锁后,其他任务会因为无法获取锁而触发重试,直到锁被释放(前一个任务执行完毕或超时),这样就能保证同一时间只有一个heavy任务在运行。这种方案不影响worker处理其他队列的任务,资源利用率更高。
内容的提问来源于stack exchange,提问作者aplusk
相关产品推荐
相关产品推荐

