如何使用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
相关产品推荐
相关产品推荐

