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.