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

如何用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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.24 23:17:22