Tutorial AWS Step Functions untuk ML: Orchestrasi Workflow ML

# Tutorial Lengkap AWS Step Functions untuk ML: Orkestrasi ML Workflows AWS Step Functions menyediakan orkestrasi workflow serverless untuk pipeline machine learning. Layanan ini memungkinkan Anda me...

By Ruby Abdullah · · tutorial
AWSStep FunctionsWorkflowMLOpsOrchestrationAutomation

Tutorial Lengkap AWS Step Functions untuk ML: Orkestrasi ML Workflows

AWS Step Functions menyediakan orkestrasi workflow serverless untuk pipeline machine learning. Layanan ini memungkinkan Anda mengkoordinasi multiple layanan AWS, menangani error dengan baik, dan membangun workflow ML kompleks dengan monitoring visual.

Mengapa Step Functions untuk ML?

Manfaat Utama:
  • Workflow visual: Lihat eksekusi pipeline secara real-time
  • Error handling: Built-in retry dan error recovery
  • Integrasi layanan: Native AWS service connectors
  • Serverless: Tidak perlu mengelola infrastruktur
  • State management: Lacak state workflow secara otomatis

Use Cases:
  • ML training pipelines
  • Data preprocessing workflows
  • Automasi deployment model
  • Orkestrasi batch inference
  • Automasi MLOps

Prerequisites

pip install boto3 sagemaker

AWS CLI sudah dikonfigurasi

aws configure

Quick Start

1. Basic ML Workflow

{

"Comment": "Simple ML Training Pipeline",

"StartAt": "PreprocessData",

"States": {

"PreprocessData": {

"Type": "Task",

"Resource": "arn:aws:lambda:us-east-1:123456789:function:preprocess",

"Next": "TrainModel"

},

"TrainModel": {

"Type": "Task",

"Resource": "arn:aws:states:::sagemaker:createTrainingJob.sync",

"Parameters": {

"TrainingJobName.$": "States.Format('training-{}', $.Execution.Name)",

"AlgorithmSpecification": {

"TrainingImage": "123456789.dkr.ecr.us-east-1.amazonaws.com/xgboost:latest",

"TrainingInputMode": "File"

},

"RoleArn": "arn:aws:iam::123456789:role/SageMakerRole",

"InputDataConfig": [

{

"ChannelName": "train",

"DataSource": {

"S3DataSource": {

"S3DataType": "S3Prefix",

"S3Uri.$": "$.traindatauri"

}

}

}

],

"OutputDataConfig": {

"S3OutputPath": "s3://bucket/output"

},

"ResourceConfig": {

"InstanceCount": 1,

"InstanceType": "ml.m5.xlarge",

"VolumeSizeInGB": 30

},

"StoppingCondition": {

"MaxRuntimeInSeconds": 3600

}

},

"Next": "EvaluateModel"

},

"EvaluateModel": {

"Type": "Task",

"Resource": "arn:aws:lambda:us-east-1:123456789:function:evaluate",

"End": true

}

}

}

2. Deploy dengan CloudFormation

AWSTemplateFormatVersion: '2010-09-09'

Resources:

MLPipelineStateMachine:

Type: AWS::StepFunctions::StateMachine

Properties:

StateMachineName: ml-training-pipeline

RoleArn: !GetAtt StepFunctionsRole.Arn

DefinitionString: !Sub |

{

"StartAt": "PreprocessData",

"States": {

"PreprocessData": {

"Type": "Task",

"Resource": "${PreprocessFunction.Arn}",

"Next": "TrainModel"

},

"TrainModel": {

"Type": "Task",

"Resource": "arn:aws:states:::sagemaker:createTrainingJob.sync",

"Parameters": {

"TrainingJobName.$": "States.Format('job-{}', $.Execution.Name)"

},

"End": true

}

}

}

Integrasi SageMaker

1. Training Job

