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

如何优化十亿行PostgreSQL表规范化转换的执行速度?

优化方案

1. 用纯SQL替代Python循环(核心优化)

当前方案把每个站点的数据拉到Python内存再写回数据库,完全是冗余的性能损耗。直接用数据库内部的关联查询+批量插入,速度能提升几个数量级:

INSERT INTO observations (station_id, {select_set})
SELECT s.id, {select_set}
FROM source src
JOIN stations s ON src.station = s.station;

将{select_set}替换为实际的字段列表(比如date, temp, humidity,...)即可。如果担心单次插入十亿行导致事务过大,可以分批处理:

-- 每次处理100万行,可根据服务器性能调整批次大小
WITH batch AS (
    SELECT src.ctid, s.id, {select_set}
    FROM source src
    JOIN stations s ON src.station = s.station
    LIMIT 1000000
)
INSERT INTO observations (station_id, {select_set})
SELECT id, {select_set} FROM batch;

DELETE FROM source WHERE ctid IN (SELECT ctid FROM batch);

循环执行上述SQL直到source表数据处理完毕,全程不需要Python参与数据流转,彻底解决Python CPU占用过高的问题。

2. 优化Python侧的COPY操作(若坚持用Python)

如果必须用Python,当前的fetchall()和逐行write_row效率极低,改成批量处理:

  • 用fetchmany(size=10000)分批拉取数据,避免内存溢出
  • 将批量数据转换成CSV格式字符串,一次性写入COPY流
def prepend_batch(prepend, batch):
    # 把批量数据转换成带前置station_id的CSV行
    return '\n'.join([f"{prepend},{','.join(map(str, row))}" for row in batch]) + '\n'

for station_id_pair in all_stations:
    station_id = station_id_pair[0][0]
    station_code = station_id_pair[0][1]
    # 用参数化查询替代字符串格式化,避免SQL注入+复用查询计划
    dbcur.execute(f"SELECT {','.join(select_set)} FROM source WHERE station=%s", (station_code,))
    with dbcur.copy(f"COPY observations ({','.join(['station_id'] + select_set)}) FROM STDIN WITH CSV") as copy:
        while True:
            batch = dbcur.fetchmany(10000)
            if not batch:
                break
            copy.write(prepend_batch(station_id, batch))
    dbconn.commit()

3. 调整PostgreSQL配置提升写入性能

  • 临时关闭自动清理:ALTER TABLE observations SET (autovacuum_enabled = false);,迁移完成后再开启
  • 增大内存参数:调高work_mem、maintenance_work_mem,减少磁盘交换
  • 优化WAL写入:调大max_wal_size,降低检查点频率
  • 临时禁用约束/索引:迁移前先删除observations表的索引和外键,完成后再重建,避免边写边维护索引的性能损耗

4. 并行处理(多进程)

利用服务器多核优势,用Python多进程同时处理不同站点,注意控制进程数(建议为CPU核心数的1/2),避免数据库连接过载:

from multiprocessing import Pool

def process_station(station_pair):
    # 每个进程单独创建数据库连接
    conn = psycopg2.connect(...)
    cur = conn.cursor()
    station_id = station_pair[0][0]
    station_code = station_pair[0][1]
    cur.execute(f"SELECT {','.join(select_set)} FROM source WHERE station=%s", (station_code,))
    with cur.copy(f"COPY observations ({','.join(['station_id'] + select_set)}) FROM STDIN WITH CSV") as copy:
        while True:
            batch = cur.fetchmany(10000)
            if not batch:
                break
            copy.write(prepend_batch(station_id, batch))
    conn.commit()
    cur.close()
    conn.close()

if __name__ == '__main__':
    with Pool(4) as p:  # 4个进程,根据实际核心数调整
        p.map(process_station, all_stations)

RAID初始化的性能影响

Dell PERC6/i的RAID10后台初始化肯定会影响性能,初始化过程会占用磁盘IO带宽,和PostgreSQL的读写操作争抢资源,导致数据库读写延迟升高。如果条件允许,建议等RAID初始化完成后再执行数据迁移;如果必须现在操作,可适当减小批量处理的大小,降低数据库IO负载。


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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.14 11:09:50