Build production ML pipelines step by step with complete code examples and real data patterns.
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.
ETL: Extract data > Transform > Load into warehouse > Dashboard reads it
ML: Extract data > Transform > Engineer features > Train model > Evaluate > Deploy > Monitor > Retrain
| Step | Input | Output | Tools |
|---|---|---|---|
| Data Ingestion | Raw sources | Unified dataset | Spark, Airflow, dbt |
| Data Validation | Unified dataset | Validated data + report | Great Expectations |
| Feature Engineering | Validated data | Feature matrix | Feast, pandas |
| Model Training | Training set | Trained model | XGBoost, PyTorch |
| Evaluation | Test set + model | Metrics report | MLflow, sklearn |
| Deployment | Best model | API endpoint | BentoML, Docker |
| Monitoring | Predictions | Drift alerts | Evidently, Grafana |
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 dataclass 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())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)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 modeldef 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()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.
| Tool | Best For | Learning Curve | Your Familiarity |
|---|---|---|---|
| Airflow | General orchestration | Medium | You know it! |
| ZenML | ML-specific pipelines | Low | Easy to learn |
| Kubeflow | K8s-native ML pipelines | High | K8s skills help |
| Metaflow | DS workflows (Netflix) | Low | Easy to learn |
| Prefect | Modern Airflow alternative | Low-Med | Similar to Airflow |