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

PostgreSQL大数据批量插入优化方案问询(Django+Python环境)

批量数据插入优化与连接异常排查

环境配置

  • OS:Ubuntu 22.04 LTS
  • 数据库:PostgreSQL 14
  • 开发框架:Python 3.11 + Django

当前现状与问题

目前采用每次插入10万行的INSERT语句,插入100万行数据平均耗时2分钟(时长可接受),但希望找到更优实现方式。近期出现耗时增加,且偶尔抛出错误:

OperationalError: (psycopg2.OperationalError) server closed the connection unexpectedly

当前实现代码

from django.db import connection
cursor = connection.cursor()

batch_size = 100000

offset = 0
while True:
    transaction_list_query = f"SELECT * FROM {source_table} LIMIT {batch_size} OFFSET {offset}; "
    cursor.execute(transaction_list_query)
    transaction_list = dictfetchall(cursor)
    if not transaction_list:
        break
    data_to_insert = []
    for transaction in transaction_list:
        # 一些计算密集型处理

    insert_query = f"INSERT INTO {per_transaction_table} ({company_ref_id_id_column}, {rrn_column},{transaction_type_ref_id_id_column}, {transactionamount_column}) VALUES  {','.join(data_to_insert)} ON CONFLICT ({rrn_column}) DO UPDATE SET {company_ref_id_id_column} = EXCLUDED.{company_ref_id_id_column};"
    cursor.execute(insert_query)
    offset += batch_size

优化实现方案

1. 替换OFFSET分页为键值分页

OFFSET在数据量增大时性能会急剧下降(数据库需扫描前置所有行定位),改用基于主键/有序唯一键的分页:

last_id = 0
while True:
    # 假设源表有自增主键id
    transaction_list_query = f"SELECT * FROM {source_table} WHERE id > {last_id} ORDER BY id LIMIT {batch_size};"
    cursor.execute(transaction_list_query)
    transaction_list = dictfetchall(cursor)
    if not transaction_list:
        break
    # 数据处理逻辑...
    # 更新批次最后一条数据的主键值
    last_id = transaction_list[-1]['id']

2. 用Django原生批量操作替代手动SQL拼接

Django的bulk_create内置优化逻辑,配合冲突更新(Django 4.0+支持)更安全高效:

from django.db import transaction
from your_app.models import PerTransaction, SourceTable  # 替换为实际模型类

with transaction.atomic():
    last_id = 0
    while True:
        # 批量拉取源数据
        source_records = SourceTable.objects.filter(id__gt=last_id).order_by('id')[:batch_size]
        if not source_records:
            break
        # 生成待插入对象列表
        insert_objs = []
        for record in source_records:
            # 计算密集型处理
            obj = PerTransaction(
                company_ref_id_id=...,
                rrn=record.rrn,
                transaction_type_ref_id_id=...,
                transactionamount=...
            )
            insert_objs.append(obj)
        # 批量插入+冲突更新
        PerTransaction.objects.bulk_create(
            insert_objs,
            update_conflicts=True,
            update_fields=['company_ref_id_id'],
            unique_fields=['rrn']
        )
        last_id = source_records.last().id

3. 数据库层面优化

  • 临时禁用索引/触发器:插入前关闭目标表的索引和触发器,插入完成后重建(注意:操作期间表查询性能会下降)
    ALTER TABLE {per_transaction_table} DISABLE TRIGGER ALL;
    -- 执行插入操作
    ALTER TABLE {per_transaction_table} ENABLE TRIGGER ALL;
    
  • 调整PostgreSQL配置:增大work_mem、maintenance_work_mem提升批量操作内存,优化shared_buffers,缩短idle_in_transaction_session_timeout避免长事务超时。

4. 计算任务优化

如果计算逻辑耗时占比高:

  • 用多进程/线程并行处理批次内数据(注意数据库连接的线程安全)
  • 将计算逻辑迁移到数据库端(编写PostgreSQL函数处理源数据,减少Python与数据库的数据传输)

连接异常排查

针对server closed the connection unexpectedly的常见解决方向:

  • 连接超时:检查PostgreSQL的idle_in_transaction_session_timeout配置,确保批次处理时及时提交事务、释放连接;避免长事务占用连接。
  • 内存不足:10万行数据批量处理可能占用大量内存,导致OOM杀死进程/连接。可缩小batch_size(如降至5万),或优化数据处理时的内存占用(避免加载不必要的字段)。
  • 数据库负载过高:插入期间监控数据库CPU、IO使用率,避开高峰时段执行,或优化插入语句的执行计划。
  • 网络问题:检查服务器间网络稳定性,排查是否有防火墙拦截、网络丢包情况。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.03 00:51:13