Intermediate
KubeFlow Pipeline SDK (KFP v2)
Build ML pipelines using the KFP v2 Python SDK with decorators, typed inputs/outputs, conditional logic, and DAG composition.
Defining Components
In KFP v2, components are Python functions decorated with @dsl.component:
from kfp import dsl
@dsl.component(base_image="python:3.11", packages_to_install=["pandas", "scikit-learn"])
def preprocess_data(input_path: str, output_path: dsl.OutputPath("Dataset")) -> int:
import pandas as pd
df = pd.read_csv(input_path)
df = df.dropna()
df.to_csv(output_path, index=False)
return len(df)
@dsl.component(base_image="python:3.11", packages_to_install=["scikit-learn", "pandas"])
def train_model(dataset: dsl.Input[dsl.Dataset], model_path: dsl.OutputPath("Model"),
learning_rate: float = 0.01) -> float:
import pandas as pd
from sklearn.ensemble import GradientBoostingClassifier
from sklearn.model_selection import train_test_split
import pickle
df = pd.read_csv(dataset.path)
X, y = df.drop("target", axis=1), df["target"]
X_train, X_test, y_train, y_test = train_test_split(X, y, test_size=0.2)
model = GradientBoostingClassifier(learning_rate=learning_rate)
model.fit(X_train, y_train)
accuracy = model.score(X_test, y_test)
with open(model_path, "wb") as f:
pickle.dump(model, f)
return accuracy
Building Pipelines
@dsl.pipeline(name="ml-training-pipeline", description="End-to-end ML training")
def training_pipeline(data_url: str, lr: float = 0.01):
preprocess_task = preprocess_data(input_path=data_url)
train_task = train_model(
dataset=preprocess_task.outputs["output_path"],
learning_rate=lr
)
# Set GPU resources for training
train_task.set_gpu_limit(1)
train_task.set_memory_limit("16Gi")
Conditional Logic and Loops
@dsl.pipeline(name="conditional-pipeline")
def conditional_pipeline(accuracy_threshold: float = 0.9):
train_task = train_model(dataset="gs://bucket/data.csv")
with dsl.Condition(train_task.output > accuracy_threshold):
deploy_model(model=train_task.outputs["model_path"])
with dsl.Condition(train_task.output <= accuracy_threshold):
retrain_with_more_data(model=train_task.outputs["model_path"])
Key concept: KFP v2 compiles pipelines to an intermediate representation (IR YAML) that is portable across execution environments. Always use typed inputs/outputs for better validation and artifact tracking.
Ready to Go Deeper?
Live instructor-led courses from our partners. Affiliate disclosure.
AI & ML Courses - 30% Off
Live instructor-led AI, machine learning, data science, and cloud courses for working professionals. Use code Limited30 at checkout.
EdurekaDataCamp - AI & Data Science
Hands-on Python, machine learning, and AI courses with interactive exercises and real projects.
DataCampedX - Top AI Courses
University-level AI courses from MIT, Harvard, Stanford. Earn certificates that employers recognize.
edX