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

16亿行PostgreSQL数据库Python批量更新慢的优化求助

16亿行PostgreSQL数据更新性能优化方案

问题背景

有一个包含约16亿行的PostgreSQL数据库,需基于Python计算结果更新特定列,但当前处理速度极慢——12小时仅完成约3%的行数。通过Windows任务管理器观察,内存与CPU尚有剩余,但磁盘占用率达100%且未达到NVME的读写速度。数据库以id为主键,当前使用psycopg2的操作代码如下:

import psycopg2


def retrieve_raw_content_batch(batch_size):
    with db_connect() as conn:
        with conn.cursor('testit') as cursor:
            cursor.execute("SELECT id, columnoftext FROM table;")
            while True:
                rows = cursor.fetchmany(batch_size)
                if not rows:
                    break

                yield rows


def update_database(upload_list):
    with db_connect() as conn:
        with conn.cursor() as cursor:
            update_query = "UPDATE table SET col1 = %s, col2 = %s WHERE id = %s"
            psycopg2.extras.execute_batch(cursor, update_query, upload_list)


def do_stuff(row_batch): 
    for rows in row_batch:
        upload_list = []
        for row in rows:
            #calculate to get id, col1, col2
            upload_list.append((id, col1, col2))
        update_database(upload_list)


def main(batch_size):
    rows_batch = retrieve_raw_content_batch(batch_size)
    do_stuff(rows_batch)

已尝试将max_wal_size调整为10GB,但对Postgres配置优化较为陌生,同时考虑是否改用COPY创建新表后再通过JOIN关联的方式替代逐行UPDATE。


核心瓶颈分析

当前方案的性能问题根源在于频繁批量UPDATE带来的高磁盘IO开销:PostgreSQL的UPDATE是原地更新逻辑,需要查找行、标记旧版本、写入新版本,会产生大量随机IO和WAL日志。这也是磁盘占满但未跑满NVME速度的原因——随机IO的性能远低于顺序IO,且WAL日志的同步写入会进一步拖慢速度。


优化方案

1. 改用COPY+替换表方案(首推)

对于超大规模数据集,直接UPDATE是效率最低的方式,更合理的思路是用顺序读写替代随机读写,具体步骤如下:

步骤1:Python生成待更新数据的CSV

将计算结果批量导出为CSV文件,避免频繁数据库交互:

import csv

def generate_update_csv(output_path, batch_size):
    with open(output_path, 'w', newline='', encoding='utf-8') as f:
        writer = csv.writer(f)
        writer.writerow(['id', 'col1', 'col2'])  # 写入表头
        
        rows_batch = retrieve_raw_content_batch(batch_size)
        for rows in rows_batch:
            for row in rows:
                # 执行你的计算逻辑,得到id, col1, col2
                writer.writerow([id, col1, col2])

步骤2:PostgreSQL端导入数据并替换原表

利用COPY的高效批量写入能力,结合JOIN生成新表后替换原表:

-- 创建临时表存储更新数据(与原表id类型保持一致)
CREATE TEMP TABLE update_data (
    id BIGINT PRIMARY KEY,
    col1 [你的数据类型],
    col2 [你的数据类型]
);

-- COPY导入CSV数据(Windows路径需用双反斜杠或正斜杠)
COPY update_data(id, col1, col2) FROM 'C:\path\to\update_data.csv' WITH (FORMAT csv, HEADER);

-- 创建包含更新后数据的新表
CREATE TABLE new_table AS
SELECT 
    t.*,
    ud.col1,
    ud.col2
FROM original_table t
LEFT JOIN update_data ud ON t.id = ud.id;

-- 重建原表的约束与索引(根据实际需求添加)
ALTER TABLE new_table ADD PRIMARY KEY (id);
-- CREATE INDEX idx_new_table_col1 ON new_table(col1);

-- 替换原表(操作前务必备份原表)
DROP TABLE original_table;
ALTER TABLE new_table RENAME TO original_table;

这种方式的优势:COPY是顺序写入,能充分利用NVME的性能;创建新表的过程是顺序扫描+顺序写入,磁盘IO效率远高于随机更新,同时WAL日志生成量大幅降低。

2. 优化现有UPDATE方案(无法替换表时使用)

如果必须保留原表结构,可从以下几点优化:

  • 复用数据库连接:当前update_database每次新建连接,开销极大,改为在do_stuff中复用同一个连接
  • 增大批次大小:将batch_size提升至10万级(根据内存调整),减少事务提交次数
  • 批量事务提交:将多个批次的更新合并到一个事务中,避免频繁提交的开销

改进后的代码示例:

import psycopg2
from psycopg2 import extras

def retrieve_raw_content_batch(batch_size):
    with db_connect() as conn:
        with conn.cursor('testit') as cursor:
            cursor.execute("SELECT id, columnoftext FROM table;")
            while True:
                rows = cursor.fetchmany(batch_size)
                if not rows:
                    break
                yield rows

def do_stuff(row_batch, commit_batch_size):
    # 复用单个连接与事务
    with db_connect() as conn:
        with conn.cursor() as cursor:
            update_query = "UPDATE table SET col1 = %s, col2 = %s WHERE id = %s"
            upload_list = []
            total_processed = 0
            
            for rows in row_batch:
                for row in rows:
                    # 计算得到id, col1, col2
                    upload_list.append((col1, col2, id))  # 注意参数顺序与UPDATE语句匹配
                    total_processed += 1
                    
                    # 达到提交批次阈值时执行更新
                    if len(upload_list) >= commit_batch_size:
                        extras.execute_batch(cursor, update_query, upload_list)
                        upload_list.clear()
                    
                    # 每N次提交后手动提交事务,避免事务过大
                    if total_processed % (commit_batch_size * 10) == 0:
                        conn.commit()
            
            # 处理剩余未提交的数据
            if upload_list:
                extras.execute_batch(cursor, update_query, upload_list)
            conn.commit()

def main(batch_size):
    rows_batch = retrieve_raw_content_batch(batch_size)
    do_stuff(rows_batch, batch_size)

3. PostgreSQL配置优化

针对磁盘IO与WAL的关键配置调整:

  • WAL参数:
    • wal_buffers = 64MB:增大WAL缓冲区,减少磁盘写入频率
    • checkpoint_completion_target = 0.9:让检查点过程更平滑,降低IO峰值
    • min_wal_size = 2GB:配合已设置的max_wal_size=10GB,减少WAL文件的创建与删除开销
  • 同步策略:若数据无需强一致性,设置synchronous_commit = off,减少磁盘同步等待时间
  • 内存缓存:Windows下设置shared_buffers = 物理内存的1/4(如32GB内存设为8GB),让更多数据缓存到内存,减少磁盘读取
  • 预加载数据:执行SELECT pg_prewarm('original_table');,将原表数据提前加载到内存,减少后续随机读取

总结

对于16亿行的超大规模数据,COPY+替换表的方案能充分发挥NVME的顺序读写性能,处理速度会比直接UPDATE提升一个数量级。若必须使用UPDATE,则重点优化连接复用、批次大小与事务策略,同时配合PostgreSQL的WAL与缓存配置调整。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.19 00:22:51