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

如何读取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)

额外优化建议

  1. 处理损坏文件:添加异常捕获,避免单个损坏文件导致任务中断
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)}")
  1. 超大数据量处理:用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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.14 08:55:18