如何读取Parquet文件并实现分块处理生成目标数据集
大体积Parquet文件分块读取处理方案
以下是三种无需转换为CSV即可实现分块读取处理的可行方案:
方案1:使用pandas原生分块读取能力
pandas 1.3.0及以上版本的read_parquet已原生支持chunksize参数,可直接返回分块迭代器,实现成本最低:
import pandas as pd # 按指定行数拆分读取,chunksize可根据可用内存调整 chunk_iterator = pd.read_parquet("your_large_file.parquet", chunksize=100000) processed_list = [] for chunk in chunk_iterator: # 此处写入单块数据的处理逻辑,比如过滤、统计、字段加工等 processed_chunk = chunk[chunk["amount"] > 0].copy() processed_list.append(processed_chunk) # 合并所有处理后的分块得到最终结果 final_df = pd.concat(processed_list, ignore_index=True)
方案2:基于PyArrow数据集接口实现灵活分块
该方案性能更优,同时支持分区Parquet读取、提前过滤行列等进阶需求,适合复杂场景:
import pyarrow.dataset as ds import pandas as pd # 加载Parquet数据集,自动识别分区结构 dataset = ds.dataset("your_large_file.parquet", format="parquet") processed_list = [] # 按自定义粒度分块读取 for batch in dataset.to_batches(batch_size=100000): chunk = batch.to_pandas() # 单块处理逻辑 processed_chunk = chunk.drop(columns=["unused_column"]) processed_list.append(processed_chunk) final_df = pd.concat(processed_list, ignore_index=True)
方案3:使用fastparquet库按行组分块
如果你默认使用fastparquet作为Parquet解析引擎,可直接用其行组迭代能力:
from fastparquet import ParquetFile import pandas as pd pf = ParquetFile("your_large_file.parquet") processed_list = [] # 按Parquet原生行组迭代读取 for row_group in pf.iter_row_groups(): chunk = row_group.to_pandas() # 单块处理逻辑 processed_chunk = chunk[chunk["category"] == "target_type"] processed_list.append(processed_chunk) final_df = pd.concat(processed_list, ignore_index=True)
注意事项
- 若最终不需要全量内存级DataFrame,可边处理分块边写入新的Parquet文件,进一步降低内存占用
chunksize/batch_size可根据当前可用内存调整,建议单块内存占用不超过可用内存的1/10- 单块处理时尽量避免保留不必要的中间对象,及时释放无用内存
内容的提问来源于stack exchange,提问作者Droid-Bird
相关产品推荐
相关产品推荐

