如何使用Polars分小批次读取Parquet大文件?BatchedParquetReader似乎无法正常工作
如何使用Polars分小批次读取Parquet大文件?BatchedParquetReader似乎无法正常工作
嘿,我完全懂你遇到的困扰——几十GB的Parquet文件硬往内存里塞肯定行不通,而且试过BatchedParquetReader却没达到预期的分批效果,这确实挺让人头疼的。下面我来给你梳理几个可行的解决办法,以及可能的问题原因:
一、优先用Polars懒加载API:scan_parquet + fetch_next_batch
这是Polars处理大文件最推荐的方式,因为scan_parquet不会一次性把整个文件读进内存,而是先构建一个查询计划,之后你可以按需分批获取数据,内存占用会被牢牢控制在单个批次的大小范围内。
举个具体的代码例子:
import polars as pl # 先创建一个懒加载的扫描器,不会立即读取数据 scanner = pl.scan_parquet("your_large_file.parquet") # 设置你想要的每个批次行数(根据你的内存情况调整,比如10万行) batch_size = 100_000 # 循环读取批次直到读完 while True: batch = scanner.fetch_next_batch(batch_size) if batch.is_empty(): break # 在这里处理你的批次数据,比如做清洗、统计或者写入数据库 print(f"处理了 {len(batch)} 行数据") # process_your_data(batch)
二、检查BatchedParquetReader的用法是否正确
如果你坚持要用BatchedParquetReader,可能是之前的用法有误导致它没真正分批。正确的打开方式应该是用iter_batches方法来迭代获取批次,而不是一次性读取:
from polars.io.parquet import BatchedParquetReader with BatchedParquetReader("your_large_file.parquet") as reader: # 按指定批次大小迭代读取 for batch in reader.iter_batches(batch_size=100_000): print(f"当前批次行数: {len(batch)}") # 处理当前批次
如果这样还是内存占用过高,你可以试试这两个优化:
- 调小
batch_size的值,比如降到5万行甚至更低 - 只读取你需要的列,通过
columns参数指定,避免加载无关数据:with BatchedParquetReader("your_large_file.parquet", columns=["col_a", "col_b"]) as reader: for batch in reader.iter_batches(batch_size=50_000): # 处理批次
三、额外提醒:Parquet文件的列存储特性
Parquet是列存储格式,如果你的文件里有特别大的列(比如超长文本列),单个批次的内存占用可能还是会很高。这种情况下,除了调小批次,还可以考虑对大列做分块处理,或者在写入Parquet的时候就设置合理的分块大小,方便后续读取。
总的来说,只要利用好Polars的懒加载机制,或者正确使用BatchedParquetReader的迭代方法,就能实现真正的分批读取,避免把整个大文件塞进内存。
备注:内容来源于stack exchange,提问作者Chris
相关产品推荐
相关产品推荐

