You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

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处理该队列:

  1. 首先在Celery配置中添加任务路由规则,指定single_task走专属队列:
app.conf.task_routes = {
    # 替换为你项目中single_task的实际导入路径
    "your_app.tasks.single_task": {"queue": "single_serial_queue"}
}
  1. 启动两个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. 提交任务时走异步调用:
# 终端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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.10.02 10:48:04