如何读取S3存储桶中所有Parquet文件并合并为Pandas DataFrame
解决方法
问题拆解
你的代码存在两个核心问题:
- 仅读取指定文件夹下的文件,未递归遍历子文件夹
- 路径下存在0字节的无效Parquet文件,导致pyarrow解析失败
步骤1:递归获取所有有效Parquet文件
先用s3fs遍历所有子文件夹,筛选出大小大于0的有效.parquet文件:
import pyarrow.parquet as pq import s3fs import pandas as pd # 初始化S3文件系统 s3 = s3fs.S3FileSystem() # 递归遍历目标路径下所有.parquet文件,过滤空文件 bucket_path = 'vivienda-test/2022/11' parquet_files = [ f"s3://{path}" for path in s3.glob(f"{bucket_path}/**/*.parquet", recursive=True) if s3.info(path)['size'] > 0 ] # 检查是否存在有效文件 if not parquet_files: print("未找到任何有效的Parquet文件") exit()
步骤2:合并为单个Pandas DataFrame
推荐两种实现方式,根据文件数量选择:
方式一:用pyarrow批量读取(性能优先)
适合文件数量多、数据量大的场景:
# 创建数据集并读取 dataset = pq.ParquetDataset(parquet_files, filesystem=s3) merged_df = dataset.read_pandas().to_pandas() print(merged_df)
方式二:循环读取合并(简单直观)
适合文件数量少、调试阶段使用:
# 初始化空DataFrame merged_df = pd.DataFrame() # 逐个读取并合并 for file in parquet_files: temp_df = pq.read_table(file, filesystem=s3).to_pandas() merged_df = pd.concat([merged_df, temp_df], ignore_index=True) print(merged_df)
额外优化建议
- 处理损坏文件:添加异常捕获,避免单个损坏文件导致任务中断
merged_df = pd.DataFrame() for file in parquet_files: try: temp_df = pq.read_table(file, filesystem=s3).to_pandas() merged_df = pd.concat([merged_df, temp_df], ignore_index=True) except Exception as e: print(f"读取文件失败 {file}: {str(e)}")
- 超大数据量处理:用Dask并行读取,避免内存溢出
import dask.dataframe as dd ddf = dd.read_parquet(f"s3://{bucket_path}/**/*.parquet", filesystem=s3, ignore_metadata_file=True) merged_df = ddf.compute()
内容的提问来源于stack exchange,提问作者Alexander Bermúdez Castañeda
相关产品推荐
相关产品推荐

