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

Django多线程场景下无法查询刚保存的RegistrationEvent实例

问题:多线程创建Django模型时事件流校验失败

我正在开发一个需追踪流程事件的项目,包含Registration和RegistrationEvent模型,后者通过ForeignKey关联前者。为RegistrationEvent编写了_ensure_correct_flow_of_events方法,在model.save时校验事件顺序(需遵循SIGNED -> STARTED -> SUCCESS -> CERTIFICATE_ISSUED,且随时可触发CANCELED),该方法通过_get_previous_event获取关联Registration的最后一个事件。

当创建SUCCESS事件后,save方法会调用Registration.threaded_issue_certificate,该方法在新线程中生成证书并创建CERTIFICATE_ISSUED事件以提升响应速度。但问题是,创建CERTIFICATE_ISSUED时,_get_previous_event无法获取刚创建的SUCCESS事件,反而返回之前的STARTED事件,触发事件流校验错误,报错日志如下:

Checking correct flow, previous event: Course started - admin registration id: 1 current event_type: 3
Checking correct flow, previous event: Course started - admin registration id: 1 current event_type: 4
Exception in thread Thread-2 (threaded_issue_certificate):
Traceback (most recent call last):
  File "/Users/zenodallavalle/miniconda3/lib/python3.10/threading.py", line 1016, in _bootstrap_inner
    self.run()
  File "/Users/zenodallavalle/miniconda3/lib/python3.10/threading.py", line 953, in run
