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

如何将S3中的多个parquet文件读取合并为pandas DataFrame并解决相关报错

算力说明

150个Parquet文件如果总大小在10G以内,普通8核16G内存的机器完全可以加载为单个DataFrame;总大小超过20G的话,建议分块处理,或使用Dask、PySpark等分布式计算框架降低单机器内存压力。


方案1:修复AWS Wrangler使用问题

你之前加了chunked=True参数才会返回生成器,不需要分块的话直接去掉该参数即可直接返回DataFrame:

import awswrangler as wr
import pandas as pd

# 直接加载为单个DataFrame
df = wr.s3.read_parquet(
    path="s3://my-s3-data/folder1/subfolder1/subfolder2/",
    dataset=True,
    columns=df_cols
)

如果需要分块避免内存不足,按以下方式迭代生成器拼接即可:

chunks = []
# chunked可指定数值,代表每个分块的最大行数,比如chunked=100000
for chunk in wr.s3.read_parquet(
    path="s3://my-s3-data/folder1/subfolder1/subfolder2/",
    dataset=True,
    columns=df_cols,
    chunked=True
):
    chunks.append(chunk)
df = pd.concat(chunks, ignore_index=True)

使用awswrangler需要提前配置好AWS凭证,可通过aws configure命令本地配置,或运行时传入boto3_session参数指定自定义会话。


方案2:修复PyArrow+s3fs报错

你遇到的AioClientCreator相关报错是boto3、s3fs、aiobotocore版本不兼容导致的,先执行版本统一安装:

pip install -U boto3 s3fs aiobotocore pyarrow

调整代码如下即可正常读取:

import s3fs
import pyarrow.parquet as pq

fs = s3fs.S3FileSystem()
# 路径明确匹配parquet后缀,避免读取到无关文件
path = "s3://my-bucket/folder1/subfolder1/subfolder2/*.parquet"
df = pq.ParquetDataset(path, filesystem=fs).read().to_pandas()

方案3:修复自定义批量读取函数问题

你遇到的找不到文件报错,是因为S3的Prefix参数不需要带s3://bucket前缀,同时你缺失了单文件读取的pd_read_s3_parquet实现,调整后的完整代码如下:

import boto3
import pandas as pd
import io

# 补充单Parquet文件读取逻辑
def pd_read_s3_parquet(key, bucket, s3_client, **args):
    obj = s3_client.get_object(Bucket=bucket, Key=key)
    return pd.read_parquet(io.BytesIO(obj['Body'].read()), **args)

def pd_read_s3_multiple_parquets(prefix, bucket, verbose=False, **args):
    if not prefix.endswith('/'):
        prefix = prefix + '/'
    s3_client = boto3.client('s3')
    s3 = boto3.resource('s3')
    # 前缀直接传文件夹相对路径,不需要带bucket和s3前缀
    s3_keys = [item.key for item in s3.Bucket(bucket).objects.filter(Prefix=prefix)
               if item.key.endswith('.parquet')]
    if not s3_keys:
        raise ValueError(f'No parquet found in bucket {bucket}, prefix {prefix}')
    if verbose:
        print(f'Found {len(s3_keys)} parquet files:')
        for p in s3_keys: 
            print(p)
    dfs = [pd_read_s3_parquet(key, bucket=bucket, s3_client=s3_client, **args) 
           for key in s3_keys]
    return pd.concat(dfs, ignore_index=True)

# 调用示例:文件在s3://my-bucket/folder1/subfolder1/subfolder2/下时按下方传参
df = pd_read_s3_multiple_parquets(prefix='folder1/subfolder1/subfolder2/', bucket='my_bucket')

内容的提问来源于stack exchange,提问作者paige the beginner

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.09.25 13:24:03