Python实现百万级PostgreSQL表与CSV增量更新及内存优化问题
百万级CSV同步PostgreSQL内存优化方案
核心优化原则:最小化内存占用,避免全量加载
1. 分块读取CSV,避免一次性加载全量数据
不要把整个CSV读进内存,按固定大小分块处理,处理完一块就释放内存:
- 用原生
csv模块(内存占用最低):
import csv from psycopg2.extras import execute_batch def process_batch(batch, db_records, cur): inserts = [] updates = [] for row in batch: acct_id = row['acct_id'] if acct_id not in db_records: inserts.append((row['acct_id'], row['col1'], row['col2'], row['update_time'])) elif row['update_time'] > db_records[acct_id]: updates.append((row['col1'], row['col2'], row['update_time'], row['acct_id'])) # 批量插入 if inserts: cur.execute(""" INSERT INTO real_acct (acct_id, col1, col2, update_time) VALUES %s """, (inserts,)) # 批量更新 if updates: execute_batch(cur, """ UPDATE real_acct SET col1 = %s, col2 = %s, update_time = %s WHERE acct_id = %s """, updates) # 主逻辑 conn = psycopg2.connect("dbname=xxx user=xxx password=xxx host=xxx") cur = conn.cursor() # 只拉取需要比对的主键和更新时间,避免全量加载表 cur.execute("SELECT acct_id, update_time FROM real_acct") db_records = {str(row[0]): row[1] for row in cur.fetchall()} # 分块读取CSV batch_size = 10000 with open('real_acct.csv', 'r') as f: reader = csv.DictReader(f) batch = [] for row in reader: batch.append(row) if len(batch) == batch_size: process_batch(batch, db_records, cur) batch = [] if batch: process_batch(batch, db_records, cur) conn.commit() cur.close() conn.close()
- 用Pandas(代码更简洁,适合需要数据清洗的场景):
import pandas as pd from psycopg2.extras import execute_batch conn = psycopg2.connect("dbname=xxx user=xxx password=xxx host=xxx") cur = conn.cursor() # 拉取比对用的主键和更新时间 cur.execute("SELECT acct_id, update_time FROM real_acct") db_records = {str(row[0]): row[1] for row in cur.fetchall()} # 指定紧凑数据类型,降低单块内存占用 dtypes = { 'acct_id': 'str', 'col1': 'int32', 'col2': 'float32', 'update_time': 'datetime64[ns]' } # 分块读取CSV for chunk in pd.read_csv('real_acct.csv', chunksize=10000, dtype=dtypes): # 区分新增和更新记录 new_rows = chunk[~chunk['acct_id'].isin(db_records.keys())] updated_rows = chunk[chunk['acct_id'].isin(db_records.keys())] updated_rows = updated_rows[updated_rows['update_time'] > updated_rows['acct_id'].map(db_records)] # 批量插入新增记录 if not new_rows.empty: cur.execute(""" INSERT INTO real_acct (acct_id, col1, col2, update_time) VALUES %s """, [tuple(row) for row in new_rows[['acct_id', 'col1', 'col2', 'update_time']].values]) # 批量更新记录 if not updated_rows.empty: update_data = [(row['col1'], row['col2'], row['update_time'], row['acct_id']) for _, row in updated_rows.iterrows()] execute_batch(cur, """ UPDATE real_acct SET col1 = %s, col2 = %s, update_time = %s WHERE acct_id = %s """, update_data) conn.commit() cur.close() conn.close()
2. 避免全量加载SQL表,只拉取必要字段
不要用pd.read_sql读取整张表,只查询主键和更新标识(比如update_time、版本号),构建映射字典用于比对,内存占用仅为全表的几分之一甚至几十分之一。
3. 用PostgreSQL原生批量操作替代单条SQL
- 用
execute_batch(psycopg2)或copy_expert实现批量写入,大幅减少SQL交互次数,同时降低内存开销。 - 最优方案:利用PostgreSQL临时表+
COPY+ON CONFLICT,把数据处理压力交给数据库,Python仅负责执行命令:
import psycopg2 conn = psycopg2.connect("dbname=xxx user=xxx password=xxx host=xxx") cur = conn.cursor() # 创建与目标表结构一致的临时表 cur.execute("CREATE TEMP TABLE temp_real_acct (LIKE real_acct INCLUDING ALL)") # 直接把CSV导入临时表 with open('real_acct.csv', 'r') as f: cur.copy_expert("COPY temp_real_acct FROM STDIN WITH CSV HEADER", f) # 一次性完成新增和更新 cur.execute(""" INSERT INTO real_acct SELECT * FROM temp_real_acct ON CONFLICT (acct_id) DO UPDATE SET col1 = EXCLUDED.col1, col2 = EXCLUDED.col2, update_time = EXCLUDED.update_time """) conn.commit() cur.execute("DROP TABLE temp_real_acct") cur.close() conn.close()
这种方式内存占用几乎为0,速度是Python处理的数倍,适合超大规模数据同步。
4. 优化数据类型
读取CSV时指定紧凑的数据类型(比如用int32替代int64,float32替代float64,category替代普通字符串),减少单块数据的内存占用。
内容的提问来源于stack exchange,提问作者Hermain Qadir
相关产品推荐
相关产品推荐