{

"TrainModel": {

"Type": "Task",

"Resource": "arn:aws:states:::sagemaker:createTrainingJob.sync",

"Parameters": {

"TrainingJobName.$": "States.Format('training-{}', $.Execution.Name)",

"AlgorithmSpecification": {

"TrainingImage": "683313688378.dkr.ecr.us-east-1.amazonaws.com/sagemaker-xgboost:1.5-1",

"TrainingInputMode": "File"

},

"RoleArn": "arn:aws:iam::123456789:role/SageMakerRole",

"InputDataConfig": [

{

"ChannelName": "train",

"DataSource": {

"S3DataSource": {

"S3DataType": "S3Prefix",

"S3Uri.$": "$.input.trainuri"

}

},

"ContentType": "text/csv"

}

],

"OutputDataConfig": {

"S3OutputPath.$": "$.input.outputuri"

},

"ResourceConfig": {

"InstanceCount": 1,

"InstanceType.$": "$.input.instancetype",

"VolumeSizeInGB": 50

},

"HyperParameters": {

"objective": "binary:logistic",

"numround": "100",

"maxdepth.$": "States.Format('{}', $.input.maxdepth)"

},

"StoppingCondition": {

"MaxRuntimeInSeconds": 7200

}

},

"ResultPath": "$.trainingresult",

"Next": "CreateModel"

}

}

2. Buat Model

{

"CreateModel": {

"Type": "Task",

"Resource": "arn:aws:states:::sagemaker:createModel",

"Parameters": {

"ModelName.$": "States.Format('model-{}', $.Execution.Name)",

"PrimaryContainer": {

"Image": "683313688378.dkr.ecr.us-east-1.amazonaws.com/sagemaker-xgboost:1.5-1",

"ModelDataUrl.$": "$.trainingresult.ModelArtifacts.S3ModelArtifacts"

},

"ExecutionRoleArn": "arn:aws:iam::123456789:role/SageMakerRole"

},

"ResultPath": "$.modelresult",

"Next": "CreateEndpointConfig"

}

}

3. Buat Endpoint

{

"CreateEndpointConfig": {

"Type": "Task",

"Resource": "arn:aws:states:::sagemaker:createEndpointConfig",

"Parameters": {

"EndpointConfigName.$": "States.Format('config-{}', $.Execution.Name)",

"ProductionVariants": [

{

"VariantName": "AllTraffic",

"ModelName.$": "$.modelresult.ModelArn",

"InitialInstanceCount": 1,

"InstanceType": "ml.m5.large"

}

]

},

"ResultPath": "$.endpointconfigresult",

"Next": "CreateEndpoint"

},

"CreateEndpoint": {

"Type": "Task",

"Resource": "arn:aws:states:::sagemaker:createEndpoint",

"Parameters": {

"EndpointName.$": "States.Format('endpoint-{}', $.Execution.Name)",

"EndpointConfigName.$": "$.endpointconfigresult.EndpointConfigArn"

},

"End": true

}

}

4. Batch Transform

{

"BatchTransform": {

"Type": "Task",

"Resource": "arn:aws:states:::sagemaker:createTransformJob.sync",

"Parameters": {

"TransformJobName.$": "States.Format('transform-{}', $.Execution.Name)",

"ModelName.$": "$.modelname",

"TransformInput": {

"DataSource": {

"S3DataSource": {

"S3DataType": "S3Prefix",

"S3Uri.$": "$.inputuri"

}

},

"ContentType": "text/csv",

"SplitType": "Line"

},

"TransformOutput": {

"S3OutputPath.$": "$.outputuri"

},

"TransformResources": {

"InstanceCount": 1,

"InstanceType": "ml.m5.large"

}

},

"End": true

}

}

Error Handling

1. Konfigurasi Retry

{

"TrainModel": {

"Type": "Task",

"Resource": "arn:aws:states:::sagemaker:createTrainingJob.sync",

"Retry": [

{

"ErrorEquals": ["SageMaker.AmazonSageMakerException"],

"IntervalSeconds": 60,

"MaxAttempts": 3,

"BackoffRate": 2.0

},

{

"ErrorEquals": ["States.Timeout"],

"IntervalSeconds": 30,

"MaxAttempts": 2,

"BackoffRate": 1.5

}

],

"Catch": [

{

"ErrorEquals": ["States.ALL"],

"ResultPath": "$.error",

"Next": "HandleTrainingError"

}

],

"Next": "EvaluateModel"

}

}

