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
- 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:
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