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

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是性价比最高的方案,不需要改现有文件结构:

  1. 先用Glue Crawler自动扫描你的S3路径,生成按group_N分区的外部表,自动开启parquet schema合并,一次配置永久生效。
  2. 直接用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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.09.30 23:54:04