Complete AWS SageMaker Pipelines Tutorial: Automating ML Workflows
SageMaker Pipelines is a purpose-built CI/CD service for machine learning that helps you automate and manage ML workflows. It enables you to create reproducible, production-ready ML pipelines with minimal code.
Why SageMaker Pipelines?
Key Benefits:- Automation: Automate end-to-end ML workflows
- Reproducibility: Track and reproduce experiments
- Integration: Native SageMaker service integration
- Visualization: DAG visualization in Studio
- Version control: Pipeline versioning and lineage
- Processing Steps
- Training Steps
- Transform Steps
- Model Steps
- Condition Steps
- Callback Steps
Prerequisites
pip install sagemaker boto3 pandas scikit-learn
Ensure SageMaker SDK >= 2.0
python -c "import sagemaker; print(sagemaker.version)"
Quick Start
1. Setup
import boto3
import sagemaker
from sagemaker.workflow.pipeline import Pipeline
from sagemaker.workflow.steps import ProcessingStep, TrainingStep
from sagemaker.workflow.parameters import ParameterString, ParameterInteger
session = sagemaker.Session()
bucket = session.defaultbucket()
role = sagemaker.getexecutionrole()
region = session.botoregionname
pipelinename = "iris-ml-pipeline"
2. Define Parameters
from sagemaker.workflow.parameters import (
ParameterString,
ParameterInteger,
ParameterFloat
)
Pipeline parameters
inputdata = ParameterString(
name="InputData",
defaultvalue=f"s3://{bucket}/iris/raw/data.csv"
)
traininginstancetype = ParameterString(
name="TrainingInstanceType",
defaultvalue="ml.m5.xlarge"
)
traininginstancecount = ParameterInteger(
name="TrainingInstanceCount",
defaultvalue=1
)
modelapprovalstatus = ParameterString(
name="ModelApprovalStatus",
defaultvalue="PendingManualApproval"
)
Processing Steps
1. Data Preprocessing
# preprocess.py
import argparse
import os
import pandas as pd
from sklearn.modelselection import traintestsplit
from sklearn.preprocessing import StandardScaler
if name == "main":
parser = argparse.ArgumentParser()
parser.addargument("--input-data", type=str)
parser.addargument("--test-size", type=float, default=0.2)
args = parser.parseargs()
# Read data
inputpath = os.path.join("/opt/ml/processing/input", "data.csv")
df = pd.readcsv(inputpath)
# Split features and target
X = df.drop("target", axis=1)
y = df["target"]
# Scale features
scaler = StandardScaler()
Xscaled = scaler.fittransform(X)
# Split data
Xtrain, Xtest, ytrain, ytest = traintestsplit(
Xscaled, y, testsize=args.testsize, randomstate=42
)
# Save outputs
traindf = pd.DataFrame(Xtrain)
traindf["target"] = ytrain.values
traindf.tocsv("/opt/ml/processing/train/train.csv", index=False, header=False)
testdf = pd.DataFrame(Xtest)
testdf["target"] = ytest.values
testdf.tocsv("/opt/ml/processing/test/test.csv", index=False, header=False)
print(f"Train size: {len(traindf)}, Test size: {len(testdf)}")
from sagemaker.processing import ProcessingInput, ProcessingOutput
from sagemaker.sklearn.processing import SKLearnProcessor
Create processor
sklearnprocessor = SKLearnProcessor(
frameworkversion="1.0-1",
role=role,
instancetype="ml.m5.large",
instancecount=1,
sagemakersession=session
)
Define processing step
stepprocess = ProcessingStep(
name="PreprocessData",
processor=sklearnprocessor,
inputs=[
ProcessingInput(
source=inputdata,
destination="/opt/ml/processing/input"
)
],
outputs=[
ProcessingOutput(
outputname="train",
source="/opt/ml/processing/train",
destination=f"s3://{bucket}/iris/processed/train"
),
ProcessingOutput(
outputname="test",
source="/opt/ml/processing/test",
destination=f"s3://{bucket}/iris/processed/test"
)
],
code="preprocess.py"
)
Training Steps
1. XGBoost Training Step
from sagemaker.estimator import Estimator
from sagemaker.inputs import TrainingInput
from sagemaker.workflow.steps import TrainingStep
Get XGBoost container
xgboostimage = sagemaker.imageuris.retrieve(
framework="xgboost",
region=region,
version="1.5-1"
)
Create estimator
xgbestimator = Estimator(
imageuri=xgboostimage,
role=role,
instancecount=traininginstancecount,
instancetype=traininginstancetype,
outputpath=f"s3://{bucket}/iris/models",
sagemakersession=session,
hyperparameters={
"objective": "multi:softmax",
"numclass": 3,
"numround": 100,
"maxdepth": 5,
"eta": 0.2
}
)
Define training step
steptrain = TrainingStep(
name="TrainModel",
estimator=xgbestimator,
inputs={
"train": TrainingInput(
s3data=stepprocess.properties.ProcessingOutputConfig.Outputs["train"].S3Output.S3Uri,
contenttype="text/csv"
),
"validation": TrainingInput(
s3data=stepprocess.properties.ProcessingOutputConfig.Outputs["test"].S3Output.S3Uri,
contenttype="text/csv"
)
}
)
2. Custom Training Step
from sagemaker.pytorch import PyTorch
pytorchestimator = PyTorch(
entrypoint="train.py",
sourcedir="src",
role=role,
instancecount=1,
instancetype="ml.p3.2xlarge",
frameworkversion="1.13",
pyversion="py39",
hyperparameters={
"epochs": 10,
"batch-size": 32
}
)
steptrainpytorch = TrainingStep(
name="TrainPyTorchModel",
estimator=pytorchestimator,
inputs={
"train": TrainingInput(
s3data=stepprocess.properties.ProcessingOutputConfig.Outputs["train"].S3Output.S3Uri
)
}
)
Evaluation Steps
# evaluate.py
import argparse
import json
import os
import tarfile
import pandas as pd
import xgboost as xgb
from sklearn.metrics import accuracyscore, f1score, classificationreport
if name == "main":
parser = argparse.ArgumentParser()
args = parser.parseargs()
# Load model
modelpath = "/opt/ml/processing/model/model.tar.gz"
with tarfile.open(modelpath) as tar:
tar.extractall(path="/opt/ml/processing/model")
model = xgb.Booster()
model.loadmodel("/opt/ml/processing/model/xgboost-model")
# Load test data
testpath = "/opt/ml/processing/test/test.csv"
testdf = pd.readcsv(testpath, header=None)
ytest = testdf.iloc[:, -1]
Xtest = testdf.iloc[:, :-1]
# Predict
dtest = xgb.DMatrix(Xtest)
predictions = model.predict(dtest)
# Calculate metrics
accuracy = accuracyscore(ytest, predictions)
f1 = f1score(ytest, predictions, average="weighted")
# Save evaluation report
report = {
"classificationmetrics": {
"accuracy": {"value": accuracy},
"f1score": {"value": f1}
}
}
outputdir = "/opt/ml/processing/evaluation"
os.makedirs(outputdir, existok=True)
with open(os.path.join(outputdir, "evaluation.json"), "w") as f:
json.dump(report, f)
print(f"Accuracy: {accuracy:.4f}, F1: {f1:.4f}")
from sagemaker.workflow.properties import PropertyFile
Create evaluation processor
evaluationprocessor = SKLearnProcessor(
frameworkversion="1.0-1",
role=role,
instancetype="ml.m5.large",
instancecount=1,
sagemakersession=session
)
Property file for condition evaluation
evaluationreport = PropertyFile(
name="EvaluationReport",
outputname="evaluation",
path="evaluation.json"
)
stepevaluate = ProcessingStep(
name="EvaluateModel",
processor=evaluationprocessor,
inputs=[
ProcessingInput(
source=steptrain.properties.ModelArtifacts.S3ModelArtifacts,
destination="/opt/ml/processing/model"
),
ProcessingInput(
source=stepprocess.properties.ProcessingOutputConfig.Outputs["test"].S3Output.S3Uri,
destination="/opt/ml/processing/test"
)
],
outputs=[
ProcessingOutput(
outputname="evaluation",
source="/opt/ml/processing/evaluation",
destination=f"s3://{bucket}/iris/evaluation"
)
],
code="evaluate.py",
propertyfiles=[evaluationreport]
)
Condition Steps
from sagemaker.workflow.conditions import ConditionGreaterThanOrEqualTo
from sagemaker.workflow.conditionstep import ConditionStep
from sagemaker.workflow.functions import JsonGet
Define condition
conditionaccuracy = ConditionGreaterThanOrEqualTo(
left=JsonGet(
stepname=stepevaluate.name,
propertyfile=evaluationreport,
jsonpath="classificationmetrics.accuracy.value"
),
right=0.9 # Minimum accuracy threshold
)
Condition step
stepcondition = ConditionStep(
name="CheckAccuracy",
conditions=[conditionaccuracy],
ifsteps=[stepregister], # If condition is true
elsesteps=[stepfail] # If condition is false
)
Model Registration
from sagemaker.workflow.modelstep import ModelStep
from sagemaker.model import Model
from sagemaker.workflow.step
collections import RegisterModel
Create model
model = Model(
imageuri=xgboostimage,
modeldata=steptrain.properties.ModelArtifacts.S3ModelArtifacts,
role=role,
sagemakersession=session
)
Register model step
stepregister = RegisterModel(
name="RegisterModel",
estimator=xgbestimator,
modeldata=steptrain.properties.ModelArtifacts.S3ModelArtifacts,
contenttypes=["text/csv"],
responsetypes=["text/csv"],
inferenceinstances=["ml.m5.large", "ml.m5.xlarge"],
transforminstances=["ml.m5.large"],
modelpackagegroupname="iris-model-package-group",
approvalstatus=modelapprovalstatus
)
Complete Pipeline
from sagemaker.workflow.pipeline import Pipeline
from sagemaker.workflow.failstep import FailStep
Fail step for low accuracy
stepfail = FailStep(
name="ModelAccuracyFail",
errormessage="Model accuracy is below threshold"
)
Create pipeline
pipeline = Pipeline(
name=pipelinename,
parameters=[
inputdata,
traininginstancetype,
traininginstancecount,
modelapprovalstatus
],
steps=[
stepprocess,
steptrain,
stepevaluate,
stepcondition
],
sagemakersession=session
)
Create/update pipeline
pipeline.upsert(rolearn=role)
Start execution
execution = pipeline.start(
parameters={
"InputData": f"s3://{bucket}/iris/raw/data.csv",
"TrainingInstanceType": "ml.m5.xlarge"
}
)
Wait for completion
execution.wait()
Get execution status
print(f"Status: {execution.describe()['PipelineExecutionStatus']}")
Pipeline Visualization
# Get pipeline definition
import json
definition = json.loads(pipeline.definition())
print(json.dumps(definition, indent=2))
List executions
executions = pipeline.listexecutions()
for exe in executions["PipelineExecutionSummaries"]:
print(f"{exe['PipelineExecutionArn']}: {exe['PipelineExecutionStatus']}")
Advanced Features
1. Caching
from sagemaker.workflow.steps import CacheConfig
cacheconfig = CacheConfig(
enablecaching=True,
expireafter="PT1H" # 1 hour
)
stepprocess = ProcessingStep(
name="PreprocessData",
processor=sklearnprocessor,
cacheconfig=cacheconfig,
# ... other parameters
)
2. Callback Steps
from sagemaker.workflow.callbackstep import CallbackStep
step
callback = CallbackStep(
name="NotifySlack",
sqsqueueurl="https://sqs.us-east-1.amazonaws.com/123456789/my-queue",
inputs={
"message": "Pipeline completed successfully",
"modelarn": stepregister.steps[0].properties.ModelPackageArn
},
outputs=[]
)
3. Lambda Steps
from sagemaker.workflow.lambdastep import LambdaStep
from sagemaker.lambda
helper import Lambda
Define Lambda function
lambdafunc = Lambda(
functionarn="arn:aws:lambda:us-east-1:123456789:function:my-function",
session=session
)
steplambda = LambdaStep(
name="InvokeLambda",
lambdafunc=lambdafunc,
inputs={
"modeldata": steptrain.properties.ModelArtifacts.S3ModelArtifacts
}
)
CI/CD Integration
1. CodePipeline Integration
# buildspec.yml
version: 0.2
phases:
install:
runtime-versions:
python: 3.9
commands:
- pip install sagemaker boto3
build:
commands:
- python pipeline.py --action create
postbuild:
commands:
- python pipeline.py --action start
2. GitHub Actions
# .github/workflows/ml-pipeline.yml
name: ML Pipeline
on:
push:
branches: [main]
jobs:
deploy:
runs-on: ubuntu-latest
steps:
- uses: actions/checkout@v3
- name: Configure AWS credentials
uses: aws-actions/configure-aws-credentials@v2
with:
aws-access-key-id: ${{ secrets.AWSACCESSKEYID }}
aws-secret-access-key: ${{ secrets.AWSSECRETACCESSKEY }}
aws-region: us-east-1
- name: Run Pipeline
run: |
pip install sagemaker boto3
python runpipeline.py
Best Practices
1. Parameterize Everything
# Good: Parameterized
traininginstancetype = ParameterString(
name="TrainingInstanceType",
defaultvalue="ml.m5.xlarge"
)
Bad: Hardcoded
instancetype = "ml.m5.xlarge"
2. Use Caching Wisely
# Enable caching for expensive steps
stepprocess = ProcessingStep(
name="PreprocessData",
cacheconfig=CacheConfig(enablecaching=True, expireafter="P7D"),
# ...
)
Disable caching for steps that should always run
stepdeploy = ModelStep(
name="DeployModel",
cacheconfig=CacheConfig(enablecaching=False),
# ...
)
Conclusion
SageMaker Pipelines enables robust MLOps:
Key takeaways:
- Parameterize for flexibility
- Use conditions for quality gates
- Enable caching for efficiency
- Integrate with CI/CD pipelines
- Monitor execution status