Python如何快速读取S3中10万+40KB小parquet文件用于数据处理
核心问题根源
你遇到的性能问题本质是海量极小parquet文件的开销浪费:每个文件仅存1行数据,单文件读取时的S3请求连接建立、parquet元数据解析的开销,是实际数据读取开销的几十上百倍,所有常规读文件逻辑都是逐个处理文件,天然会被S3的单请求延迟卡住,这是你现有方案耗时极高的核心原因。
优化方案
1. 代码侧优化(可将耗时压缩到1~2分钟)
你当前最快的Boto3方案是单线程同步请求,每秒最多处理3~5个文件,改成多线程批量异步拉取即可获得几十倍的性能提升:S3读取是纯IO密集型任务,用线程池批量发起请求可以完全压满网络带宽,避免单请求等待的空耗。
参考实现:
from concurrent.futures import ThreadPoolExecutor import boto3 import pandas as pd from io import BytesIO import logging # 线程数可根据实际情况调到100~300,S3默认支持单账号每秒数千请求 MAX_WORKERS = 200 s3_client = boto3.client('s3') MY_BUCKET = "你的桶名" def read_single_parquet(file_key: str): try: resp = s3_client.get_object(Bucket=MY_BUCKET, Key=file_key) return pd.read_parquet(BytesIO(resp["Body"].read())) except Exception: # 异常文件返回空DF过滤即可 return pd.DataFrame() def read_parquet_objects(self, objects_dict: dict) -> dict: df_holder = {} with ThreadPoolExecutor(max_workers=MAX_WORKERS) as executor: for group_name in objects_dict.keys(): logging.info(f"Start reading parquets for: {group_name}") # 批量提交所有文件的读取任务 futures = [executor.submit(read_single_parquet, file_key) for file_key in objects_dict[group_name]] # 收集所有返回结果 df_list = [f.result() for f in futures if not f.result().empty] # 合并时自动兼容不同schema df_holder[group_name] = pd.concat(df_list, axis=0, ignore_index=True, join="outer") return df_holder
如果还需要更快,可以换成aiobotocore的异步IO实现,性能还能再提升30%左右。
2. AWS原生服务方案(可满足<30秒的耗时要求)
要达到秒级读取,直接用Athena是性价比最高的方案,不需要改现有文件结构:
- 先用Glue Crawler自动扫描你的S3路径,生成按
group_N分区的外部表,自动开启parquet schema合并,一次配置永久生效。 - 直接用AWS Wrangler执行SQL查询拉取结果:
import awswrangler as wr # 直接查指定group的全量数据,返回合并好的pandas DataFrame df = wr.athena.read_sql_query( sql="SELECT * FROM 你的表名 WHERE group_N='group1'", database="你的Glue库名" )
Athena底层会自动批量合并小文件读取,不需要你自己处理IO、schema合并逻辑,即使是数万级的小文件,查询+拉取结果的总耗时也不会超过20秒,完全符合你的要求。
如果是定期运行的任务,可以额外加一个Lambda触发器:每次有新文件上传到S3时,自动触发Lambda将同group的小文件合并成一个大parquet存到另一个路径,后续读取大文件的耗时仅需毫秒级,一劳永逸解决小文件问题。
之前方案的问题说明
- Dask、PyArrow、原生AWS Wrangler慢的原因:这些库默认读取大量小文件时,会先逐个请求文件的元数据,每个文件要发2次S3请求,额外开销比你直接用Boto3拉取全量内容还大,所以速度更慢。
- PySpark报错原因:GC溢出是因为Driver端需要收集所有小文件的元数据,你给的512m Driver内存完全不够,连接池超时是因为默认S3连接池配置太小,同时发起的请求过多拿不到连接,即使调整配置跑通,性能也远不如Athena,没有优化必要。
内容的提问来源于stack exchange,提问作者E. Faslo
相关产品推荐
相关产品推荐

