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 Pipelines, menyediakan eksekusi serverless dengan integrasi Google Cloud.
Mengapa Vertex AI Pipelines?
Manfaat Utama:- Serverless: Tidak perlu mengelola infrastruktur
- Reproducible: Workflow dengan version control
- Scalable: Menangani ML jobs skala besar
- Integration: Layanan Google Cloud native
- Reusable: Komponen pipeline modular
- Automated ML training
- Data preprocessing workflows
- Pipeline deployment model
- Feature engineering
- Automasi MLOps
Prerequisites
pip install google-cloud-aiplatform kfp
Autentikasi
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
Definisikan 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
Definisikan pipeline
@dsl.pipeline(
name="simple-ml-pipeline",
description="Pipeline ML training sederhana"
)
def mlpipeline(inputdata: str, modeloutput: str):
preprocesstask = preprocessdata(
inputpath=inputdata,
outputpath="gs://bucket/processed/data.csv"
)
traintask = trainmodel(
datapath=preprocesstask.output,
modelpath=modeloutput
)
Compile dan jalankan
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)
# Evaluasi
predictions = clf.predict(Xtest)
accuracy = accuracyscore(ytest, predictions)
f1 = f1score(ytest, predictions, average="weighted")
# Log metrics
metrics.logmetric("accuracy", accuracy)
metrics.logmetric("f1score", f1)
# Simpan 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
# Simpan component ke YAML
from kfp.components import createcomponentfromfunc
traincomponent = createcomponentfromfunc(
trainsklearnmodel,
baseimage="python:3.9",
packagestoinstall=["pandas", "scikit-learn"]
)
Simpan ke file
traincomponent.componentspec.save("traincomponent.yaml")
Load component
from kfp.components import loadcomponentfromfile
loadedcomponent = loadcomponentfromfile("traincomponent.yaml")
Complete ML Pipeline
1. Definisi Pipeline Lengkap
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"
)
# Buat atau dapatkan 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="Pipeline ML lengkap dengan preprocessing, training, dan 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
)
Eksekusi Kondisional
1. Menggunakan 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="Akurasi model di bawah 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"])
Eksekusi Paralel
1. Parallel Tasks
@dsl.pipeline(name="parallel-pipeline")
def paralleltrainingpipeline():
# Ini berjalan paralel
modela = trainmodela(...)
modelb = trainmodelb(...)
modelc = trainmodelc(...)
# Ini menunggu semua selesai
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)
Penjadwalan Pipelines
1. Buat Schedule
from google.cloud import aiplatform
Buat pipeline job template
job = aiplatform.PipelineJob(
displayname="scheduled-pipeline",
templatepath="pipeline.json",
parametervalues={"inputdata": "gs://bucket/data.csv"}
)
Buat schedule
schedule = aiplatform.PipelineJobSchedule.create(
displayname="daily-training-schedule",
pipelinejob=job,
cron="0 6 ", # Harian jam 6 pagi
timezone="Asia/Jakarta"
)
print(f"Schedule dibuat: {schedule.resourcename}")
2. Kelola 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()
Hapus schedule
schedule.delete()
Parameter Pipeline
1. Definisikan 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. Jalankan dengan 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 dan Debugging
1. Dapatkan Status Pipeline
# Dapatkan pipeline job
job = aiplatform.PipelineJob.get("projects/123/locations/us-central1/pipelineJobs/456")
print(f"State: {job.state}")
print(f"Start time: {job.createtime}")
Dapatkan detail task
for task in job.taskdetails:
print(f"{task.taskname}: {task.state}")
2. Lihat Logs
# Dapatkan logs dari 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:
# Component ini akan di-cache jika input tidak berubah
return process(input)
Disable caching untuk task spesifik
@dsl.pipeline
def pipeline():
task = cacheablecomponent(input="data")
task.setcachingoptions(enablecaching=False)
2. Spesifikasi Resource
@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)
Kesimpulan
Vertex AI Pipelines menyediakan:
Key takeaways:
- Bangun komponen modular dan reusable
- Gunakan conditions untuk workflow dinamis
- Jadwalkan pipelines untuk automasi
- Monitor eksekusi pipeline
- Cache komponen untuk efisiensi