如何避免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
相关产品推荐
相关产品推荐