2. State Error Handler

{

"HandleTrainingError": {

"Type": "Task",

"Resource": "arn:aws:lambda:us-east-1:123456789:function:handle-error",

"Parameters": {

"error.$": "$.error",

"executionid.$": "$.Execution.Id",

"statename": "TrainModel"

},

"Next": "NotifyFailure"

},

"NotifyFailure": {

"Type": "Task",

"Resource": "arn:aws:states:::sns:publish",

"Parameters": {

"TopicArn": "arn:aws:sns:us-east-1:123456789:ml-alerts",

"Message.$": "States.Format('Training gagal: {}', $.error.Cause)"

},

"End": true

}

}

Logika Kondisional

1. Choice State

{

"EvaluateModel": {

"Type": "Task",

"Resource": "arn:aws:lambda:us-east-1:123456789:function:evaluate",

"ResultPath": "$.evaluation",

"Next": "CheckAccuracy"

},

"CheckAccuracy": {

"Type": "Choice",

"Choices": [

{

"Variable": "$.evaluation.accuracy",

"NumericGreaterThanEquals": 0.9,

"Next": "DeployModel"

},

{

"Variable": "$.evaluation.accuracy",

"NumericGreaterThanEquals": 0.8,

"Next": "ManualApproval"

}

],

"Default": "RetrainModel"

}

}

2. Manual Approval

{

"ManualApproval": {

"Type": "Task",

"Resource": "arn:aws:states:::sqs:sendMessage.waitForTaskToken",

"Parameters": {

"QueueUrl": "https://sqs.us-east-1.amazonaws.com/123456789/approval-queue",

"MessageBody": {

"executionid.$": "$.Execution.Id",

"modelmetrics.$": "$.evaluation",

"tasktoken.$": "$.Task.Token"

}

},

"Next": "DeployModel"

}

}

Parallel Processing

1. Parallel Branches

{

"ParallelProcessing": {

"Type": "Parallel",

"Branches": [

{

"StartAt": "TrainModelA",

"States": {

"TrainModelA": {

"Type": "Task",

"Resource": "arn:aws:states:::sagemaker:createTrainingJob.sync",

"Parameters": {

"TrainingJobName.$": "States.Format('model-a-{}', $.Execution.Name)"

},

"End": true

}

}

},

{

"StartAt": "TrainModelB",

"States": {

"TrainModelB": {

"Type": "Task",

"Resource": "arn:aws:states:::sagemaker:createTrainingJob.sync",

"Parameters": {

"TrainingJobName.$": "States.Format('model-b-{}', $.Execution.Name)"

},

"End": true

}

}

}

],

"Next": "SelectBestModel"

}

}

2. Map State untuk Batch

{

"ProcessDatasets": {

"Type": "Map",

"ItemsPath": "$.datasets",

"MaxConcurrency": 5,

"Iterator": {

"StartAt": "ProcessSingleDataset",

"States": {

"ProcessSingleDataset": {

"Type": "Task",

"Resource": "arn:aws:lambda:us-east-1:123456789:function:process-dataset",

"End": true

}

}

},

"ResultPath": "$.processeddatasets",

"Next": "MergeResults"

}

}

Pipeline Lengkap

