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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.01 17:05:21