AWS技术问询:SageMaker Pipeline数据编排自动化及Step Functions串联Notebook
问题1:如何在AWS SageMaker Pipeline中编排并自动化数据迁移与数据转换?
核心实现思路
SageMaker Pipeline通过Pipeline Step串联不同任务,数据迁移可结合AWS原生服务(如S3、Lambda)完成,数据转换依托SageMaker Processing Job实现,全程通过Pipeline定义实现自动化调度与触发。
具体步骤
定义数据迁移任务
- 若从外部数据源(如RDS、跨账户S3)迁移,可通过
LambdaStep调用Lambda函数,利用AWS SDK完成数据同步至SageMaker指定S3路径。 - 示例Lambda核心逻辑(Python):
import boto3 s3 = boto3.client('s3') def lambda_handler(event, context): # 从源S3桶复制到SageMaker数据桶 s3.copy_object( Bucket='sagemaker-data-bucket-xxx', CopySource={'Bucket': 'source-data-bucket', 'Key': 'raw-data.csv'}, Key='raw/raw-data.csv' ) return {'status': 'success', 's3_path': 's3://sagemaker-data-bucket-xxx/raw/raw-data.csv'}
- 若从外部数据源(如RDS、跨账户S3)迁移,可通过
定义数据转换Processing Step
- 使用SageMaker Processing Job,选择Spark/Scikit-learn等框架编写转换脚本,在Pipeline中定义
ProcessingStep关联该任务。 - 示例Pipeline中Processing Step定义:
from sagemaker.processing import ScriptProcessor, ProcessingInput, ProcessingOutput from sagemaker.workflow.steps import ProcessingStep # 初始化Spark脚本处理器 script_processor = ScriptProcessor( image_uri='711395599931.dkr.ecr.us-east-1.amazonaws.com/sagemaker-spark-processing:2.4-cpu-py37-v1.0', role='SageMakerExecutionRole', instance_count=1, instance_type='ml.m5.xlarge' ) # 定义数据转换步骤 processing_step = ProcessingStep( name='Data-Transform-Step', processor=script_processor, inputs=[ProcessingInput(source='s3://sagemaker-data-bucket-xxx/raw/', destination='/opt/ml/processing/input')], outputs=[ProcessingOutput(source='/opt/ml/processing/output', destination='s3://sagemaker-data-bucket-xxx/transformed/')], code='s3://sagemaker-code-bucket/transform_script.py' )
- 使用SageMaker Processing Job,选择Spark/Scikit-learn等框架编写转换脚本,在Pipeline中定义
串联Pipeline并启用自动化
- 将迁移Step(LambdaStep)和转换Step(ProcessingStep)按顺序加入Pipeline,可通过CloudWatch Events设置定时触发,或绑定S3对象创建事件触发。
- 示例Pipeline定义:
from sagemaker.workflow.pipeline import Pipeline from sagemaker.workflow.lambda_step import LambdaStep, LambdaOutput # 定义Lambda输出 lambda_output = LambdaOutput(output_name='s3_path') lambda_step = LambdaStep( name='Data-Migrate-Step', lambda_func_arn='arn:aws:lambda:us-east-1:xxx:function:data-migrate-func', outputs=[lambda_output] ) # 构建并提交Pipeline pipeline = Pipeline( name='Data-Migrate-Transform-Pipeline', parameters=[], steps=[lambda_step, processing_step], sagemaker_session=sagemaker_session ) pipeline.upsert(role_arn='SageMakerExecutionRole')
问题2:将Azure Databricks Notebook迁移至AWS,用SageMaker+Step Functions按顺序执行三类Notebook
关键实现路径
将Databricks Notebook转换为可执行脚本,通过Step Functions状态机编排SageMaker Processing/Training任务,实现数据处理→特征工程→ML算法的顺序执行。
具体实现步骤
Notebook迁移与适配
- 将Databricks Notebook导出为
.py脚本,替换Azure专属依赖:把ADLS访问代码改为S3(用boto3或Spark S3 API),移除Databricks专属API调用,调整Spark配置适配SageMaker镜像。 - 示例代码适配:
# 原Databricks代码 # df = spark.read.csv("abfss://container@storageaccount.dfs.core.windows.net/raw.csv") # 改为AWS兼容代码 df = spark.read.csv("s3://sagemaker-data-bucket/raw.csv")
- 将Databricks Notebook导出为
Step Functions状态机编排
- 采用SageMaker Processing/Training Job作为执行单元,在Step Functions中定义串联任务:
状态机JSON示例(简化版):{ "Comment": "Databricks Notebook Migration Pipeline", "StartAt": "Data-Processing-Step", "States": { "Data-Processing-Step": { "Type": "Task", "Resource": "arn:aws:states:::sagemaker:createProcessingJob.sync", "Parameters": { "ProcessingJobName": "data-processing-job", "ProcessingInputs": [{ "InputName": "raw-data", "S3Input": { "S3Uri": "s3://sagemaker-data-bucket/raw/", "LocalPath": "/opt/ml/processing/input", "S3DataType": "S3Prefix", "S3InputMode": "File" } }], "ProcessingOutputConfig": { "Outputs": [{ "OutputName": "processed-data", "S3Output": { "S3Uri": "s3://sagemaker-data-bucket/processed/", "LocalPath": "/opt/ml/processing/output", "S3UploadMode": "EndOfJob" } }] }, "AppSpecification": { "ImageUri": "711395599931.dkr.ecr.us-east-1.amazonaws.com/sagemaker-spark-processing:2.4-cpu-py37-v1.0", "ContainerEntrypoint": ["spark-submit"], "ContainerArguments": ["s3://sagemaker-code-bucket/data_processing.py"] }, "RoleArn": "arn:aws:iam::xxx:role/SageMakerExecutionRole", "InstanceCount": 1, "InstanceType": "ml.m5.xlarge" }, "Next": "Feature-Engineering-Step" }, "Feature-Engineering-Step": { "Type": "Task", "Resource": "arn:aws:states:::sagemaker:createProcessingJob.sync", # 替换为特征工程脚本与输出路径 "Next": "ML-Algorithm-Step" }, "ML-Algorithm-Step": { "Type": "Task", "Resource": "arn:aws:states:::sagemaker:createTrainingJob.sync", # 配置SageMaker训练任务,关联ML算法脚本 "End": true } } }
- 采用SageMaker Processing/Training Job作为执行单元,在Step Functions中定义串联任务:
配套资源准备
- SageMaker Execution Role:需包含S3访问、Step Functions调用、SageMaker Processing/Training权限。
- 自定义镜像(可选):若有特殊依赖,构建并推送镜像至AWS ECR,供SageMaker任务调用。
- CloudWatch日志:配置状态机日志,跟踪每个步骤的执行状态与报错信息。
内容的提问来源于stack exchange,提问作者Vishnu
相关产品推荐
相关产品推荐

