Polars读取Parquet指定列内存暴涨?求原理与文档指引
问题分析与解决方案
为什么读取600MB的3列会耗尽64GB内存?
你遇到的核心问题和Parquet的row group机制以及Polars流式处理的工作逻辑直接相关:
- Parquet的最小读取单元是row group:Parquet文件按row group存储数据,即便你只选择3列,Polars也必须加载整个row group中的这3列数据到内存。如果你的Parquet文件row group设置过大(比如单个row group中这3列的数据就超过几十GB),哪怕开启流式处理,也会因为需要加载完整row group而耗尽内存。
- 流式处理的局限性:
collect(streaming=True)确实会分批处理数据,但批次大小由Parquet的row group大小决定——row group是Parquet的原子读取单元,无法拆分单个row group进行流式处理。
Polars读取Parquet的核心机制
- 懒执行计划:
scan_parquet不会立即加载数据,而是创建一个LazyFrame,记录后续所有操作(比如select)形成执行计划,直到调用collect才会实际执行数据读取和处理。 - 列裁剪优化:Polars会利用Parquet文件的元数据,精准定位到指定列的存储位置,只会读取这些列的数据,不会加载其他列——你的代码逻辑没问题,问题不在列裁剪本身。
- 流式处理逻辑:开启
streaming=True后,Polars会把数据分成多个批次处理,每个批次对应一个row group,处理完一个批次后释放内存再处理下一个。但如果单个row group的大小超过可用内存,这个机制就会失效。
解决办法
- 检查并调整Parquet的row group大小
- 用PyArrow查看现有文件的row group信息:
import pyarrow.parquet as pq parquet_file = pq.ParquetFile("/home/ubuntu/parquet_files/your_file.parquet") # 遍历查看每个row group的行数和目标列大小 for i in range(parquet_file.num_row_groups): rg_meta = parquet_file.metadata.row_group(i) print(f"Row Group {i} 行数: {rg_meta.num_rows}") for col in rg_meta.columns: if col.path_in_schema in ["L_SHIPDATE", "L_LINESTATUS", "L_RETURNFLAG"]: print(f" 列 {col.path_in_schema} 大小: {col.total_bytes / 1024 / 1024:.2f} MB") - 如果row group过大,重新写入Parquet时设置更小的row group(比如10万行):
# 用Polars重新写入拆分row group df = pl.scan_parquet(directory).collect(streaming=True) df.write_parquet("new_parquet_files/", row_group_size=100_000)
- 用PyArrow查看现有文件的row group信息:
- 验证执行计划
- 用
explain()查看Polars的执行计划,确认确实只读取了指定列:
输出的df.select([ pl.col("L_SHIPDATE"), pl.col("L_LINESTATUS"), pl.col("L_RETURNFLAG") ]).explain()ParquetScan节点中,selected_columns应该显示你指定的3列。
- 用
- 调整Polars内存参数
- 可以设置Polars的内存上限,强制它更严格地控制内存使用:
pl.Config.set_memory_limit("50GB")
- 可以设置Polars的内存上限,强制它更严格地控制内存使用:
关键文档指引
- 懒执行机制:Polars的核心设计是
LazyFrame,所有文件读取和数据处理操作都会延迟到collect阶段,这样Polars可以进行全局优化(比如列裁剪、过滤下推)。 - Parquet读取参数:
scan_parquet支持row_groups(指定读取特定row group)、n_rows(限制读取行数)、low_memory(强制低内存模式)等参数,可帮助你精准控制内存使用。 - 流式处理:
collect(streaming=True)适用于处理超内存数据集,但前提是单个row group的大小不超过可用内存,否则需要先拆分row group。
内容的提问来源于stack exchange,提问作者Bhaskar Dabhi
相关产品推荐
相关产品推荐

