如何让Amazon SageMaker ProcessingJob按S3文件夹分片输入?
解决SageMaker ProcessingJob按文件夹分片处理的问题
因为SageMaker默认的ShardByS3Key是按单个文件分片,FullyReplicated会把所有数据复制到每个实例,都没法满足你按文件夹(每个文件夹里的三个文件必须一起处理)分片的需求,这里给你几个可行的实现思路:
方案1:预生成文件夹列表,自定义实例分片逻辑
第一步:获取所有目标文件夹列表
先写一个简单的脚本(可以用本地环境或者一次性的小ProcessingJob),用boto3遍历S3桶,提取所有job*格式的文件夹前缀。比如:import boto3 s3 = boto3.client('s3') bucket_name = "your-bucket-name" prefix = "" # 如果文件夹在根目录就留空,有父目录就填对应路径 # 列出所有文件夹前缀 folders = set() paginator = s3.get_paginator('list_objects_v2') for page in paginator.paginate(Bucket=bucket_name, Prefix=prefix, Delimiter='/'): for common_prefix in page.get('CommonPrefixes', []): folder = common_prefix['Prefix'] folders.add(folder) # 把文件夹列表保存到S3,比如存成folders.txt,每行一个文件夹路径 folder_list = list(folders) with open('/tmp/folders.txt', 'w') as f: f.write('\n'.join(folder_list)) s3.upload_file('/tmp/folders.txt', bucket_name, 'processing/folders.txt')第二步:主处理脚本实现分片
在你的ProcessingJob脚本里,通过SageMaker的环境变量获取实例信息(总实例数、当前实例索引),然后从S3下载folders.txt,拆分列表给当前实例处理:import os import boto3 s3 = boto3.client('s3') bucket_name = "your-bucket-name" # 下载文件夹列表 s3.download_file(bucket_name, 'processing/folders.txt', '/tmp/folders.txt') with open('/tmp/folders.txt', 'r') as f: all_folders = [line.strip() for line in f if line.strip()] # 获取实例信息 num_instances = int(os.environ['SM_NUM_HOSTS']) current_instance = os.environ['SM_CURRENT_HOST'] host_list = os.environ['SM_HOSTS'].split(',') instance_index = host_list.index(current_instance) # 按实例数均分文件夹列表 assigned_folders = all_folders[instance_index::num_instances] # 处理每个文件夹里的三个文件 for folder in assigned_folders: # 下载a.jpg、b.json、c.proto s3.download_file(bucket_name, f"{folder}a.jpg", '/tmp/a.jpg') s3.download_file(bucket_name, f"{folder}b.json", '/tmp/b.json') s3.download_file(bucket_name, f"{folder}c.proto", '/tmp/c.proto') # 这里写你的处理逻辑,比如读取三个文件做计算 process_files('/tmp/a.jpg', '/tmp/b.json', '/tmp/c.proto') # 处理完后上传结果到S3 s3.upload_file('/tmp/result.json', bucket_name, f"{folder}result.json")配置ProcessingJob
提交ProcessingJob时,把folders.txt作为输入之一,同时设置多实例参数(比如InstanceCount=3),脚本会自动按实例分配文件夹。
方案2:用SageMaker Pipeline实现单文件夹单任务并行
如果你的文件夹数量不是特别多(比如几千个以内),可以把每个文件夹作为一个独立的ProcessingJob,用SageMaker Pipeline的并行执行能力来处理:
- 先获取所有文件夹列表,然后生成对应的ProcessingJob定义
- 用
Parallel步骤把这些任务并行启动,每个任务只处理一个文件夹的三个文件 - 这种方式不需要手动分片,SageMaker会自动调度资源,缺点是任务数量过多时可能会遇到配额限制
方案3:自定义输入分片器(进阶)
你可以自定义SageMaker的输入分片逻辑,通过实现一个自定义的输入处理器,不过这个复杂度较高,适合有一定SageMaker底层开发经验的场景。核心思路是重写分片逻辑,把同一个文件夹下的文件归为同一个分片,然后分配给实例。
内容的提问来源于stack exchange,提问作者Vova Anisimov
相关产品推荐
相关产品推荐

