内存有限下用Polars高效分块读取大CSV并入库
纯Polars高效处理方案
针对你的场景,核心是利用Polars的Lazy API流式处理+预定义数据类型+批量写入SQLite,避免全量加载数据到内存,同时解决关联小表的需求。以下是具体方案:
一、核心优化思路
- 放弃
low_memory=True:Polars默认的列式读取内存效率远高于行式读取,开启low_memory会强制行式解析,反而导致内存过载。 - 预定义CSV字段类型:避免Polars反复扫描数据推断类型,减少内存占用和耗时。
- Lazy模式下关联小表:小DataFrame(3.6万行)可以直接广播关联,Lazy操作不会立即加载全量数据。
- 利用
sink_sqlite分批写入:无需手动分块,Polars会自动流式写入SQLite,控制内存使用。
二、具体代码实现
1. 读取字段转换小表
import polars as pl # 读取用于字段转换的小表 mapping_df = pl.read_database( query="SELECT * FROM your_mapping_table", connection_uri="sqlite:///your_source_db.db" ) mapping_lf = mapping_df.lazy()
2. 流式处理大CSV并写入SQLite
# 扫描大CSV,预定义dtypes避免类型推断开销 large_lf = pl.scan_csv( "large_data.csv", # 根据实际字段定义类型,示例: dtypes={ "user_id": pl.Int64, "event_time": pl.Datetime, "event_type": pl.String, "value": pl.Float64 }, infer_schema_length=1000, # 减少扫描行数推断类型 parse_dates=["event_time"], # 直接指定日期列,避免额外解析 skip_rows_with_error=True # 跳过坏行,避免处理中断 ) # 关联转换表(Lazy模式下仅生成执行计划,不加载数据) processed_lf = large_lf.join( mapping_lf, left_on="event_type", # 替换为你的实际关联键 right_on="map_type", how="left" ) # 批量写入SQLite,自动流式处理 processed_lf.sink_sqlite( destination="target_db.db", table_name="large_data_table", if_exists="replace", batch_size=100_000 # 可根据内存调整,建议50万以内 )
3. 手动分块降级方案(若自动sink仍有问题)
如果sink_sqlite内存压力仍大,可手动分块处理:
# 获取总行数(近似值,快速计算) total_rows = large_lf.select(pl.count()).collect().item() batch_size = 100_000 for offset in range(0, total_rows, batch_size): # 切片获取批次数据(Lazy模式下仅定位,不加载) batch_lf = large_lf.slice(offset, batch_size).join(mapping_lf, left_on="event_type", right_on="map_type", how="left") # 追加写入SQLite batch_lf.sink_sqlite( "target_db.db", table_name="large_data_table", if_exists="append", batch_size=batch_size )
三、你之前遇到问题的原因
- scan_csv+low_memory死机:
low_memory=True强制行式读取,破坏了Polars列式存储的内存优势,导致内存爆炸。 - read_csv(n_rows+skip)速度慢:Eager模式下每次读取都要重新扫描文件到指定跳过位置,IO开销极大;Lazy模式的
slice直接在扫描时定位,无需重复读取文件头。 - read_csv_batched/Lazy切片性能差:未预定义dtypes导致每次分块都要重新推断类型,额外消耗CPU和时间;加上批次大小不合理,频繁IO导致性能下降。
内容的提问来源于stack exchange,提问作者tgrandje
相关产品推荐
相关产品推荐

