Chapter 03 of 12

The ML Pipeline

Build production ML pipelines step by step with complete code examples and real data patterns.

What is an ML Pipeline?

An ML pipeline is an automated sequence of steps that transforms raw data into deployed predictions. Think of it as your ETL pipeline extended to include model training, evaluation, and deployment.

Your ETL vs ML Pipeline

ETL: Extract data > Transform > Load into warehouse > Dashboard reads it
ML: Extract data > Transform > Engineer features > Train model > Evaluate > Deploy > Monitor > Retrain

Pipeline Components

StepInputOutputTools
Data IngestionRaw sourcesUnified datasetSpark, Airflow, dbt
Data ValidationUnified datasetValidated data + reportGreat Expectations
Feature EngineeringValidated dataFeature matrixFeast, pandas
Model TrainingTraining setTrained modelXGBoost, PyTorch
EvaluationTest set + modelMetrics reportMLflow, sklearn
DeploymentBest modelAPI endpointBentoML, Docker
MonitoringPredictionsDrift alertsEvidently, Grafana

Complete Pipeline Code Example

Step 1: Data Ingestion

import pandas as pd
import numpy as np

class DataIngestion:
    def generate_sample_data(self, n=50000):
        np.random.seed(42)
        data = pd.DataFrame({
            "customer_id": [f"C{i:05d}" for i in range(n)],
            "total_orders": np.random.poisson(12, n),
            "avg_order_value": np.random.normal(65, 25, n).clip(5),
            "days_since_last_order": np.random.exponential(20, n).astype(int),
            "support_tickets": np.random.poisson(2, n),
            "satisfaction_score": np.random.normal(3.8, 0.8, n).clip(1, 5),
        })
        churn_prob = (
            0.3 * (data["days_since_last_order"] > 30).astype(float) +
            0.2 * (data["support_tickets"] > 3).astype(float) +
            0.3 * (data["satisfaction_score"] < 3).astype(float)
        )
        data["churned"] = (np.random.random(n) < churn_prob).astype(int)
        return data

Step 2: Data Validation

class DataValidator:
    def validate(self, df):
        checks = {
            "has_rows": len(df) > 1000,
            "no_duplicates": df["customer_id"].nunique() == len(df),
            "nulls_ok": df.isnull().mean().max() < 0.05,
            "target_balanced": 0.05 < df["churned"].mean() < 0.50,
            "score_range": df["satisfaction_score"].between(0, 5).all(),
        }
        for name, passed in checks.items():
            icon = "PASS" if passed else "FAIL"
            print(f"  [{icon}] {name}")
        return all(checks.values())

Step 3: Feature Engineering

class FeatureEngineer:
    def transform(self, df):
        df = df.copy()
        df["orders_per_month"] = df["total_orders"] / 12
        df["spend_per_order"] = df["avg_order_value"]
        df["recency_risk"] = (df["days_since_last_order"] > 30).astype(int)
        df["high_support"] = (df["support_tickets"] > 3).astype(int)
        df["low_satisfaction"] = (df["satisfaction_score"] < 3).astype(int)
        return df.fillna(0)

Step 4: Training with MLflow

import mlflow
from sklearn.ensemble import GradientBoostingClassifier
from sklearn.model_selection import train_test_split
from sklearn.metrics import accuracy_score, f1_score

class ModelTrainer:
    def train(self, df):
        features = [c for c in df.columns if c not in ["customer_id", "churned"]]
        X, y = df[features], df["churned"]
        X_train, X_test, y_train, y_test = train_test_split(X, y, test_size=0.2)
        
        with mlflow.start_run(run_name="gradient_boost"):
            model = GradientBoostingClassifier(n_estimators=200, learning_rate=0.05)
            model.fit(X_train, y_train)
            
            y_pred = model.predict(X_test)
            mlflow.log_metric("accuracy", accuracy_score(y_test, y_pred))
            mlflow.log_metric("f1", f1_score(y_test, y_pred))
            mlflow.sklearn.log_model(model, "model")
            
            return model

Step 5: Run the Complete Pipeline

def run_pipeline():
    print("== Step 1: Ingestion ==")
    data = DataIngestion().generate_sample_data()
    
    print("== Step 2: Validation ==")
    assert DataValidator().validate(data), "Validation failed!"
    
    print("== Step 3: Features ==")
    features = FeatureEngineer().transform(data)
    
    print("== Step 4: Training ==")
    model = ModelTrainer().train(features)
    
    print("Pipeline complete! Check MLflow UI at localhost:5000")

run_pipeline()
Key Pattern

Each step is a separate, testable class. This is the same modular design you use in data engineering. Each step can be independently tested, swapped, or reused.

Pipeline Orchestration Tools

ToolBest ForLearning CurveYour Familiarity
AirflowGeneral orchestrationMediumYou know it!
ZenMLML-specific pipelinesLowEasy to learn
KubeflowK8s-native ML pipelinesHighK8s skills help
MetaflowDS workflows (Netflix)LowEasy to learn
PrefectModern Airflow alternativeLow-MedSimilar to Airflow