Celery同名同参数任务需等待前置任务完成的实现问询
实现同一用户的Celery任务串行执行
要让send_message_to_user任务在相同user_id参数下等待前序任务完成再执行,核心思路是给每个用户分配一个分布式锁,确保同一时间只有一个针对该用户的任务在运行。下面是具体的实现方案:
1. 准备工作
首先确保你的Celery环境配置了Redis作为Broker和Backend(锁依赖Redis实现),并安装redis库:
pip install redis
2. 修改tasks.py实现带锁的任务
在tasks.py中,我们用Redis实现分布式锁,并结合Celery的重试机制来处理锁冲突:
from celery import Celery import redis # 可以根据你的实际配置修改Broker和Backend地址 app = Celery('tasks', broker='redis://localhost:6379/0', backend='redis://localhost:6379/0') # 初始化Redis客户端用于管理锁 redis_client = redis.Redis(host='localhost', port=6379, db=0) @app.task(bind=True) def send_message_to_user(self, user_id, message): # 生成唯一的锁Key,绑定当前用户ID lock_key = f"send_message:{user_id}:lock" # 获取锁:nx=True表示仅当Key不存在时才设置,ex=300设置锁的过期时间(秒),防止任务崩溃导致死锁 lock_acquired = redis_client.set(lock_key, self.request.id, nx=True, ex=300) if not lock_acquired: # 锁未获取到,说明同用户的任务正在执行,触发重试 raise self.retry( exc=Exception(f"Task for user {user_id} is still running"), countdown=5 # 5秒后重试 ) try: # -------------------------- # 这里写你实际发送消息的业务逻辑 print(f"Processing message for user {user_id}: {message}") # 示例:模拟发送耗时操作 # import time # time.sleep(10) # -------------------------- finally: # 任务完成后,仅当锁属于当前任务时才释放,避免误删其他任务的锁 if redis_client.get(lock_key) == self.request.id.encode(): redis_client.delete(lock_key)
3. 方案说明
- 分布式锁:通过Redis的
setnx特性(nx=True)确保同一用户的锁同一时间只能被一个任务持有。 - 锁过期时间:设置
ex=300(5分钟)是为了避免任务意外崩溃(比如worker挂了)导致锁一直存在,后续任务无法执行。你可以根据实际任务的最长执行时间调整这个值。 - 重试机制:当锁获取失败时,Celery会自动重试任务,直到获取到锁并执行。
- 安全释放锁:通过对比锁中的值(当前任务ID)来释放锁,防止其他重试的任务误删正在执行的任务的锁。
替代方案(不推荐高并发场景)
如果你的用户数量不多,也可以通过队列路由实现:给每个用户创建一个专属队列,然后用单个worker消费该队列。但这种方法在用户量大时会导致队列数量爆炸,维护成本高,所以更推荐上面的分布式锁方案。
内容的提问来源于stack exchange,提问作者Edouard Malet
相关产品推荐
相关产品推荐

