Complete Vertex AI Pipelines Tutorial: Orchestrating ML Workflows
Vertex AI Pipelines enables you to orchestrate ML workflows as directed acyclic graphs (DAGs). Built on Kubeflow Pipelines, it provides serverless execution with Google Cloud integration.
Why Vertex AI Pipelines?
Key Benefits:- Serverless: No infrastructure to manage
- Reproducible: Version-controlled workflows
- Scalable: Handles large-scale ML jobs
- Integration: Native Google Cloud services
- Reusable: Modular pipeline components
- Automated ML training
- Data preprocessing workflows
- Model deployment pipelines
- Feature engineering
- MLOps automation
Prerequisites
pip install google-cloud-aiplatform kfp
Authenticate
gcloud auth login
gcloud config set project your-project-id
Quick Start
1. Simple Pipeline
from kfp import dsl
from kfp.dsl import component
from google.cloud import aiplatform
Define components
@component
def preprocessdata(inputpath: str, outputpath: str):
import pandas as pd
df = pd.readcsv(inputpath)
df = df.dropna()
df.tocsv(outputpath, index=False)
return outputpath
@component
def trainmodel(datapath: str, modelpath: str) -> float:
import pandas as pd
from sklearn.ensemble import RandomForestClassifier
from sklearn.modelselection import traintestsplit
import joblib
df = pd.readcsv(datapath)
X = df.drop("target", axis=1)
y = df["target"]
Xtrain, Xtest, ytrain, ytest = traintestsplit(X, y, testsize=0.2)
model = RandomForestClassifier(nestimators=100)
model.fit(Xtrain, ytrain)
accuracy = model.score(Xtest, ytest)
joblib.dump(model, modelpath)
return accuracy
Define pipeline
@dsl.pipeline(
name="simple-ml-pipeline",
description="A simple ML training pipeline"
)
def mlpipeline(inputdata: str, modeloutput: str):
preprocesstask = preprocessdata(
inputpath=inputdata,
outputpath="gs://bucket/processed/data.csv"
)
traintask = trainmodel(
datapath=preprocesstask.output,
modelpath=modeloutput
)
Compile and run
from kfp import compiler
compiler.Compiler().compile(
pipelinefunc=mlpipeline,
packagepath="pipeline.json"
)
Submit pipeline
aiplatform.init(project="your-project", location="us-central1")
job = aiplatform.PipelineJob(
displayname="ml-pipeline-run",
templatepath="pipeline.json",
parametervalues={
"inputdata": "gs://bucket/raw/data.csv",
"modeloutput": "gs://bucket/models/model.joblib"
}
)
job.run()
Pipeline Components
1. Python Function Components
from kfp.dsl import component, Input, Output, Dataset, Model, Metrics
@component(
baseimage="python:3.9",
packagestoinstall=["pandas", "scikit-learn"]
)
def trainsklearnmodel(
trainingdata: Input[Dataset],
model: Output[Model],
metrics: Output[Metrics],
nestimators: int = 100,
maxdepth: int = 10
):
import pandas as pd
from sklearn.ensemble import RandomForestClassifier
from sklearn.modelselection import traintestsplit
from sklearn.metrics import accuracyscore, f1score
import joblib
# Load data
df = pd.readcsv(trainingdata.path)
X = df.drop("target", axis=1)
y = df["target"]
Xtrain, Xtest, ytrain, ytest = traintestsplit(X, y, testsize=0.2)
# Train
clf = RandomForestClassifier(nestimators=nestimators, maxdepth=maxdepth)
clf.fit(Xtrain, ytrain)
# Evaluate
predictions = clf.predict(Xtest)
accuracy = accuracyscore(ytest, predictions)
f1 = f1score(ytest, predictions, average="weighted")
# Log metrics
metrics.logmetric("accuracy", accuracy)
metrics.logmetric("f1score", f1)
# Save model
joblib.dump(clf, model.path)
2. Container Components
from kfp.dsl import containercomponent, ContainerSpec
@containercomponent
def customtrainingcomponent(
inputpath: str,
outputpath: str,
epochs: int
):
return ContainerSpec(
image="gcr.io/your-project/training:latest",
command=["python", "train.py"],
args=[
"--input", inputpath,
"--output", outputpath,
"--epochs", str(epochs)
]
)
3. Reusable Components
# Save component to YAML
from kfp.components import createcomponentfromfunc
traincomponent = createcomponentfromfunc(
trainsklearnmodel,
baseimage="python:3.9",
packagestoinstall=["pandas", "scikit-learn"]
)
Save to file
traincomponent.componentspec.save("traincomponent.yaml")
Load component
from kfp.components import loadcomponentfromfile
loadedcomponent = loadcomponentfromfile("traincomponent.yaml")
Complete ML Pipeline
1. Full Pipeline Definition
from kfp import dsl"from kfp.dsl import component, Input, Output, Dataset, Model, Metrics, Artifact
@component(packagestoinstall=["pandas", "google-cloud-bigquery"])
def extractdata(
bqtable: str,
outputdata: Output[Dataset]
):
from google.cloud import bigquery
import pandas as pd
client = bigquery.Client()
query = f"SELECT FROM
{bqtable}df = client.query(query).todataframe()
df.tocsv(outputdata.path, index=False)
@component(packagestoinstall=["pandas", "scikit-learn"])
def preprocess(
inputdata: Input[Dataset],
outputdata: Output[Dataset],
scalerartifact: Output[Artifact]
):
import pandas as pd
from sklearn.preprocessing import StandardScaler
import joblib
df = pd.readcsv(inputdata.path)
# Preprocess
numericcols = df.selectdtypes(include=["number"]).columns
scaler = StandardScaler()
df[numericcols] = scaler.fittransform(df[numericcols])
df.tocsv(outputdata.path, index=False)
joblib.dump(scaler, scalerartifact.path)
@component(packagestoinstall=["pandas", "scikit-learn"])
def train(
trainingdata: Input[Dataset],
model: Output[Model],
metrics: Output[Metrics],
nestimators: int,
maxdepth: int
):
import pandas as pd
from sklearn.ensemble import RandomForestClassifier
from sklearn.modelselection import traintestsplit
from sklearn.metrics import accuracyscore, f1score, precisionscore, recallscore
import joblib
df = pd.readcsv(trainingdata.path)
X = df.drop("target", axis=1)
y = df["target"]
Xtrain, Xtest, ytrain, ytest = traintestsplit(X, y, testsize=0.2)
clf = RandomForestClassifier(nestimators=nestimators, maxdepth=maxdepth)
clf.fit(Xtrain, ytrain)
predictions = clf.predict(Xtest)
metrics.logmetric("accuracy", accuracyscore(ytest, predictions))
metrics.logmetric("f1score", f1score(ytest, predictions, average="weighted"))
metrics.logmetric("precision", precisionscore(ytest, predictions, average="weighted"))
metrics.logmetric("recall", recallscore(ytest, predictions, average="weighted"))
joblib.dump(clf, model.path)
@component(packagestoinstall=["google-cloud-aiplatform"])
def deploymodel(
model: Input[Model],
project: str,
location: str,
endpointname: str
):
from google.cloud import aiplatform
aiplatform.init(project=project, location=location)
# Upload model
uploadedmodel = aiplatform.Model.upload(
displayname="sklearn-model",
artifacturi=model.uri,
servingcontainerimageuri="us-docker.pkg.dev/vertex-ai/prediction/sklearn-cpu.1-0:latest"
)
# Create or get endpoint
endpoints = aiplatform.Endpoint.list(filter=f'displayname="{endpointname}"')
if endpoints:
endpoint = endpoints[0]
else:
endpoint = aiplatform.Endpoint.create(displayname=endpointname)
# Deploy
endpoint.deploy(
model=uploadedmodel,
machinetype="n1-standard-4",
minreplicacount=1,
maxreplicacount=3
)
@dsl.pipeline(
name="complete-ml-pipeline",
description="Complete ML pipeline with preprocessing, training, and deployment"
)
def completepipeline(
bqtable: str,
nestimators: int = 100,
maxdepth: int = 10,
project: str = "your-project",
location: str = "us-central1",
endpointname: str = "ml-endpoint"
):
# Extract
extracttask = extractdata(bqtable=bqtable)
# Preprocess
preprocesstask = preprocess(inputdata=extracttask.outputs["outputdata"])
# Train
traintask = train(
trainingdata=preprocesstask.outputs["outputdata"],
nestimators=nestimators,
maxdepth=maxdepth
)
# Deploy
deploytask = deploymodel(
model=traintask.outputs["model"],
project=project,
location=location,
endpointname=endpointname
)
Conditional Execution
1. Using Conditions
from kfp import dsl
@dsl.pipeline(name="conditional-pipeline")
def conditionalpipeline(accuracythreshold: float = 0.85):
traintask = traincomponent(...)
with dsl.Condition(traintask.outputs["accuracy"] >= accuracythreshold):
deploytask = deploycomponent(
model=traintask.outputs["model"]
)
with dsl.Condition(traintask.outputs["accuracy"] < accuracythreshold):
notifytask = sendnotification(
message="Model accuracy below threshold"
)
2. Exit Handler
@dsl.pipeline(name="pipeline-with-exit-handler")
def pipelinewithcleanup():
with dsl.ExitHandler(cleanupcomponent()):
traintask = traincomponent(...)
deploytask = deploycomponent(model=traintask.outputs["model"])
Parallel Execution
1. Parallel Tasks
@dsl.pipeline(name="parallel-pipeline")
def paralleltrainingpipeline():
# These run in parallel
modela = trainmodela(...)
modelb = trainmodelb(...)
modelc = trainmodelc(...)
# This waits for all
ensemble = createensemble(
modela=modela.outputs["model"],
modelb=modelb.outputs["model"],
modelc=modelc.outputs["model"]
)
2. ParallelFor Loop
from kfp import dsl
@dsl.pipeline(name="parallel-for-pipeline")
def hyperparametersearchpipeline(
learningrates: list = [0.01, 0.001, 0.0001]
):
with dsl.ParallelFor(learningrates) as lr:
traintask = trainwithlr(learningrate=lr)
Scheduling Pipelines
1. Create Schedule
from google.cloud import aiplatform
Create pipeline job template
job = aiplatform.PipelineJob(
displayname="scheduled-pipeline",
templatepath="pipeline.json",
parametervalues={"inputdata": "gs://bucket/data.csv"}
)
Create schedule
schedule = aiplatform.PipelineJobSchedule.create(
displayname="daily-training-schedule",
pipelinejob=job,
cron="0 6 ", # Daily at 6 AM
timezone="UTC"
)
print(f"Schedule created: {schedule.resourcename}")
2. Manage Schedules
# List schedules
schedules = aiplatform.PipelineJobSchedule.list()
for s in schedules:
print(f"{s.displayname}: {s.cron}")
Pause schedule
schedule.pause()
Resume schedule
schedule.resume()
Delete schedule
schedule.delete()
Pipeline Parameters
1. Define Parameters
@dsl.pipeline(
name="parameterized-pipeline",
pipelineroot="gs://bucket/pipeline-root"
)
def parameterizedpipeline(
inputdata: str,
nestimators: int = 100,
maxdepth: int = 10,
learningrate: float = 0.01,
deploy: bool = True
):
traintask = traincomponent(
data=inputdata,
nestimators=nestimators,
maxdepth=maxdepth,
learningrate=learningrate
)
if deploy:
deploytask = deploycomponent(model=traintask.outputs["model"])
2. Run with Parameters
job = aiplatform.PipelineJob(
displayname="pipeline-run",
templatepath="pipeline.json",
parametervalues={
"inputdata": "gs://bucket/data.csv",
"nestimators": 200,
"maxdepth": 15,
"learningrate": 0.001,
"deploy": True
}
)
job.run(sync=True)
Monitoring and Debugging
1. Get Pipeline Status
# Get pipeline job
job = aiplatform.PipelineJob.get("projects/123/locations/us-central1/pipelineJobs/456")
print(f"State: {job.state}")
print(f"Start time: {job.createtime}")
Get task details
for task in job.taskdetails:
print(f"{task.taskname}: {task.state}")
2. View Logs
# Get logs from Cloud Logging
from google.cloud import logging
client = logging.Client()
logger = client.logger("vertex-ai-pipelines")
for entry in logger.listentries(filter=f'resource.labels.jobid="{job.name}"'):
print(entry.payload)
Best Practices
1. Component Caching
@component
def cacheablecomponent(input: str) -> str:
# This component will be cached if inputs don't change
return process(input)
Disable caching for specific task
@dsl.pipeline
def pipeline():
task = cacheablecomponent(input="data")
task.setcachingoptions(enablecaching=False)
2. Resource Specifications
@dsl.pipeline
def resourcepipeline():
traintask = traincomponent(...)
# Set resources
traintask.setcpulimit("4")
traintask.setmemorylimit("16G")
traintask.addnodeselectorconstraint("cloud.google.com/gke-accelerator", "nvidia-tesla-t4")
traintask.setgpu_limit(1)
Conclusion
Vertex AI Pipelines provides:
Key takeaways:
- Build modular, reusable components
- Use conditions for dynamic workflows
- Schedule pipelines for automation
- Monitor pipeline execution
- Cache components for efficiency