[21/Mar/2024 15:16:41] "POST /admin/main/registrationevent/add/ HTTP/1.1" 302 0
    self._target(*self._args, **self._kwargs)
  File "/Users/zenodallavalle/Downloads/test/main/models.py", line 79, in threaded_issue_certificate
    return self.issue_certificate()
  File "/Users/zenodallavalle/Downloads/test/main/models.py", line 84, in issue_certificate
    RegistrationEvent.objects.create(
  File "/Users/zenodallavalle/Downloads/test/env/lib/python3.10/site-packages/django/db/models/manager.py", line 87, in manager_method
    return getattr(self.get_queryset(), name)(*args, **kwargs)
  File "/Users/zenodallavalle/Downloads/test/env/lib/python3.10/site-packages/django/db/models/query.py", line 679, in create
    obj.save(force_insert=True, using=self.db)
  File "/Users/zenodallavalle/Downloads/test/main/models.py", line 175, in save
    self._ensure_correct_flow_of_events(is_new=is_new)
  File "/Users/zenodallavalle/Downloads/test/main/models.py", line 155, in _ensure_correct_flow_of_events
    raise ValueError(
ValueError: After started next event must be 'success' or 'canceled'

附上models.py代码以复现问题:

from django.db import models
from django.contrib.auth.models import User
from logging import getLogger

from main.utils import make_thread

logger = getLogger(__name__)
import threading


def make_thread(fn):
    def _make_thread(*args, **kwargs):
        thread = threading.Thread(target=fn, args=args, kwargs=kwargs)
        thread.start()
        return thread

    return _make_thread


class Registration(models.Model):
    course_user = models.ForeignKey(
        User,
        on_delete=models.CASCADE,
        related_name="registrations",
    )

    created_by = models.ForeignKey(
        User,
        null=True,
        on_delete=models.SET_NULL,
    )
    created_at = models.DateTimeField(auto_now_add=True)
    updated_at = models.DateTimeField(auto_now=True)

    def _check_user_not_signed_for_other_courses(self):
        other_registrations = self.course_user.registrations.exclude(pk=self.pk)
        if any([not r.ended for r in other_registrations]):
            raise ValueError("User already signed for another course")

    def _ensure_created_by_is_not_null(self):
        if self.created_by is None:
            raise ValueError("Created by is null")

    def save(self, *args, **kwargs) -> None:
        is_new = self._state.adding
        if is_new:
            self._ensure_created_by_is_not_null()
            self._check_user_not_signed_for_other_courses()
        ret = super().save(*args, **kwargs)
        if is_new:
            RegistrationEvent.objects.create(
                course_registration=self,
                event_type=RegistrationEvent.EventType.SIGNED,
            )
        return ret

    def __str__(self):
        return f"{self.course_user} registration id: {self.pk}"

    def __repr__(self):
        return f"<Registration: {self.course_user} registration id: {self.pk}>"

    @property
    def ended(self):
        return self.events.filter(
            event_type__in=(
                RegistrationEvent.EventType.CERTIFICATE_ISSUED,
                RegistrationEvent.EventType.CANCELED,
                RegistrationEvent.EventType.FAILED,
            )
        ).exists()

    @property
    def last_event(self):
        return self.events.order_by("-created_at").first()

    @make_thread
    def threaded_issue_certificate(self):
        return self.issue_certificate()

    def issue_certificate(self):
        # Do something here
        # Register it as an event
        RegistrationEvent.objects.create(
            course_registration=self,
            event_type=RegistrationEvent.EventType.CERTIFICATE_ISSUED,
        )


class RegistrationEvent(models.Model):
    class EventType(models.IntegerChoices):
        SIGNED = 1, "Signed up"
        STARTED = 2, "Course started"
        SUCCESS = 3, "Course success"
        CERTIFICATE_ISSUED = 4, "Certificate issued"
        CANCELED = 5, "Cancelled"
        FAILED = 6, "Course failed"

    course_registration = models.ForeignKey(
        Registration,
        on_delete=models.CASCADE,
        related_name="events",
    )
    event_type = models.IntegerField(choices=EventType.choices)

    created_at = models.DateTimeField(auto_now_add=True, verbose_name="Creato il")
    updated_at = models.DateTimeField(auto_now=True, verbose_name="Aggiornato il")

    @property
    def event_type_description(self):
        return self.EventType(self.event_type).label

    def __str__(self):
        return f"{self.event_type_description} - {self.course_registration}"

    def __repr__(self):
        return f"<RegistrationEvent: {self.event_type_description} - {self.course_registration}>"

    def _get_previous_event(self, is_new):
        qs = self.course_registration.events.all()
        if not is_new:
            qs = qs.exclude(created_at__gte=self.created_at)
        return qs.order_by("-created_at").first()

    def _ensure_correct_flow_of_events(self, is_new):
        # If the course is completed (certificate issued or cancelled), no more events can be added
        if is_new:
            if self.course_registration.ended:
                raise ValueError(
                    "Il corso è completato, non è possibile aggiungere eventi"
                )

        if self.event_type == self.EventType.CANCELED:
            return  # No further checks needed

        previous_event = self._get_previous_event(is_new=is_new)
        print(
            "Checking correct flow, previous event:",
            previous_event,
            "current event_type:",
            self.event_type,
        )

        if not previous_event:
            if self.event_type != self.EventType.SIGNED:
                raise ValueError("First event must be 'signed'")
        elif previous_event.event_type == self.EventType.SIGNED:
            if self.event_type != self.EventType.STARTED:
                raise ValueError("After signed next event must be 'started'")
        elif previous_event.event_type == self.EventType.STARTED:
            if self.event_type not in (
                self.EventType.SUCCESS,
                self.EventType.CANCELED,
            ):
                raise ValueError(
                    "After started next event must be 'success' or 'canceled'"
                )
        elif previous_event.event_type == self.EventType.SUCCESS:
            if self.event_type != self.EventType.CERTIFICATE_ISSUED:
                raise ValueError(
                    "After success next event must be 'certificate issued'"
                )

    def _issue_certificate_if_needed(self, is_new):
        if not is_new:
            return
        if not self.event_type == self.EventType.SUCCESS:
            return
        self.course_registration.threaded_issue_certificate()

    def save(self, *args, **kwargs):
        is_new = self._state.adding

        if is_new:
            self._ensure_correct_flow_of_events(is_new=is_new)

        ret = super().save(*args, **kwargs)

        self._issue_certificate_if_needed(is_new=is_new)
        return ret
问题原因

核心问题在于Django的数据库连接线程隔离机制以及事务提交时机:

  • 线程间数据库连接独立:Django为每个线程维护独立的数据库连接。主线程创建SUCCESS事件后,虽然已经调用super().save()写入数据库,但如果当前处于事务中(比如Django admin的默认事务行为),事务还未提交,新线程的数据库连接无法看到未提交的SUCCESS事件。
  • 事务提交顺序冲突:RegistrationEvent.save()的流程是:先校验事件流→保存SUCCESS事件→启动新线程。但主线程的事务(如admin操作的事务)可能在新线程执行查询时还未提交,新线程查询到的是事务提交前的旧数据,自然拿不到刚创建的SUCCESS事件,只能读取到之前的STARTED事件,触发校验失败。
  • 校验逻辑依赖实时数据库数据:_get_previous_event直接从数据库查询前序事件,而线程间的事务隔离导致新线程无法读取主线程未提交的数据,出现数据不一致。

解决思路参考

  • 等待事务提交后启动线程:使用Django的transaction.on_commit()钩子,确保主线程事务提交后再启动证书生成线程,让新线程能读取到最新的SUCCESS事件。
  • 直接传递前序事件信息:启动线程时,把刚创建的SUCCESS事件实例或ID传入,线程中无需查询数据库,直接基于已知的SUCCESS事件创建CERTIFICATE_ISSUED,跳过依赖数据库的前序事件查询。
  • 调整校验逻辑:允许CERTIFICATE_ISSUED事件直接关联指定的SUCCESS事件,而非依赖查询数据库的最后一个事件,避免线程间数据不一致问题。

内容的提问来源于stack exchange,提问作者Zeno Dalla Valle

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.27 16:14:52