Django+Celery+Redis环境下同一任务重复执行问题求助
问题描述
用Django、Celery配合Redis作为消息中间件实现后台任务,通过TryoutSubmission模型的post_save信号触发任务执行,但同一任务会被重复执行上千次(其他debug/add任务运行正常)。已尝试Celery官方「确保任务仅执行一次」的方案,未能解决问题。
相关代码
models.py
from django.db.models.signals import post_save from django.dispatch import receiver from student.tasks import create_student_subject_tryout_result, add @receiver(post_save, sender=TryoutSubmission, dispatch_uid='create_student_subject_tryout_result') def result_calculation(sender, instance, **kwargs): if instance.status == 'C': print('Calculating result') create_student_subject_tryout_result.delay(instance.student.id, instance.tryout.id)
celery.py
import os from celery import Celery os.environ.setdefault('DJANGO_SETTINGS_MODULE', 'eec.settings') app = Celery('eec') app.config_from_object('django.conf:settings', namespace='CELERY') app.conf.broker_transport_options = {'visibility_timeout': 3600} app.autodiscover_tasks() @app.task(bind=True) def debug_task(self): print(f'Request: {self.request!r}')
tasks.py
from celery import shared_task import tryout.models @shared_task(bind=True) def create_student_subject_tryout_result(self, student_id, tryout_id): tryout_submission=tryout.models.TryoutSubmission.objects.get( student_id=student_id, tryout_id=tryout_id ) tryout_questions = tryout_submission.tryout.tryoutquestion_set.all().count() answered_qs = tryout_submission.tryout.tryoutanswersubmission_set.filter( is_answered=True).count() correct_ans = tryout_submission.tryout.tryoutanswersubmission_set.filter( is_correct=True).count() tryout_submission.total_questions = tryout_questions tryout_submission.answered_questions = answered_qs tryout_submission.correct_answers = correct_ans tryout_submission.total_time = tryout_submission.end_time - tryout_submission.start_time tryout_submission.save() return "Result created"
settings.py
CELERY_RESULT_BACKEND = 'django-db' CELERY_CACHE_BACKEND = 'django-cache' CELERY_BROKER_URL = 'redis://localhost:6379' CELERY_ACCEPT_CONTENT = ['application/json'] CELERY_TASK_SERIALIZER = 'json' CELERY_RESULT_SERIALIZER = 'json' CELERY_TIMEZONE = 'Asia/Kolkata' CELERY_BEAT_SCHEDULER = 'django_celery_beat.schedulers:DatabaseScheduler'
问题排查与解决方案
核心问题是任务执行时触发了新的post_save信号,形成无限循环:任务create_student_subject_tryout_result中调用了tryout_submission.save(),这会再次触发TryoutSubmission的post_save信号,信号接收器又会发送新的任务,反复循环导致任务被执行上千次。
解决方案1:使用update()替代save()(推荐)
直接通过ORM的update()方法更新数据库,不会触发post_save信号,彻底避免循环:
修改tasks.py中的任务逻辑:
@shared_task(bind=True) def create_student_subject_tryout_result(self, student_id, tryout_id): tryout_submission = tryout.models.TryoutSubmission.objects.get( student_id=student_id, tryout_id=tryout_id ) tryout_questions = tryout_submission.tryout.tryoutquestion_set.all().count() answered_qs = tryout_submission.tryout.tryoutanswersubmission_set.filter( is_answered=True).count() correct_ans = tryout_submission.tryout.tryoutanswersubmission_set.filter( is_correct=True).count() total_time = tryout_submission.end_time - tryout_submission.start_time # 用update()直接更新,不触发post_save信号 tryout.models.TryoutSubmission.objects.filter( student_id=student_id, tryout_id=tryout_id ).update( total_questions=tryout_questions, answered_questions=answered_qs, correct_answers=correct_ans, total_time=total_time ) return "Result created"
解决方案2:在保存时跳过信号触发
如果必须使用save(),可以给模型添加自定义逻辑,允许保存时跳过信号:
- 修改
TryoutSubmission模型:
class TryoutSubmission(models.Model): # 原有字段... def save(self, skip_post_save=False, *args, **kwargs): self._skip_post_save = skip_post_save super().save(*args, **kwargs)
- 修改信号接收器,判断是否跳过:
@receiver(post_save, sender=TryoutSubmission, dispatch_uid='create_student_subject_tryout_result') def result_calculation(sender, instance, **kwargs): # 如果是跳过信号的保存,直接返回 if hasattr(instance, '_skip_post_save') and instance._skip_post_save: return if instance.status == 'C': print('Calculating result') create_student_subject_tryout_result.delay(instance.student.id, instance.tryout.id)
- 在任务中调用
save()时传入参数:
# 替换原有tryout_submission.save() tryout_submission.save(skip_post_save=True)
额外检查点
- 确认
dispatch_uid唯一,避免信号被重复注册 - 检查Celery worker是否有多个实例重复消费任务(但其他任务正常,此可能性较低)
内容的提问来源于stack exchange,提问作者Kuldeep
相关产品推荐
相关产品推荐

