Vertex AI Pipelines Tutorial: ML Pipeline Orchestration

# Tutorial Lengkap Vertex AI Pipelines: Orkestrasi Workflow ML Vertex AI Pipelines memungkinkan Anda mengorkestrasi workflow ML sebagai directed acyclic graphs (DAGs). Dibangun di atas Kubeflow Pipel...

By Ruby Abdullah · · tutorial
GCPVertex AIPipelinesKubeflowMLOpsAutomation

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

Use Cases:
  • 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:

  • Orchestration: DAG-based workflows
  • Serverless: Managed execution
  • Reproducibility: Version-controlled pipelines
  • Scalability: Distributed processing
  • Integration: Google Cloud services
  • Key takeaways:

    • Build modular, reusable components
    • Use conditions for dynamic workflows
    • Schedule pipelines for automation
    • Monitor pipeline execution
    • Cache components for efficiency

    Related Articles

    Vertex AI Model Monitoring Tutorial: Production Model Observability

    Tutorial Lengkap Vertex AI Model Monitoring: Monitoring ML Berkelanjutan Vertex AI Model Monitoring secara otomatis mend...

    Vertex AI Feature Store Tutorial: Centralized Feature Management

    Tutorial Lengkap Vertex AI Feature Store: Manajemen Fitur Terpusat Vertex AI Feature Store adalah repositori terpusat un...

    Complete Vertex AI Tutorial: Google Cloud Unified ML Platform

    Tutorial Lengkap Vertex AI: Platform ML Terpadu di Google Cloud Vertex AI adalah platform machine learning terpadu Googl...

    Azure ML Pipelines Tutorial: ML Pipeline Automation

    Tutorial Lengkap Azure ML Pipelines: CI/CD untuk Machine Learning Azure ML Pipelines memungkinkan Anda membangun workflo...