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

如何读取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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.09.27 18:45:05