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

使用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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.19 14:17:08