Python环境下Celery任务队列先入先出(FIFO)不生效问题求助
问题原因
- 你当前直接调用
single_task(参数)属于同步执行函数逻辑,任务没有提交到Celery的分布式队列,两个终端的调用是在各自进程独立运行,自然不会互相等待。异步提交任务需要使用single_task.delay(参数)方法。 - 默认Celery Worker的并发数大于1,会同时拉取多个任务并行执行,哪怕任务都进了队列也不会按顺序串行执行。
解决方案
场景1:全队列都需要严格FIFO串行
直接调整Worker启动参数,将并发数设为1即可,启动命令如下:celery -A 你的Django项目名 worker -l info -c 1
所有提交到队列的任务会按先入先出的顺序逐个执行。
场景2:仅single_task需要串行,其他任务保持并发
给single_task分配独立专用队列,单独启动一个并发为1的Worker处理该队列:
- 首先在Celery配置中添加任务路由规则,指定single_task走专属队列:
app.conf.task_routes = { # 替换为你项目中single_task的实际导入路径 "your_app.tasks.single_task": {"queue": "single_serial_queue"} }
- 启动两个Worker,一个处理普通任务(保持正常并发),一个专门处理串行任务队列(并发设为1):
# 处理普通任务的Worker,按需设置并发数 celery -A 你的Django项目名 worker -l info -c 4 -Q celery # 处理串行任务的Worker,并发数固定为1 celery -A 你的Django项目名 worker -l info -c 1 -Q single_serial_queue
- 提交任务时走异步调用:
# 终端1提交 single_task.delay(True) # 终端2提交 single_task.delay(False)
场景3:需要更严格的串行保障(避免误启动多个串行Worker)
可以引入分布式锁保证同一时间只有一个single_task在执行,即便不小心启动了多个处理串行队列的Worker也不会出现并行问题:
修改任务代码如下:
import redis import time from celery.exceptions import Retry # 初始化Redis客户端,和你Celery用的Redis实例保持一致即可 redis_client = redis.Redis(host="127.0.0.1", port=6379, db=0, decode_responses=True) SINGLE_TASK_LOCK_KEY = "single_task_running_lock" @app.task(bind=True, max_retries=None) def single_task(self, delay): """Run task.""" # 尝试获取分布式锁,锁过期时间设置为大于任务最大可能执行时长,避免死锁 lock_acquired = redis_client.set(SINGLE_TASK_LOCK_KEY, "locked", nx=True, ex=60) if not lock_acquired: # 未获取到锁则1秒后重试 raise self.retry(countdown=1) try: if delay: time.sleep(10) print("ran ...") return True finally: # 任务执行完成后释放锁 redis_client.delete(SINGLE_TASK_LOCK_KEY)
内容的提问来源于stack exchange,提问作者soubhagya
相关产品推荐
相关产品推荐

