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
相关产品推荐
相关产品推荐

