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

如何避免Celery队列中重复添加相同任务?

解决FastAPI+Celery重复任务问题的实现方案

核心思路

借助Redis作为共享状态存储,跟踪每个请求字符串对应的任务状态,区分三种处理场景:

  • 缓存已有计算结果:直接返回状态码200和结果
  • 任务正在队列等待或执行中:返回状态码202,不重复提交任务
  • 无缓存且无正在处理的任务:提交Celery任务并标记状态,返回202

具体实现步骤

1. 封装Redis状态操作工具函数

统一处理任务状态的标记、查询与清理,保证操作原子性:

import redis

# 复用项目中已有的Redis连接实例
redis_client = redis.Redis(host="localhost", port=6379, db=0)
PROCESSING_KEY_PREFIX = "processing:"

def is_task_running(target_str: str) -> bool:
    """检查指定字符串对应的任务是否正在处理"""
    return redis_client.exists(PROCESSING_KEY_PREFIX + target_str)

def mark_task_running(target_str: str, task_id: str, expire_sec: int = 15):
    """标记任务为处理中,设置过期时间避免异常残留"""
    # SETNX保证只有当key不存在时才设置,防止并发请求重复标记
    redis_client.setnx(PROCESSING_KEY_PREFIX + target_str, task_id)
    redis_client.expire(PROCESSING_KEY_PREFIX + target_str, expire_sec)

def clear_task_running_mark(target_str: str):
    """移除任务的处理中标记"""
    redis_client.delete(PROCESSING_KEY_PREFIX + target_str)

2. 修改FastAPI路由逻辑

在原有逻辑中加入任务状态检查:

from fastapi import FastAPI, status
from your_celery_module import app as celery_app

app = FastAPI()

@app.get("/query/{target_str}")
async def query_result(target_str: str):
    # 第一步:检查缓存是否已有结果
    cache_key = f"result:{target_str}"
    cached_data = redis_client.hgetall(cache_key)
    if cached_data:
        return {"data": cached_data, "status": "completed"}, status.HTTP_200_OK
    
    # 第二步:检查任务是否正在处理
    if is_task_running(target_str):
        return {"status": "processing"}, status.HTTP_202_ACCEPTED
    
    # 第三步:提交任务并标记状态
    task = celery_app.send_task("compute_result", args=[target_str])
    # 过期时间设为任务最长耗时+5秒,避免标记残留
    mark_task_running(target_str, task.id, expire_sec=15)
    return {"task_id": task.id, "status": "pending"}, status.HTTP_202_ACCEPTED

3. 调整Celery任务逻辑

任务执行完成后自动清理状态标记,并写入缓存:

from celery import Celery

app = Celery("tasks", broker="redis://localhost:6379/0", backend="redis://localhost:6379/0")

@app.task(bind=True)
def compute_result(self, target_str: str):
    try:
        # 替换为实际的耗时计算逻辑
        computed_result = {
            "input": target_str,
            "value": f"processed_{target_str}",
            "completed_at": "2024-xx-xx xx:xx:xx"
        }
        
        # 将结果写入Redis缓存
        cache_key = f"result:{target_str}"
        redis_client.hset(cache_key, mapping=computed_result)
        redis_client.expire(cache_key, 3600)  # 缓存过期时间按需调整
        
        return computed_result
    finally:
        # 无论任务成功失败,都清理处理中标记
        clear_task_running_mark(target_str)

关键细节说明

  • 用SETNX原子操作避免并发请求同时提交任务,解决竞态问题
  • 给处理中标记设置过期时间,防止worker崩溃等异常场景下标记永久残留
  • 若需要更精细化的任务状态追踪,也可结合Celery的AsyncResult查询任务状态,但Redis直接标记的性能更优,适合高并发场景

内容的提问来源于stack exchange,提问作者Robert Alexander

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.07 00:05:59