Django+PostgreSQL事务并发递增字段触发OperationalError问题排查
问题
我有一个关联到分组(bunch)的模型,同一分组内的每条记录需要拥有唯一的version(按数据库中记录的创建顺序生成)。我认为最佳实现方式是使用事务,但在并行执行事务块时遇到了问题:移除transaction.atomic()块后代码可正常运行,但version字段无法保证唯一性。
我编写了测试代码来验证数据库中记录version字段的并发递增逻辑:
def _save_instance(instance): time = random.randint(1, 50) sleep(time/1000) instance.text = str(time) instance.save() def _parallel(): instances = MyModel.objects.all() # clear version print('-- clear old numbers -- ') instances.update(version=None) processes = [] for instance in instances: p = Process(target=_save_instance, args=(instance,)) processes.append(p) print('-- launching -- ') for p in processes: p.start() for p in processes: p.join() sleep(1) ... # assertions to check if versions are correct in one bunch print('parallel Ok!')
MyModel的save()方法定义如下:
... def save(self, *args, **kwargs) -> None: with transaction.atomic(): if not self.number and self.banch_id: max_number = MyModel.objects.filter( banch_id=self.banch_id ).aggregate(max_number=models.Max('version'))['max_number'] self.version = max_number + 1 if max_number else 1 super().save(*args, **kwargs)
当我用随机数量(30-300条)的记录运行测试代码时,出现以下错误:
django.db.utils.OperationalError: server closed the connection unexpectedly This probably means the server terminated abnormally before or while processing the request.
之后所有进程都会卡住,只能通过KeyboardInterrupt终止脚本。完整的进程栈追踪信息如下:
Process Process-14: Traceback (most recent call last): File "/usr/local/lib/python3.6/site-packages/django/db/backends/utils.py", line 86, in _execute return self.cursor.execute(sql, params) psycopg2.OperationalError: server closed the connection unexpectedly This probably means the server terminated abnormally before or while processing the request. The above exception was the direct cause of the following exception: Traceback (most recent call last): File "/usr/local/lib/python3.6/multiprocessing/process.py", line 258, in _bootstrap self.run() File "/usr/local/lib/python3.6/multiprocessing/process.py", line 93, in run self._target(*self._args, **self._kwargs) File "/app/scripts/test_concurrent_saving.py", line 17, in _save_instance instance.save() File "/app/apps/incident/models.py", line 385, in save ).aggregate(max_number=models.Max('version'))['max_number'] File "/usr/local/lib/python3.6/site-packages/django/db/models/query.py", line 384, in aggregate return query.get_aggregation(self.db, kwargs) File "/usr/local/lib/python3.6/site-packages/django/db/models/sql/query.py", line 503, in get_aggregation result = compiler.execute_sql(SINGLE) File "/usr/local/lib/python3.6/site-packages/django/db/models/sql/compiler.py", line 1152, in execute_sql cursor.execute(sql, params) File "/usr/local/lib/python3.6/site-packages/django/db/backends/utils.py", line 100, in execute return super().execute(sql, params) File "/usr/local/lib/python3.6/site-packages/raven/contrib/django/client.py", line 123, in execute return real_execute(self, sql, params) File "/usr/local/lib/python3.6/site-packages/django/db/backends/utils.py", line 68, in execute return self._execute_with_wrappers(sql, params, many=False, executor=self._execute) File "/usr/local/lib/python3.6/site-packages/django/db/backends/utils.py", line 77, in _execute_with_wrappers return executor(sql, params, many, context) File "/usr/local/lib/python3.6/site-packages/django/db/backends/utils.py", line 86, in _execute return self.cursor.execute(sql, params) File "/usr/local/lib/python3.6/site-packages/django/db/utils.py", line 90, in __exit__ raise dj_exc_value.with_traceback(traceback) from exc_value File "/usr/local/lib/python3.6/site-packages/django/db/backends/utils.py", line 86, in _execute return self.cursor.execute(sql, params) django.db.utils.OperationalError: server closed the connection unexpectedly This probably means the server terminated abnormally before or while processing the request.
原因分析与解决方案
问题原因
- 多进程数据库连接复用失效:Django的数据库连接绑定到进程,子进程会继承父进程的连接,但这个继承的连接在子进程中是无效状态。当子进程用无效连接执行事务时,会触发数据库端主动关闭连接,引发异常。
- 事务并发锁竞争与连接耗尽:每个save操作都开启事务,并发执行时会对分组数据加共享锁,大量并发请求会导致锁等待队列过长,数据库连接池被耗尽,最终引发数据库服务异常终止。
- 竞态条件导致逻辑失效:当前
MAX(version)+1的逻辑即使在事务中也存在竞态——多个事务同时读取到相同的max_number,生成重复version值;而事务的引入又进一步加剧了连接和锁的压力。
解决方案
1. 修复多进程数据库连接问题
在子进程执行数据库操作前,强制重置数据库连接:
def _save_instance(instance): # 重置子进程的数据库连接 from django.db import connection connection.close() time = random.randint(1, 50) sleep(time/1000) instance.text = str(time) instance.save()
2. 用数据库原子操作替代应用层计算
避免在应用层计算最大值,改用数据库级别的原子递增逻辑,彻底消除竞态:
from django.db.models import F def save(self, *args, **kwargs) -> None: if not self.version and self.banch_id: # 原子更新分组的最大version,无记录则设为1 updated = MyModel.objects.filter( banch_id=self.banch_id ).update(version=F('version') + 1) if not updated: self.version = 1 else: self.version = MyModel.objects.filter( banch_id=self.banch_id ).aggregate(max_version=models.Max('version'))['max_version'] super().save(*args, **kwargs)
3. 优化并发测试与数据库配置
- 降低并发进程数,避免瞬间打满数据库连接池
- 若业务允许,改用
threading替代multiprocessing,Django对多线程的连接管理更友好 - 在
settings.py中调整数据库连接池参数,比如增大MAX_CONNS、合理设置CONN_MAX_AGE
4. 添加数据库层面的唯一约束
从数据库层面保证version的唯一性,即使应用层逻辑出错也能阻止脏数据:
class MyModel(models.Model): banch_id = models.ForeignKey(Banch, on_delete=models.CASCADE) version = models.IntegerField(null=True) # 其他字段... class Meta: unique_together = ('banch_id', 'version')
内容的提问来源于stack exchange,提问作者Ярига Олег
相关产品推荐
相关产品推荐

