如何向SageMaker SKLearnProcessor传递依赖文件并在Pipeline中使用
解决方案:SKLearnProcessor传递依赖文件与安装第三方库
你可以直接通过SKLearnProcessor.run()方法的source_dir或者dependencies参数完成依赖文件传递,无需切换到FrameworkProcessor,因为SKLearnProcessor本身就是FrameworkProcessor的子类,原生支持依赖处理逻辑:
- 如果你的所有依赖文件和主脚本都在同一个目录下,直接指定
source_dir参数即可,SageMaker会自动打包上传整个目录,运行时主脚本的工作目录就是该目录,可以直接导入同目录下的其他Python文件,同时如果目录下存在requirements.txt,会自动执行安装:
# 假设文件结构: # ./src/ # ├─ preprocessing.py # ├─ file1.py # ├─ file2.py # └─ requirements.txt sklearn_processor.run( code='preprocessing.py', source_dir='./src', # 新增这行 inputs=[ProcessingInput( source=input_data, destination='/opt/ml/processing/input')], outputs=[ProcessingOutput(output_name='train_data', source='/opt/ml/processing/train'), ProcessingOutput(output_name='test_data', source='/opt/ml/processing/test')], arguments=['--train-test-split-ratio', '0.2'] )
- 如果依赖文件分散在不同路径,不想统一放到一个目录,可以直接用
dependencies参数传入依赖文件路径列表:
sklearn_processor.run( code='preprocessing.py', dependencies=['file1.py', 'file2.py', 'requirements.txt'], # 新增这行 inputs=[ProcessingInput( source=input_data, destination='/opt/ml/processing/input')], outputs=[ProcessingOutput(output_name='train_data', source='/opt/ml/processing/train'), ProcessingOutput(output_name='test_data', source='/opt/ml/processing/test')], arguments=['--train-test-split-ratio', '0.2'] )
问题1解答
ScriptProcessor、SKLearnProcessor以及所有框架专属的Processor(比如PyTorchProcessor、TensorFlowProcessor)都是FrameworkProcessor的子类,全部支持dependencies、source_dir、code三个参数,用法和FrameworkProcessor完全一致,不需要额外做适配。只有最底层的基础Processor类不支持这些参数,需要手动配置输入。
问题2解答
不需要SageMaker Project就可以创建Pipeline
SageMaker Pipeline是独立的服务组件,创建和运行都不需要绑定SageMaker Project,只要你的执行角色有对应的SageMaker权限就可以直接创建。
ProcessingStep集成到Pipeline的示例代码:
from sagemaker.workflow.steps import ProcessingStep from sagemaker.workflow.pipeline import Pipeline # 先定义好你需要用的SKLearnProcessor实例,和之前的定义一致 sklearn_processor = SKLearnProcessor( framework_version='0.20.0', role=role, instance_type='ml.m5.xlarge', instance_count=1 ) # 定义处理步骤 step_process = ProcessingStep( name="DataPreprocessing", processor=sklearn_processor, inputs=[ProcessingInput( source=input_data, destination='/opt/ml/processing/input' )], outputs=[ ProcessingOutput(output_name='train_data', source='/opt/ml/processing/train'), ProcessingOutput(output_name='test_data', source='/opt/ml/processing/test') ], code="./src/preprocessing.py", source_dir="./src", # 同样支持传入source_dir、dependencies参数 arguments=['--train-test-split-ratio', '0.2'] ) # 定义Pipeline pipeline = Pipeline( name="data-preprocessing-pipeline", steps=[step_process], role=role ) # 提交创建/更新Pipeline pipeline.upsert(role_arn=role) # 运行Pipeline execution = pipeline.start()
内容的提问来源于stack exchange,提问作者shaik moeed
相关产品推荐
相关产品推荐

