使用psycopg COPY命令复制千万级数据集过慢,求性能优化方案
PostgreSQL跨库批量数据复制性能优化方案
原代码采用逐行循环读写的方式处理1000万+行数据,由于每次IO操作粒度极小,导致Python层调用和系统IO开销剧增,最终复制耗时过长。以下是针对性的优化方案:
核心优化:块级批量读写替代逐行循环
直接使用shutil.copyfileobj实现块级数据流对接,替代逐行遍历,大幅减少IO调用次数和Python层开销,这是最直接的性能提升手段:
import shutil import psycopg from psycopg import sql with psycopg.connect(dsn_src) as conn1, psycopg.connect(dsn_tgt) as conn2: # 安全构造COPY语句,避免SQL注入 copy_to_sql = sql.SQL("COPY {} TO STDOUT (FORMAT BINARY)").format(sql.Identifier(table)) copy_from_sql = sql.SQL("COPY {} FROM STDIN (FORMAT BINARY)").format(sql.Identifier(table)) with conn1.cursor().copy(copy_to_sql) as copy1: with conn2.cursor().copy(copy_from_sql) as copy2: # 采用1MB块大小读写,可根据内存情况调整(如2MB/4MB) shutil.copyfileobj(copy1, copy2, length=1024*1024) conn2.commit()
进阶优化方案
1. 利用PostgreSQL原生工具(性能最优)
直接使用pg_dump和pg_restore命令行工具,完全绕过Python中转,这是跨库复制大表的性能天花板:
# 命令行直接执行 pg_dump -d "postgresql://user:pass@source_host/db" -t target_table --format=c | pg_restore -d "postgresql://user:pass@target_host/db" -t target_table
若需在Python中调用,可通过subprocess执行:
import subprocess # 构造命令 pg_dump_cmd = [ "pg_dump", "-d", dsn_src, "-t", table, "--format=c" ] pg_restore_cmd = [ "pg_restore", "-d", dsn_tgt, "-t", table ] # 建立管道对接数据流 dump_proc = subprocess.Popen(pg_dump_cmd, stdout=subprocess.PIPE) restore_proc = subprocess.Popen(pg_restore_cmd, stdin=dump_proc.stdout) # 关闭子进程stdout,避免资源泄漏 dump_proc.stdout.close() # 等待复制完成 restore_proc.wait()
2. 调整数据库连接与配置参数
- 增大内存分配:连接时临时调高
work_mem,提升大表处理的内存缓存能力:# 连接时设置work_mem为64MB conn1 = psycopg.connect(dsn_src, options="-c work_mem=64MB") conn2 = psycopg.connect(dsn_tgt, options="-c work_mem=64MB") - 临时关闭目标库约束:复制前禁用目标表的触发器和约束(需确保数据一致性),复制完成后恢复:
with conn2.cursor() as cur: cur.execute(sql.SQL("ALTER TABLE {} DISABLE TRIGGER ALL").format(sql.Identifier(table))) # 复制完成后恢复 with conn2.cursor() as cur: cur.execute(sql.SQL("ALTER TABLE {} ENABLE TRIGGER ALL").format(sql.Identifier(table))) - 优化WAL日志配置:临时增大目标库的
max_wal_size、延长checkpoint_timeout,减少写入时的刷盘频率,复制完成后改回原配置。
3. 分批次并行复制(超大规模表)
对于亿级行的超大规模表,可按主键/时间范围分批次复制,同时支持多进程并行处理(需注意避免锁冲突):
import shutil import psycopg from psycopg import sql from multiprocessing import Pool batch_size = 200000 total_rows = 0 # 获取总行数 with psycopg.connect(dsn_src) as conn: with conn.cursor() as cur: cur.execute(sql.SQL("SELECT COUNT(*) FROM {}").format(sql.Identifier(table))) total_rows = cur.fetchone()[0] # 定义单批次复制函数 def copy_batch(offset): with psycopg.connect(dsn_src) as conn1, psycopg.connect(dsn_tgt) as conn2: copy_to_sql = sql.SQL("COPY (SELECT * FROM {} LIMIT %s OFFSET %s) TO STDOUT (FORMAT BINARY)").format(sql.Identifier(table)) copy_from_sql = sql.SQL("COPY {} FROM STDIN (FORMAT BINARY)").format(sql.Identifier(table)) with conn1.cursor().copy(copy_to_sql, params=(batch_size, offset)) as copy1: with conn2.cursor().copy(copy_from_sql) as copy2: shutil.copyfileobj(copy1, copy2, length=2*1024*1024) conn2.commit() # 生成批次偏移量列表 offsets = list(range(0, total_rows, batch_size)) # 启用4进程并行复制(可根据CPU核数调整) with Pool(processes=4) as pool: pool.map(copy_batch, offsets)
4. 网络层面优化
- 优先使用内网IP连接两个数据库,避免公网带宽瓶颈
- 调整PostgreSQL的
tcp_keepalives参数,避免长连接超时:conn = psycopg.connect(dsn_src, keepalives_idle=60, keepalives_interval=10)
内容的提问来源于stack exchange,提问作者rcs
相关产品推荐
相关产品推荐

