如何在SageMaker SKLearnProcessing作业中并行化独立循环任务
在Amazon SageMaker Processing中并行化SKLearn作业的方案
针对你200次独立迭代的任务,有两种核心并行思路可以最大化利用SageMaker的算力:单实例内多进程并行、多实例分布式并行,两种可以结合使用。
一、单实例内多进程并行(利用单实例的多核CPU)
ml.m5.4xlarge实例有16个vCPU,你可以直接在script.py中用Python的并行库(比如concurrent.futures.ProcessPoolExecutor)把任务分散到多个进程中执行,无需修改SageMaker启动代码。
修改后的script.py示例:
from concurrent.futures import ProcessPoolExecutor import os def run_function(i): # 保留你原有的run_function逻辑 # 注意:每个任务的输出文件要命名唯一,比如用i作为后缀,避免互相覆盖 # 示例:保存结果到/opt/ml/processing/output/result_{i}.txt with open(f'/opt/ml/processing/output/result_{i}.txt', 'w') as f: f.write(f'Task {i} completed') if __name__ == "__main__": # 假设你的任务列表是200个元素,这里用range(200)代替 task_list = list(range(200)) # 根据实例CPU核心数设置进程数,预留2个核心给系统进程 cpu_count = int(os.environ.get('CPU_COUNT', 16)) max_workers = cpu_count - 2 # 启动进程池执行任务 with ProcessPoolExecutor(max_workers=max_workers) as executor: executor.map(run_function, task_list)
二、多实例分布式并行(横向扩展实例数)
如果单实例的多核还不够,你可以增加SageMaker Processing的实例数量,把任务拆分到多个实例上执行,每个实例处理一部分任务,同时每个实例内部再用多进程加速。
1. 修改Notebook中的启动代码
调整instance_count参数,比如设置为4(可根据需求调整,最多支持数十个实例):
sklearn_processor = SKLearnProcessor( framework_version="1.0-1", role=role, instance_type="ml.m5.4xlarge", instance_count=4, # 增加实例数 sagemaker_session = Session() ) out_path = 's3://' + os.path.join(bucket, prefix,'outpath') sklearn_processor.run( code="script.py", outputs=[ ProcessingOutput(output_name="load_training_data", source = '/opt/ml/processing/output', # 修正原代码的语法错误:多了一个} destination = out_path), ], arguments=["--some-args", "args"] )
2. 修改script.py实现任务分片
SageMaker会给每个实例注入环境变量SM_HOSTS(所有实例的列表)和SM_CURRENT_HOST(当前实例名称),用这些变量来拆分任务:
from concurrent.futures import ProcessPoolExecutor import os import json def run_function(i): # 保留你原有的run_function逻辑 with open(f'/opt/ml/processing/output/result_{i}.txt', 'w') as f: f.write(f'Task {i} completed on host {os.environ["SM_CURRENT_HOST"]}') if __name__ == "__main__": # 获取集群实例信息 hosts = json.loads(os.environ['SM_HOSTS']) current_host = os.environ['SM_CURRENT_HOST'] host_index = hosts.index(current_host) total_instances = len(hosts) # 生成200个任务的列表 task_list = list(range(200)) # 按实例数分片:每个实例处理索引对总实例数取模等于自身索引的任务 assigned_tasks = [task for idx, task in enumerate(task_list) if idx % total_instances == host_index] # 单实例内用多进程加速 cpu_count = int(os.environ.get('CPU_COUNT', 16)) max_workers = cpu_count - 2 with ProcessPoolExecutor(max_workers=max_workers) as executor: executor.map(run_function, assigned_tasks)
关键注意事项
- 输出文件唯一性:每个任务的输出必须用唯一文件名(比如任务ID作为后缀),否则多个进程/实例的输出会互相覆盖。
- 资源监控:如果
run_function占用大量内存,要适当减少进程数,避免实例内存耗尽。 - 成本控制:实例数越多算力越强,但成本也越高,建议根据任务总耗时需求调整实例数量。
内容的提问来源于stack exchange,提问作者SKSKSKSK
相关产品推荐
相关产品推荐

