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

