如何仅使用fastparquet实现分块写入parquet文件?
分块写入Parquet实现方案
- 你当前的分块逻辑是合理的,
np.array_split可以将全量DataFrame均匀切分为指定数量的子块,适合全量数据可一次性加载到内存的场景。如果全量数据过大无法整表读入内存,可以在读源数据阶段就通过pandas.read_csv/pandas.read_sql的chunksize参数直接获取分块迭代器,减少内存占用。 - 完整的循环写入逻辑需要配合fastparquet的
append参数实现,第一次写入时关闭追加模式,后续分块开启追加模式即可,所有分块的表结构(列名、字段类型)必须完全一致,否则会出现写入报错。
示例代码如下:import numpy as np from fastparquet import write n_chunks = 10 merged_dataframes = np.array_split(df, n_chunks) for idx, chunk_df in enumerate(merged_dataframes): # 首块覆盖写入,后续块追加写入 append_flag = False if idx == 0 else True write( "./local_output.parquet", chunk_df, append=append_flag, compression="snappy", index=False ) - 如果需要生成多分片Parquet数据集而非单个Parquet文件,无需使用追加模式,直接为每个分块生成独立的分片文件即可,上传到S3后可直接作为数据集被读取:
import os from fastparquet import write tmp_parts_dir = "./parquet_parts" os.makedirs(tmp_parts_dir, exist_ok=True) n_chunks = 10 merged_dataframes = np.array_split(df, n_chunks) for idx, chunk_df in enumerate(merged_dataframes): part_file_path = f"{tmp_parts_dir}/part-{str(idx).zfill(5)}.snappy.parquet" write(part_file_path, chunk_df, compression="snappy", index=False) - 写入完成后可直接通过boto3的上传接口将本地文件/目录传输到S3,不需要依赖被限制的AWS Wrangler、S3FS等工具。
注意:如果分块过程中存在字段类型不一致的情况,需要在写入前对每个分块做统一的类型校准,避免schema不匹配导致写入失败。
内容的提问来源于stack exchange,提问作者caasswa
相关产品推荐
相关产品推荐

