Python Polars处理大CSV写入PostgreSQL时内存泄漏崩溃求助
问题描述
我编写了一个Python Polars脚本,通过LazyFrame扫描CSV文件,选择指定列、过滤特定行后将结果上传至PostgreSQL数据库。该脚本在处理小型测试CSV时运行正常,但处理真实业务场景下的大CSV时,Python进程会持续占用内存直至崩溃。已在配备64GB内存的Windows设备和32GB内存的macOS设备上测试过,可排除硬件不足问题。尝试过设置low_memory=True、使用collect(engine='streaming')、sink_csv等方法,均无法解决内存占用过高的问题。
原代码如下:
import polars as pl from sqlalchemy import create_engine result = ( pl.scan_csv("/path/to/file/201601.csv", separator="|", encoding='utf8-lossy', try_parse_dates=True ) # 选择指定列 .select( ["ProgramClassID", "NetworkAffiliationID", "LogServiceID", "LogEntryDate","StartTime", "EndTime","Duration", "ProgramTitle", "Subtitle"] ) # 过滤掉LogServiceID不等于3680的行 .filter( pl.col("LogServiceID") != 3680) # 按日期排序 .sort( "LogEntryDate") # 加载到内存 .collect() ) # 数据库连接信息 uri='postgresql://postgres:pwd@0.0.0.0:5432/crtc_program_logs' # 写入数据库 result.write_database(table_name="polars_test_02", connection=uri, if_table_exists="replace" )
问题根源
核心问题出在排序操作和全量内存加载:
- LazyFrame的
.sort()无法流式处理,必须将全量数据加载到内存才能完成排序,这是大文件内存溢出的主要原因。 .collect()会把整个数据集加载到内存,后续write_database再全量写入,进一步加剧内存压力。
修复方案
1. 移除/转移排序操作
如果业务允许,直接去掉.sort();如果最终需要有序数据,建议在PostgreSQL中通过ORDER BY查询,或创建索引——数据库层面的排序比Python内存排序高效得多。
2. 流式写入数据库,避免全量加载
不要调用.collect(),直接在LazyFrame上调用.write_database(),Polars会自动分批次处理数据,无需将全量数据存入内存。
3. 优化CSV读取参数
- 手动指定列的数据类型:通过
dtypes参数给数值列指定合适的类型(如pl.Int32),减少Polars自动推断的内存开销。 - 按需开启日期解析:如果
try_parse_dates导致额外内存占用,可手动指定日期格式,或先按字符串读取再转换。
优化后的代码
import polars as pl ( pl.scan_csv( "/path/to/file/201601.csv", separator="|", encoding='utf8-lossy', try_parse_dates=True, # 手动指定数据类型,降低内存占用 dtypes={ "ProgramClassID": pl.Int32, "NetworkAffiliationID": pl.Int32, "LogServiceID": pl.Int32, "Duration": pl.Int32 } ) .select([ "ProgramClassID", "NetworkAffiliationID", "LogServiceID", "LogEntryDate","StartTime", "EndTime","Duration", "ProgramTitle", "Subtitle" ]) .filter(pl.col("LogServiceID") != 3680) # 直接流式写入数据库,无需加载全量数据 .write_database( table_name="polars_test_02", connection='postgresql://postgres:pwd@0.0.0.0:5432/crtc_program_logs', if_table_exists="replace", # 控制批次大小,进一步降低内存峰值 batch_size=100_000 ) )
额外建议
- 若必须在写入前排序:可按
LogEntryDate范围拆分CSV,分别排序后写入数据库,最后在库内合并;或使用partition_by结合排序,但需严格控制单分区数据量。 - 监控内存开销:通过
pl.Config.set_tbl_rows(10)查看数据预览,或用profile()方法分析各步骤的内存占用。
内容的提问来源于stack exchange,提问作者queen_macaroni
相关产品推荐
相关产品推荐

