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

如何使用awswrangler将S3多个Parquet文件合并写入单个文件

解决方案

原有代码生成大量小文件的原因是:dataset=True模式下,每次循环调用wr.s3.to_parquet的append操作都会生成独立的新Parquet文件,循环次数和输出文件数量成正比。

方案1:直接写入S3单个文件(无需本地磁盘存储全量文件,推荐)

这个方案用pyarrow的ParquetWriter配合S3文件系统逐块写入,全程仅保留单个分块在内存,不会出现OOM,最终输出S3单个文件:

import awswrangler as wr
import pyarrow as pa
import pyarrow.parquet as pq
from botocore.session import Session

try:
    # 初始化S3文件系统
    s3 = pa.fs.S3FileSystem(region=Session().get_config_variable('region'))
    writer = None
    
    # 分块读取S3上的小Parquet文件
    dfs = wr.s3.read_parquet(
        path=input_folder, 
        path_suffix=['.parquet'], 
        chunked=True, 
        use_threads=True,
        output_format='pyarrow' # 直接输出pyarrow表,减少转换开销
    )

    for idx, table in enumerate(dfs):
        # 第一次读取时拿到schema初始化writer,指定最终输出的S3路径
        if idx == 0:
            writer = pq.ParquetWriter(
                where=f"{target_path.lstrip('s3://')}", # pyarrow的S3路径不需要s3://前缀
                schema=table.schema,
                filesystem=s3,
                compression='snappy' # 按需调整压缩算法
            )
        # 逐块写入
        writer.write_table(table)
        logger.info(f"已写入第{idx+1}块数据")
    
    # 关闭writer完成写入,此时S3上就生成了单个Parquet文件
    writer.close()
    logger.info(f"单个Parquet文件已写入完成,路径:{target_path}")

except Exception as e:
    if writer is not None:
        writer.close()
    logger.error(e, exc_info=True)

方案2:本地中转后上传(适合本地磁盘可容纳最终单个文件的场景)

如果本地磁盘空间足够,也可以先把所有分块写入本地单个Parquet文件,再上传到S3,代码更简单:

import awswrangler as wr
import os

local_temp_path = "./temp_merged.parquet"

try:
    dfs = wr.s3.read_parquet(
        path=input_folder, 
        path_suffix=['.parquet'], 
        chunked=True, 
        use_threads=True
    )
    first_chunk = True
    for df in dfs:
        # 逐块追加写入本地单个Parquet文件
        df.to_parquet(local_temp_path, engine='pyarrow', append=not first_chunk)
        first_chunk = False
        logger.info("已追加写入1块数据到本地临时文件")
    
    # 上传本地文件到S3
    wr.s3.upload(local_file=local_temp_path, path=target_path)
    logger.info(f"单个Parquet文件已上传完成,路径:{target_path}")
    # 可选:删除本地临时文件
    os.remove(local_temp_path)

except Exception as e:
    logger.error(e, exc_info=True)
    # 异常时清理本地临时文件
    if os.path.exists(local_temp_path):
        os.remove(local_temp_path)

注意事项

  • Parquet单个文件建议控制在1~2GB以内,最大不要超过5GB,否则后续读取该文件时也容易出现OOM问题。
  • 两种方案都保留了chunked=True参数,不会加载全量数据到内存,符合内存限制要求。

内容的提问来源于stack exchange,提问作者Nirav Nagda

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.10.03 16:57:03