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

如何仅使用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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.09.24 14:15:06