AWS Step Functions for ML Tutorial: ML Workflow Orchestration

# 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

Complete AWS Step Functions for ML Tutorial: Orchestrating ML Workflows

AWS Step Functions provides serverless workflow orchestration for machine learning pipelines. It enables you to coordinate multiple AWS services, handle errors gracefully, and build complex ML workflows with visual monitoring.

Why Step Functions for ML?

Key Benefits:
  • Visual workflows: See pipeline execution in real-time
  • Error handling: Built-in retry and error recovery
  • Service integration: Native AWS service connectors
  • Serverless: No infrastructure to manage
  • State management: Track workflow state automatically

Use Cases:
  • ML training pipelines
  • Data preprocessing workflows
  • Model deployment automation
  • Batch inference orchestration
  • MLOps automation

Prerequisites

pip install boto3 sagemaker

AWS CLI configured

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

}

}

}

SageMaker Integration

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. Create 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. Create 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. Retry Configuration

{

"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. Error Handler State

{

"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 failed: {}', $.error.Cause)"

},

"End": true

}

}

Conditional Logic

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 for 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"

}

}

Complete Pipeline

{

"Comment": "Complete ML Pipeline with 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 completed successfully"

},

"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": "Model accuracy below threshold"

},

"End": true

}

}

}

Best Practices

1. Input/Output Management

{

"TrainModel": {

"Type": "Task",

"Resource": "...",

"InputPath": "$.trainingconfig",

"ResultPath": "$.trainingresult",

"OutputPath": "$"

}

}

2. Execution with SDK

import boto3

import json

sfn = boto3.client("stepfunctions")

Start execution

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"]

Check status

status = sfn.describeexecution(executionArn=execution_arn)

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

Conclusion

Step Functions for ML provides:

  • Visual orchestration: Monitor pipeline execution
  • Error handling: Automatic retries and recovery
  • Parallelism: Concurrent processing
  • Integration: Native AWS service connectors
  • Serverless: No infrastructure management
  • Key takeaways:

    • Use sync integrations for waiting
    • Implement proper error handling
    • Leverage parallel states for efficiency
    • Use choice states for conditional logic
    • Monitor executions via console

    Related Articles

    AWS SageMaker Pipelines Tutorial: ML Pipeline Automation

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

    Windmill: Turn Python, TypeScript, Bash, and SQL Scripts into Workflows, APIs, and UIs

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

    Temporal Tutorial: Durable Execution for Reliable Workflows

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

    ZenML: Modular and Cloud-Agnostic MLOps Pipeline Framework

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