Django Celery异步任务与事务:重复唯一字段实例创建问题
会出现的问题
当第二个任务执行时,因为第一个任务已经成功保存了带有相同唯一标识符的MyModel实例,第二个任务走到serializer.save()步骤时,会触发数据库唯一约束冲突错误(比如PostgreSQL的UniqueViolation、MySQL的Duplicate entry),直接导致任务执行失败;如果配置了Celery重试机制,还可能触发无意义的重试,但本质问题是重复插入违反了唯一字段的约束。
这里要注意:虽然每个任务都用了transaction.atomic(),但这是各自独立的事务——第一个任务的事务提交后,第二个任务的事务再尝试插入相同标识符的记录,数据库层面会直接拒绝这个操作,事务回滚,任务报错。
解决方法
针对这个重复提交导致的幂等性问题,有几种不同层级的解决方案,你可以根据业务场景选择:
1. 任务提交前做前置校验(简单但需注意并发)
在提交Celery任务之前,先查询数据库是否已经存在该标识符的实例:
# 提交任务前的代码 if not MyModel.objects.filter(identifier=model_identifier).exists(): create_model.delay(model_identifier) else: # 处理已存在的情况,比如返回提示、直接复用已有实例等 pass
⚠️ 注意:这种方法在高并发场景下可能失效——如果两个请求同时校验,都查到不存在,然后同时提交任务,还是会出现冲突。所以适合并发量不高的场景,或者作为辅助手段。
2. 在任务内部实现幂等性(最可靠的数据库层面保障)
把校验和写入放在同一个事务里,确保原子性:
方案A:先查询再创建
调整任务代码,在事务内先检查是否存在,不存在再创建:@shared_task def create_model(model_identifier): with transaction.atomic(): # 先查询,避免重复创建 if MyModel.objects.filter(identifier=model_identifier).exists(): # 可以返回已有实例ID,或者直接结束任务 return MyModel.objects.get(identifier=model_identifier).id # 不存在则继续创建流程 serializer = MyModelSerializer(data=model_data) serializer.is_valid(raise_exception=True) # 建议加上raise_exception,明确校验错误 serializer.save() # ...更多操作为了避免高并发下的间隙锁问题,还可以用
select_for_update()加锁:# 替换上面的查询逻辑 try: instance = MyModel.objects.select_for_update().get(identifier=model_identifier) return instance.id except MyModel.DoesNotExist: # 执行创建逻辑 pass方案B:使用
get_or_create简化逻辑
Django的get_or_create方法本身就是原子性的,会自动处理“查询-不存在则创建”的流程,并且依赖数据库的唯一约束保证幂等:@shared_task def create_model(model_identifier): with transaction.atomic(): # 假设model_data里除了identifier还有其他字段 instance, created = MyModel.objects.get_or_create( identifier=model_identifier, defaults={k: v for k, v in model_data.items() if k != 'identifier'} ) if not created: # 如果是已存在的实例,可根据需求处理后续操作(比如跳过、更新字段等) return instance.id # ...只有创建了新实例才执行的后续操作方案C:捕获唯一约束异常
如果业务允许,可以直接捕获数据库的唯一约束异常,优雅处理:from django.db import IntegrityError @shared_task def create_model(model_identifier): try: with transaction.atomic(): serializer = MyModelSerializer(data=model_data) serializer.is_valid(raise_exception=True) serializer.save() # ...更多操作 except IntegrityError as e: # 判断是否是唯一约束冲突(不同数据库的错误信息格式不同,需要适配) if "unique constraint" in str(e).lower(): # 处理已存在的情况,比如返回已有实例 return MyModel.objects.get(identifier=model_identifier).id # 其他完整性错误则重新抛出 raise
3. 使用分布式锁(任务层面的并发控制)
借助Redis等工具实现分布式锁,确保同一个model_identifier对应的任务同一时间只有一个在执行:
import redis from celery.exceptions import Ignore redis_client = redis.Redis() @shared_task def create_model(model_identifier): lock_key = f"create_model_lock:{model_identifier}" # 尝试获取锁,超时时间设为任务最长执行时间 acquired = redis_client.set(lock_key, "locked", nx=True, ex=300) if not acquired: # 锁已被持有,直接结束任务或者重试 raise Ignore("Task already in progress for this identifier") try: with transaction.atomic(): # ...创建逻辑 finally: # 释放锁 redis_client.delete(lock_key)
这种方法适合需要严格控制任务并发执行的场景,避免多个任务同时处理同一个标识符的情况。
内容的提问来源于stack exchange,提问作者user2880391

