Python读取PostgreSQL千万行数据写CSV内存溢出优化方案
内存占用过高的核心原因
你当前代码的核心问题是创建的是psycopg2默认的客户端侧缓冲游标,这种游标不管你怎么设置itersize参数,都会在执行查询后把整个结果集的所有数据全部拉取到Python进程的内存中缓存。1000万行数据哪怕单条长度很短,叠加Python对象的内存开销,很容易占到数GB内存,直接触发容器OOM被杀。
你设置的itersize=20000参数仅对服务端命名游标生效,普通游标下这个参数不会起任何作用。
可落地的优化方案
按优化效果从高到低排序:
- 优先用PostgreSQL原生COPY命令做流式导出:这是性能最好、内存占用最低的方案,不需要Python层逐行加载、转换数据,数据库会直接把查询结果以CSV格式字节流的方式写入你打开的文件,全程内存占用稳定在MB级别,导出速度比逐行fetch快5~10倍。核心是用psycopg2游标自带的
copy_expert方法,直接对接文件指针,不会缓存全量结果。 - 如果必须在Python层处理每一行数据,改用服务端命名游标:创建游标时传入自定义
name参数,psycopg2就会创建服务端游标,此时你设置的itersize才会生效,驱动会按照你指定的批次大小分批从数据库拉取数据,每次仅在内存里保留当前批次的行,不会全量缓存结果。注意使用服务端游标时不要把连接设为autocommit模式,保持默认的事务状态即可。 - 去掉不必要的内存开销:写文件时指定明确的编码,不要在迭代结果的过程中持有多余的行引用,不需要处理数据的话不要做多余的行格式转换,进一步压缩内存占用。
- 额外说明:你原来的异常处理代码有语法bug,
print('Unable to connect!\n{0}').format(e)的format方法写到了print括号外面,运行时会抛错,调整为f-string或者把format放到print括号内即可。
参考改后代码
最优方案:COPY流式导出(无逐行处理需求时用)
import os import sys import json import psycopg2 def create_server_connection(): DB_CONNECTION_PARAMS = os.environ["DB_REPLICA_CONNECTION"] json_object = json.loads(DB_CONNECTION_PARAMS) try: conn = psycopg2.connect( database=json_object["PGDATABASE"], user=json_object["PGUSER"], password=json_object["PGPASSWORD"], host=json_object["PGHOST"], port=json_object["PGPORT"] ) except psycopg2.OperationalError as e: print(f'Unable to connect!\n{e}') sys.exit(1) return conn if __name__ == "__main__": with create_server_connection() as connection: with connection.cursor() as cursor, open(file_name, 'w', newline='', encoding='utf-8') as fp: # 如果你要从外部sql文件读查询,就替换下面的查询逻辑 # query = open(sql_file_path, 'r', encoding='utf-8').read().strip(';') # export_sql = f"COPY ({query}) TO STDOUT WITH (FORMAT CSV, HEADER)" export_sql = "COPY (SELECT * FROM events) TO STDOUT WITH (FORMAT CSV, HEADER)" cursor.copy_expert(export_sql, fp) connection.commit()
备选方案:服务端游标分批拉取(需要逐行处理数据时用)
import os import sys import json import csv import psycopg2 # create_server_connection方法和上面一致,不需要改动 if __name__ == "__main__": with create_server_connection() as connection: # 传入name参数创建服务端游标,此时itersize才生效 with connection.cursor(name="events_export_cursor") as cursor: cursor.itersize = 20000 query = open(sql_file_path, mode='r', encoding='utf-8').read() cursor.execute(query) with open(file_name, 'w', newline='', encoding='utf-8') as fp: csv_writer = csv.writer(fp) for row in cursor: # 在这里加你需要的行处理逻辑 csv_writer.writerow(row) connection.commit()
用上面任意一种方案改完,整个导出流程的内存占用会稳定在几十MB级别,不需要给AWS Batch容器配大内存就能稳定跑完全量1000万行的导出。
内容的提问来源于stack exchange,提问作者Karry Bee
相关产品推荐
相关产品推荐

