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

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.
原因分析与解决方案

问题原因

  1. 多进程数据库连接复用失效:Django的数据库连接绑定到进程,子进程会继承父进程的连接,但这个继承的连接在子进程中是无效状态。当子进程用无效连接执行事务时,会触发数据库端主动关闭连接,引发异常。
  2. 事务并发锁竞争与连接耗尽:每个save操作都开启事务,并发执行时会对分组数据加共享锁,大量并发请求会导致锁等待队列过长,数据库连接池被耗尽,最终引发数据库服务异常终止。
  3. 竞态条件导致逻辑失效:当前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,提问作者Ярига Олег

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.17 07:40:39