You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.08.27 01:54:28