如何用64个Multiprocessing Worker并行处理S3中大量对象?
解决方案:并行处理S3存储桶中的所有对象
你的现有代码存在两个核心问题:
- 每个工作进程都会完整遍历一遍S3对象列表,仅处理其中第
i个对象,造成大量重复IO和资源浪费 - 当对象数量超过64时,只有前64个对象会被处理,剩余对象完全被忽略
要实现遍历所有对象并充分利用64个进程,正确的思路是先一次性获取所有目标对象的列表,再将对象作为任务分发给进程池并行处理,每个进程负责处理单个或多个对象,进程池会自动做负载均衡。
修改后的代码示例
import multiprocessing import pandas as pd import boto3 def process_single_s3_obj(obj, some_list, dest_bucket, key, secret): """处理单个S3对象:读取、过滤、写入目标桶""" try: # 读取S3对象 response = obj.get() df = pd.read_csv(response['Body']) # 数据处理逻辑 df = df[df['some_col'].isin(some_list)] # 生成目标文件路径(可根据原文件名调整命名规则) dest_key = f"some_other_folder/{obj.key.split('/')[-1]}" df.to_csv(f"s3://{dest_bucket}/{dest_key}", index=False, storage_options={'key': key, 'secret': secret}) except Exception as e: print(f"处理对象 {obj.key} 失败: {str(e)}") if __name__ == '__main__': # 初始化S3客户端 s3 = boto3.resource('s3') bucket = s3.Bucket('your_source_bucket_name') dest_bucket = 'some_bucket_name' key = 'your_aws_access_key' secret = 'your_aws_secret_key' some_list = ['target_value1', 'target_value2'] # 替换为你的过滤列表 # 一次性获取所有需要处理的S3对象(自动处理分页) target_objs = list(bucket.objects.filter(Prefix='folder_name/')) # 创建64进程的进程池 with multiprocessing.Pool(64) as a_pool: # 使用starmap传递多参数任务:每个对象对应一个处理任务 a_pool.starmap(process_single_s3_obj, [(obj, some_list, dest_bucket, key, secret) for obj in target_objs])
关键优化点
- 避免重复遍历:仅在主进程中遍历一次S3对象列表,所有进程共享这个列表,减少S3 API调用次数
- 任务粒度合理:以单个S3对象为任务单位,进程池会自动将所有对象分配给64个进程,实现负载均衡
- 异常处理:添加异常捕获,避免单个对象处理失败导致整个进程崩溃
- 命名规则灵活:根据原对象的Key生成目标路径,避免原代码中
file_{i}.csv的命名冲突(当多个进程处理多个对象时会覆盖文件)
额外建议
- 尽量不要硬编码AWS密钥,优先使用环境变量或AWS配置文件(
~/.aws/credentials),这样更安全且便于部署 - 如果对象数量极大(十万级以上),可以考虑使用
imap_unordered替代starmap,这样能在任务完成时立即获取结果,避免占用过多内存存储所有任务参数 - 对于超大CSV文件,可以考虑分块读取处理,避免内存溢出
内容的提问来源于stack exchange,提问作者qwerty
相关产品推荐
相关产品推荐

