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

