Tutorial Vertex AI Pipelines: Orkestrasi ML Pipeline

# 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

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

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

  • Orkestrasi: Workflow berbasis DAG
  • Serverless: Eksekusi managed
  • Reproducibility: Pipeline dengan version control
  • Scalability: Processing terdistribusi
  • Integration: Layanan Google Cloud
  • Key takeaways:

    • Bangun komponen modular dan reusable
    • Gunakan conditions untuk workflow dinamis
    • Jadwalkan pipelines untuk automasi
    • Monitor eksekusi pipeline
    • Cache komponen untuk efisiensi

    Artikel Terkait

    Tutorial Vertex AI Model Monitoring: Observabilitas Model Produksi

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

    Tutorial Vertex AI Feature Store: Manajemen Feature Terpusat

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

    Tutorial Lengkap Vertex AI: Platform ML Terpadu Google Cloud

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

    Tutorial Azure ML Pipelines: Automasi Pipeline ML

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