You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

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(),可以给模型添加自定义逻辑,允许保存时跳过信号:

  1. 修改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)
  1. 修改信号接收器,判断是否跳过:
@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)
  1. 在任务中调用save()时传入参数:
# 替换原有tryout_submission.save()
tryout_submission.save(skip_post_save=True)

额外检查点

  • 确认dispatch_uid唯一,避免信号被重复注册
  • 检查Celery worker是否有多个实例重复消费任务(但其他任务正常,此可能性较低)

内容的提问来源于stack exchange,提问作者Kuldeep

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.08.25 02:54:26