Celery任务重复执行是因分布式环境下多worker竞争同一任务且无互斥锁,Redis仅保证task_id唯一不保证互斥;需用SETNX实现带过期时间的分布式锁,并结合幂等设计保障最终一致性。

为什么Celery任务会重复执行
不是代码写错了,而是分布式环境下天然存在的竞争问题:多个worker同时从Redis队列里取到同一个task_id,又没加锁,结果都开始执行。尤其在retries开启、网络抖动、worker重启时特别常见。
Redis本身不提供“执行中任务”的全局视图,Celery的task_id只保证唯一性,不保证互斥——这点很多人误以为它自带防重。
用Redis SETNX实现最简分布式锁
不用引入额外库,靠Redis原生命令就能落地。核心是SET key value EX seconds NX:只有key不存在时才设值,并自动过期。
-
NX确保只有一个worker能抢到锁 -
EX避免死锁(比如worker崩溃没释放) - value必须是唯一标识(如
uuid.uuid4().hex),释放锁时要校验,防止误删别人锁
示例(放在task开头):
import redis
import uuid
<p>r = redis.Redis()</p><p>def my_task(user_id):
lock_key = f"lock:my_task:{user_id}"
lock_value = uuid.uuid4().hex</p><h1>尝试上锁,10秒超时</h1><pre class="brush:php;toolbar:false;">if r.set(lock_key, lock_value, ex=10, nx=True):
try:
# 执行业务逻辑
process_user(user_id)
finally:
# 安全释放:先GET再DEL,且只删自己的value
if r.get(lock_key) == lock_value.encode():
r.delete(lock_key)
else:
# 锁已被占,直接退出或记录日志
returnCelery task装饰器封装锁逻辑
每次手动写锁太容易漏,直接封装成可复用的decorator。注意两点:一是锁key要能动态拼接参数,二是必须支持timeout和自动续期(长任务场景)。
- 用
functools.wraps保留原函数签名,否则Celery inspect查不到参数 - 锁超时时间建议略大于任务预估耗时,比如任务通常3秒,设为5秒
- 如果任务可能超过锁超时(如IO阻塞),需配合
redis.lock的extend或用Redlock算法
轻量封装示意:
from functools import wraps
<p>def distributed_lock(lock_key_func, timeout=5):
def decorator(task_func):
@wraps(task_func)
def wrapper(*args, <strong>kwargs):
key = lock_key_func(*args, *<em>kwargs)
lock_value = str(uuid.uuid4())
if r.set(key, lock_value, ex=timeout, nx=True):
try:
return task_func(</em>args, </strong>kwargs)
finally:
if r.get(key) == lock_value.encode():
r.delete(key)
else:
return None # 或 raise Ignore
return wrapper
return decorator</p><h1>使用</h1><p>@distributed_lock(lambda user_id: f"lock:send_email:{user_id}", timeout=8)
def send_email_task(user_id):
...
</p>
比锁更稳妥的方案:幂等+唯一键
锁解决的是“并发执行”,但真正要防的是“重复效果”。很多场景下,与其死磕锁,不如让任务本身幂等。
- 数据库操作加
unique_together或ON CONFLICT DO NOTHING(PostgreSQL) - 发消息前先查
TaskResult.objects.filter(task_id=..., status='SUCCESS') - Celery 5.3+ 支持
acks_late=True+reject_on_worker_lost=True,减少重复入队 - Redis里用
HSET job_status:{task_id} status "started"做状态标记,任务开头先查再写
锁和幂等不是二选一,而是组合使用:锁控执行入口,幂等兜底最终状态。
最容易被忽略的是锁key的粒度——用task_id锁整个任务毫无意义,必须落到业务维度,比如user_id、order_id;还有就是忘记给锁加过期时间,一旦worker卡住,后续所有同key任务全堵死。
Python免费学习笔记(深入):立即使用
在学习笔记中,你将探索 Python 的核心概念和高级技巧!











