Polars拆分排序LazyFrame为小分块 保留完整id组避免内存溢出
Polars按id组拆分大文件实现方案
核心逻辑:由于数据集已按id排序、且每个id固定对应连续13条记录,无需使用内存开销更高的groupby操作,直接按固定行偏移做惰性切片即可,全程不会将全量数据加载到内存,同时自动处理最后一个分块行数不足的边界场景。
最小可复现示例实现
示例数据集共20组id、合计260行,要求拆分后前2个文件各包含8组id(104行),最后1个文件包含剩余4组id(52行),实现代码如下:
import polars as pl import numpy as np # 构造示例数据集 df = pl.DataFrame({"id":13*list(range(20)), "value":np.arange(13*20)}).sort(by="id") ldf = df.lazy() # 拆分参数配置 rows_per_id_group = 13 # 每个id对应的固定行数 ids_per_chunk = 8 # 单个分块包含的id组数量,实际生产环境替换为5000即可 chunk_row_count = rows_per_id_group * ids_per_chunk total_rows = ldf.select(pl.len()).collect().item() # 惰性获取总行数,无额外内存开销 # 循环生成所有分块 for chunk_idx, offset in enumerate(range(0, total_rows, chunk_row_count)): # 惰性切片获取当前分块:当剩余行数不足chunk_row_count时,Polars自动截断到文件末尾 current_chunk = ldf.slice(offset, chunk_row_count) # 替换为实际的文件写出逻辑,示例中先打印分块行数做校验 chunk_rows = current_chunk.select(pl.len()).collect().item() print(f"分块{chunk_idx+1} 行数:{chunk_rows}") # 生产环境写出示例:current_chunk.write_parquet(f"dataset_part_{chunk_idx}.parquet")
运行后输出符合预期:
分块1 行数:104 分块2 行数:104 分块3 行数:52
GB级大文件生产环境注意事项
- 不要先加载全量数据到内存再转LazyFrame,直接使用
pl.scan_parquet()/pl.scan_csv()从源文件创建惰性计算框架,从源头控制内存占用。 - 按测算的单文件100MB的目标,直接将
ids_per_chunk参数设置为5000即可,其余逻辑不需要调整。 - 写出文件优先选择带压缩的Parquet列存格式,读写效率是CSV的数倍,且磁盘占用更低。
- 如需校验分块完整性,可在写出前增加校验逻辑:
assert current_chunk.select(pl.col("id").n_unique() <= ids_per_chunk).collect().item(),无报错即代表没有出现id组被拆分到不同文件的问题。
方案选型说明
不推荐使用groupby实现拆分:对GB级数据集做groupby操作时,即使是惰性模式,Polars也需要维护全量id的分组索引,内存开销远高于固定偏移切片方案。在数据已按id排序、每个id对应行数固定的前提下,直接切片是性能最高、内存占用最低的实现方式。
内容的提问来源于stack exchange,提问作者TomNorway
相关产品推荐
相关产品推荐

