如何用Pandas与Concurrent.Futures优化S3 Parquet文件读取效率
并行读取S3中Parquet文件的问题解决与优化
解决ProcessPoolExecutor启动错误
你遇到的启动错误核心原因是Python多进程在Windows类环境下需要明确的程序入口,必须用if __name__ == '__main__':包裹启动多进程的主逻辑,否则子进程会重复执行模块顶层代码,导致启动异常。
修正后的script1示例:
import boto3 import pandas as pd from concurrent.futures import ProcessPoolExecutor import io # 封装单个文件读取函数 def read_s3_parquet(s3_path): s3 = boto3.client('s3') bucket_name, key = s3_path.split('/', 2)[1], s3_path.split('/', 2)[2] try: response = s3.get_object(Bucket=bucket_name, Key=key) return pd.read_parquet(io.BytesIO(response['Body'].read())) except Exception as e: print(f"读取文件 {s3_path} 失败: {str(e)}") return pd.DataFrame() # 返回空DataFrame避免后续合并报错 if __name__ == '__main__': # 替换为你的1000个Parquet文件路径列表 file_paths = ['s3://your-bucket/path/file1.parquet', ...] # 8核虚拟机建议设为6-7个进程,留资源给系统 with ProcessPoolExecutor(max_workers=7) as executor: dfs = list(executor.map(read_s3_parquet, file_paths)) # 合并所有读取结果 combined_df = pd.concat(dfs, ignore_index=True) # 调用script2处理数据 import script2 script2.process_data(combined_df)
与script2的正确集成方式
- 保持职责分离:script1专注并行读取与数据合并,将完整DataFrame传递给script2的处理函数
- 若无需全量数据,可边读边处理:在并行回调中调用script2的单文件处理逻辑,减少内存占用
# script1中修改为边读边处理 def process_single_file(s3_path): df = read_s3_parquet(s3_path) if not df.empty: import script2 script2.process_single_df(df) # script2需实现单文件处理逻辑 if __name__ == '__main__': with ProcessPoolExecutor(max_workers=7) as executor: executor.map(process_single_file, file_paths)
读取错误的处理策略
- 捕获具体异常:在
read_s3_parquet中针对s3.exceptions.NoSuchKey(文件不存在)、PermissionError(权限不足)、pd.errors.EmptyDataError(空文件)等做针对性处理 - 记录错误日志:将失败的文件路径和错误信息写入日志文件,方便后续排查
- 跳过错误文件:返回空DataFrame或标记错误,避免整个并行任务中断
虚拟机配置的方案适用性
你的虚拟机有8个虚拟处理器、32GB内存,ProcessPoolExecutor是完全合适的方案:
- 多进程可充分利用多核CPU,Parquet解析属于CPU密集型操作,并行读取能显著提升速度
- 32GB内存足够支撑1000个Parquet文件的合并(只要单个文件体积不是过大),若内存紧张可改用分块处理或Dask框架
- 设置
max_workers为6-7,避免占用全部CPU资源导致系统卡顿
内容的提问来源于stack exchange,提问作者sergioMoreno
相关产品推荐
相关产品推荐

