Django中如何实现Celery同类型任务互斥执行?不同类型可并发
实现方案
要实现"taskA同一时间仅允许一个实例执行,重复触发时返回错误提示"的需求,有以下几种可行方案,按可靠性和适用场景排序:
方案一:基于Redis的分布式锁(推荐,适合多Worker/分布式场景)
利用Redis的原子性操作实现分布式锁,确保同一时间只有一个taskA被触发,即使在多Worker环境下也能可靠工作。
步骤:
- 安装Redis客户端:
pip install redis - 在Django配置中添加Redis连接信息(如
settings.py):
REDIS_HOST = "localhost" REDIS_PORT = 6379 REDIS_DB = 0 REDIS_PASSWORD = "" # 如有密码请填写
- 修改任务和触发逻辑:
import redis from django.conf import settings from celery import shared_task from django.http import JsonResponse # 初始化Redis客户端 redis_client = redis.Redis( host=settings.REDIS_HOST, port=settings.REDIS_PORT, db=settings.REDIS_DB, password=settings.REDIS_PASSWORD, decode_responses=True ) @shared_task def taskA(): try: # 这里编写taskA的业务逻辑 pass finally: # 无论任务成功还是失败,都释放锁 redis_client.delete("taskA_running_lock") # 触发taskA的视图函数 def trigger_taskA(request): # 尝试获取锁,ex设置为任务预计最长执行时间(单位:秒),防止任务异常退出导致死锁 lock_acquired = redis_client.set("taskA_running_lock", "1", nx=True, ex=300) if not lock_acquired: return JsonResponse({"error": "taskA正在执行中,请稍后再试"}, status=400) # 锁获取成功,触发任务 taskA.delay() return JsonResponse({"message": "taskA已开始执行"})
注意事项:
- 调整
ex参数的值,确保其大于taskA的最长执行时间,避免锁提前过期导致重复触发。 - 如果任务执行时间可能超过锁的过期时间,可以在任务中定期调用
redis_client.expire("taskA_running_lock", 300)刷新锁的有效期。
方案二:用Celery Inspect检查运行中任务(适合单Worker/低并发场景)
通过Celery的inspectAPI查询当前运行中的任务,判断是否有taskA在执行。此方案无需额外依赖,但在多Worker或高并发场景下可能存在状态延迟或并发竞态问题。
代码示例:
from celery import current_app from django.http import JsonResponse from .tasks import taskA def trigger_taskA(request): inspect = current_app.control.inspect() running_tasks = inspect.active() # 获取所有Worker上正在运行的任务 # 遍历检查是否有taskA在运行 taskA_is_running = False if running_tasks: for worker_tasks in running_tasks.values(): for task in worker_tasks: # 替换为你的taskA的完整路径(如"myapp.tasks.taskA") if task["name"] == "path.to.your.taskA": taskA_is_running = True break if taskA_is_running: break if taskA_is_running: return JsonResponse({"error": "taskA正在执行中,请稍后再试"}, status=400) taskA.delay() return JsonResponse({"message": "taskA已开始执行"})
局限性:
inspect的结果存在一定延迟,可能导致短时间内的并发请求绕过检查。- 多Worker环境下,Worker之间的状态同步可能不及时,存在误判风险。
方案三:基于Django ORM的状态锁(适合依赖数据库的场景)
通过Django模型记录taskA的运行状态,结合数据库事务和行级锁实现并发控制。
步骤:
- 创建状态模型(
models.py):
from django.db import models class TaskARunningStatus(models.Model): is_running = models.BooleanField(default=False) updated_at = models.DateTimeField(auto_now=True) class Meta: verbose_name = "taskA运行状态" verbose_name_plural = verbose_name
- 执行迁移:
python manage.py makemigrations && python manage.py migrate - 修改任务和触发逻辑:
from django.db import transaction from django.http import JsonResponse from celery import shared_task from .models import TaskARunningStatus @shared_task def taskA(status_pk): try: # 编写taskA的业务逻辑 pass finally: # 任务完成后更新状态 status = TaskARunningStatus.objects.get(pk=status_pk) status.is_running = False status.save() def trigger_taskA(request): # 确保只有一条状态记录 status, _ = TaskARunningStatus.objects.get_or_create(pk=1) with transaction.atomic(): # 行级锁,防止并发修改 locked_status = TaskARunningStatus.objects.select_for_update().get(pk=1) if locked_status.is_running: return JsonResponse({"error": "taskA正在执行中,请稍后再试"}, status=400) locked_status.is_running = True locked_status.save() taskA.delay(status_pk=1) return JsonResponse({"message": "taskA已开始执行"})
注意事项:
- 依赖数据库的事务和行级锁,性能略低于Redis锁。
- 需确保数据库支持行级锁(如MySQL InnoDB、PostgreSQL)。
内容的提问来源于stack exchange,提问作者musical_ant
相关产品推荐
相关产品推荐

