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

AWS技术问询:SageMaker Pipeline数据编排自动化及Step Functions串联Notebook

问题1:如何在AWS SageMaker Pipeline中编排并自动化数据迁移与数据转换?

核心实现思路

SageMaker Pipeline通过Pipeline Step串联不同任务,数据迁移可结合AWS原生服务(如S3、Lambda)完成,数据转换依托SageMaker Processing Job实现,全程通过Pipeline定义实现自动化调度与触发。

具体步骤

  1. 定义数据迁移任务

    • 若从外部数据源(如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'}
      
  2. 定义数据转换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'
      )
      
  3. 串联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算法的顺序执行。

具体实现步骤

  1. 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")
      
  2. 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
          }
        }
      }
      
  3. 配套资源准备

    • SageMaker Execution Role:需包含S3访问、Step Functions调用、SageMaker Processing/Training权限。
    • 自定义镜像(可选):若有特殊依赖,构建并推送镜像至AWS ECR,供SageMaker任务调用。
    • CloudWatch日志:配置状态机日志,跟踪每个步骤的执行状态与报错信息。

内容的提问来源于stack exchange,提问作者Vishnu

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.07 17:25:37