{

"Comment": "Complete ML Pipeline dengan Step Functions",

"StartAt": "ValidateInput",

"States": {

"ValidateInput": {

"Type": "Task",

"Resource": "arn:aws:lambda:us-east-1:123456789:function:validate",

"Next": "PreprocessData"

},

"PreprocessData": {

"Type": "Task",

"Resource": "arn:aws:states:::glue:startJobRun.sync",

"Parameters": {

"JobName": "preprocess-job",

"Arguments": {

"--inputpath.$": "$.inputpath",

"--outputpath.$": "$.outputpath"

}

},

"Next": "TrainModel"

},

"TrainModel": {

"Type": "Task",

"Resource": "arn:aws:states:::sagemaker:createTrainingJob.sync",

"Parameters": {

"TrainingJobName.$": "States.Format('train-{}', $.Execution.Name)"

},

"Retry": [{"ErrorEquals": ["States.ALL"], "MaxAttempts": 2}],

"Catch": [{"ErrorEquals": ["States.ALL"], "Next": "NotifyFailure"}],

"Next": "EvaluateModel"

},

"EvaluateModel": {

"Type": "Task",

"Resource": "arn:aws:lambda:us-east-1:123456789:function:evaluate",

"Next": "CheckMetrics"

},

"CheckMetrics": {

"Type": "Choice",

"Choices": [

{

"Variable": "$.accuracy",

"NumericGreaterThanEquals": 0.9,

"Next": "RegisterModel"

}

],

"Default": "NotifyLowAccuracy"

},

"RegisterModel": {

"Type": "Task",

"Resource": "arn:aws:states:::sagemaker:createModel",

"Next": "DeployEndpoint"

},

"DeployEndpoint": {

"Type": "Task",

"Resource": "arn:aws:states:::sagemaker:createEndpoint",

"Next": "NotifySuccess"

},

"NotifySuccess": {

"Type": "Task",

"Resource": "arn:aws:states:::sns:publish",

"Parameters": {

"TopicArn": "arn:aws:sns:us-east-1:123456789:ml-notifications",

"Message": "Pipeline berhasil selesai"

},

"End": true

},

"NotifyFailure": {

"Type": "Task",

"Resource": "arn:aws:states:::sns:publish",

"Parameters": {

"TopicArn": "arn:aws:sns:us-east-1:123456789:ml-alerts",

"Message.$": "$.error"

},

"End": true

},

"NotifyLowAccuracy": {

"Type": "Task",

"Resource": "arn:aws:states:::sns:publish",

"Parameters": {

"Message": "Akurasi model di bawah threshold"

},

"End": true

}

}

}

Best Practices

1. Manajemen Input/Output

{

"TrainModel": {

"Type": "Task",

"Resource": "...",

"InputPath": "$.trainingconfig",

"ResultPath": "$.trainingresult",

"OutputPath": "$"

}

}

2. Eksekusi dengan SDK

import boto3

import json

sfn = boto3.client("stepfunctions")

Mulai eksekusi

response = sfn.startexecution(

stateMachineArn="arn:aws:states:us-east-1:123456789:stateMachine:ml-pipeline",

name="execution-001",

input=json.dumps({

"inputpath": "s3://bucket/input",

"outputpath": "s3://bucket/output"

})

)

executionarn = response["executionArn"]

Cek status

status = sfn.describeexecution(executionArn=execution_arn)

print(f"Status: {status['status']}")

Kesimpulan

Step Functions untuk ML menyediakan:

  • Orkestrasi visual: Monitor eksekusi pipeline
  • Error handling: Retry dan recovery otomatis
  • Paralelisme: Processing bersamaan
  • Integrasi: Native AWS service connectors
  • Serverless: Tidak perlu manajemen infrastruktur
  • Key takeaways:

    • Gunakan sync integrations untuk menunggu
    • Implementasikan error handling yang proper
    • Manfaatkan parallel states untuk efisiensi
    • Gunakan choice states untuk logika kondisional
    • Monitor eksekusi via console

    Artikel Terkait

    Tutorial AWS SageMaker Pipelines: ML Pipeline Automation

    Tutorial Lengkap AWS SageMaker Pipelines: Automasi ML Workflows SageMaker Pipelines adalah layanan CI/CD yang dibuat khu...

    Windmill: Ubah Script Python, TypeScript, Bash, dan SQL Jadi Workflow, API, dan UI

    Windmill: Ubah Script Python, TypeScript, Bash, dan SQL Jadi Workflow, API, dan UI Halo temen-temen, di tutorial kali in...

    Tutorial Temporal: Durable Execution untuk Workflow yang Andal

    Temporal dengan Python: Durable Execution untuk Workflow yang Andal Temporal adalah platform untuk durable execution: ia...

    ZenML: Framework Pipeline MLOps yang Modular dan Cloud-Agnostic

    ZenML: Framework Pipeline MLOps yang Modular dan Cloud-Agnostic Pendahuluan Membangun model machine learning yang akurat...