Data Pipelines for Machine Learning: From Raw Data to Model-Ready Features
Tutorials usually start with a clean, ready-made dataset. Real projects don't. Real data arrives late, with missing values, typos, inconsistent labels and hidden traps. The work of turning that mess into reliable model inputs, again and again, is done by a data pipeline, and it's where most of the effort in real machine learning goes.
This post walks through the stages of a machine learning data pipeline, the mistakes that quietly ruin models, and how to build a robust pipeline in Python with scikit-learn. It follows on from the machine learning series, but it stands on its own.
The model is the small part
It's tempting to think machine learning is mostly about choosing algorithms. In practice, data scientists often say the majority of their time goes into getting, checking and preparing data. The model is one step in a longer chain:
A machine learning pipeline: the model is only one step out of eight.
- Ingest: collect data from databases, files, APIs or event streams.
- Validate: check that the data looks the way you expect before trusting it.
- Clean: fix or remove bad values, duplicates and inconsistencies.
- Transform: turn raw columns into features a model can use.
- Split: separate training and test data, before anything learns from the data.
- Train the model.
- Evaluate it honestly, as covered in the overfitting and metrics posts.
- Serve predictions on new data, then monitor and retrain as the world changes.
Validate before you trust
Data breaks silently. An upstream system changes a date format, a sensor starts sending zeros, or a column suddenly fills with blanks. The model will happily train on garbage and produce confident garbage. A few automated checks catch most of these problems early:
- Schema: are all expected columns present, with the right types?
- Completeness: what share of each column is missing, and has it changed?
- Ranges: are values plausible? Ages shouldn't be negative, and percentages shouldn't exceed 100.
- Categories: have new or misspelled categories appeared?
- Volume: did today's batch arrive with roughly the expected number of rows?
Libraries such as Great Expectations and Pandera turn checks like these into automated tests that run every time new data arrives.
Clean the mess
Cleaning fixes what validation finds. The most common jobs are handling missing values and making categories consistent:
Typical cleaning: fill missing values and make categories consistent, so “Premium ” and “premium” count as the same plan.
Missing values can be filled (imputed) with a typical value such as the median, flagged with an extra "was missing" column, or, if very few rows are affected, dropped. Which is right depends on why the data is missing. Text categories need standardizing, because to a computer "Premium", "premium" and "Premium " (with a trailing space) are three different things.
Transform raw columns into features
Most models need numbers on comparable scales. Two transformations cover a large share of everyday work:
- Scaling puts numeric columns on a similar scale, for example mean 0 and standard deviation 1, so a column measured in thousands doesn't drown out one measured in single digits.
- One-hot encoding turns a category into several 0/1 columns, one per category, so the model doesn't mistakenly treat "premium" as numerically bigger than "basic".
In scikit-learn, a ColumnTransformer applies the right steps to each type of column, and a Pipeline chains everything, including the model, into one object:
A ColumnTransformer sends each type of column through the right preparation steps, then combines them for the model.
Feature engineering can go much further, for example turning a signup date into "months as a customer", or counting support calls in the last 30 days. Good features often improve a model more than a fancier algorithm does.
The silent killer: data leakage
Data leakage happens when information that won't be available at prediction time sneaks into training. The model looks brilliant during development and then fails in production. It comes in two main forms:
- Target leakage: a feature that is really a consequence of the answer. For example, a "cancellation fee charged" column in a churn model: the fee is only charged after a customer leaves, so it gives the answer away.
- Train-test contamination: preparing data using statistics from the whole dataset, such as scaling with the overall average, before splitting. The test set then isn't truly unseen.
Real results from the code below. One leaky column made a mediocre model look nearly perfect.
The warning sign is a score that seems too good to be true. When you see one, ask of every feature: "Would I actually know this value at the moment I need to make the prediction?" Putting all preparation steps inside a Pipeline, and splitting the data first, prevents train-test contamination automatically, because each step learns only from the training data.
The complete pipeline in Python
The program below creates a realistic messy customer dataset with missing values, inconsistent plan names and a leaky column. It then validates it, cleans it, builds a single pipeline, demonstrates the leakage trap, and uses the same pipeline to score a brand-new customer. Each numbered comment matches a pipeline stage.
import numpy as np
import pandas as pd
from sklearn.model_selection import train_test_split, cross_val_score
from sklearn.compose import ColumnTransformer
from sklearn.pipeline import Pipeline
from sklearn.impute import SimpleImputer
from sklearn.preprocessing import StandardScaler, OneHotEncoder
from sklearn.linear_model import LogisticRegression
# 1. Ingest: a messy customer dataset (generated here so the example runs anywhere)
rng = np.random.default_rng(0)
n = 1000
df = pd.DataFrame({
"tenure_months": rng.integers(1, 72, n).astype(float),
"monthly_bill": rng.normal(70, 25, n).round(2),
"plan": rng.choice(["basic", "standard", "premium", "Premium "], n),
"support_calls": rng.poisson(2, n).astype(float),
})
risk = -0.04 * df.tenure_months + 0.02 * df.monthly_bill + 0.5 * df.support_calls - 1
df["churned"] = (rng.random(n) < 1 / (1 + np.exp(-risk))).astype(int)
# Real-world mess: missing values and a column that leaks the answer
df.loc[rng.random(n) < 0.08, "monthly_bill"] = np.nan
df.loc[rng.random(n) < 0.05, "tenure_months"] = np.nan
df["cancellation_fee_charged"] = df.churned * (rng.random(n) < 0.95)
# 2. Validate: check the data before trusting it
print("Rows:", len(df), "| missing values:", df.isna().sum().to_dict())
print("Plan values:", sorted(df.plan.unique()))
# 3. Clean: fix inconsistent categories
df["plan"] = df.plan.str.strip().str.lower()
# 4. Transform: one reusable pipeline for numbers and categories
numeric = ["tenure_months", "monthly_bill", "support_calls"]
categorical = ["plan"]
preprocess = ColumnTransformer([
("num", Pipeline([("impute", SimpleImputer(strategy="median")),
("scale", StandardScaler())]), numeric),
("cat", OneHotEncoder(handle_unknown="ignore"), categorical),
])
model = Pipeline([("prep", preprocess), ("clf", LogisticRegression(max_iter=1000))])
# 5. Split, then train and evaluate
y = df.churned
X_train, X_test, y_train, y_test = train_test_split(
df.drop(columns="churned"), y, test_size=0.2, random_state=42, stratify=y)
# The leak: a feature that is only known AFTER a customer churns
leaky = Pipeline([("prep", ColumnTransformer(
preprocess.transformers + [("leak", "passthrough", ["cancellation_fee_charged"])])),
("clf", LogisticRegression(max_iter=1000))])
print("With leaky feature: ", round(cross_val_score(leaky, X_train, y_train, cv=5).mean(), 3))
print("Without leaky feature:", round(cross_val_score(model, X_train, y_train, cv=5).mean(), 3))
model.fit(X_train, y_train)
print("Test accuracy:", round(model.score(X_test, y_test), 3))
# 6. Serve: the same pipeline handles a brand-new, messy record
new_customer = pd.DataFrame([{"tenure_months": np.nan, "monthly_bill": 95.0,
"plan": "PREMIUM", "support_calls": 5.0,
"cancellation_fee_charged": 0}])
new_customer["plan"] = new_customer.plan.str.strip().str.lower()
print("Churn probability for new customer:", round(model.predict_proba(new_customer)[0, 1], 2))
Running it prints:
Rows: 1000 | missing values: {'tenure_months': 43, 'monthly_bill': 93, 'plan': 0, 'support_calls': 0, 'churned': 0, 'cancellation_fee_charged': 0}
Plan values: ['Premium ', 'basic', 'premium', 'standard']
With leaky feature: 0.97
Without leaky feature: 0.681
Test accuracy: 0.685
Churn probability for new customer: 0.84
Three things to notice. First, validation caught both problems: 136 missing values and a "Premium " category with a stray space. Second, the leaky feature inflated accuracy from 68% to 97%; in production, where that column is always empty for current customers, the model would fall apart. Third, the new customer had a missing tenure and an uppercase plan name, and the pipeline handled both without any extra code. (The honest score of about 68% is modest because this dataset was generated with plenty of randomness; roughly half of its customers churn, so guessing would score about 51%.)
Training and serving: one pipeline, no surprises
A classic production bug is training-serving skew: data is prepared one way during training and slightly differently when the model is live, for example with a different fill value or a missing category. Predictions quietly degrade. The fix is to save the whole pipeline, preparation included, as one object and use it everywhere:
Use the exact same pipeline for training and for live predictions, so new data is prepared exactly like the training data.
In scikit-learn, you can save the fitted pipeline with joblib.dump(model, "churn_model.joblib") and load it in your application. Larger teams often go further with a feature store, a shared system that computes each feature once and serves the identical values for training and live predictions.
Batch or streaming?
- Batch pipelines run on a schedule, such as nightly, processing data in large chunks. They are simpler and suit most reporting and many predictions, such as weekly churn scores. Tools include Apache Airflow, dbt and Spark.
- Streaming pipelines process events as they happen, within seconds. They suit fraud detection or real-time recommendations. Tools include Apache Kafka and Flink.
Start with batch unless you truly need real-time answers. It's far easier to build, test and fix.
A pipeline checklist
- Split your data before fitting any preparation step.
- Put every preparation step inside one Pipeline object.
- Validate incoming data automatically, every run.
- Question any feature you wouldn't know at prediction time.
- Save and reuse the same fitted pipeline for serving.
- Monitor live data and model performance, and retrain when they drift.
Key terms, decoded
- Data pipeline
- An automated series of steps that moves data from its sources to a usable form.
- Data validation
- Automatic checks that data has the expected structure and plausible values.
- Imputation
- Filling in missing values, for example with the median.
- Feature engineering
- Creating model inputs from raw data.
- One-hot encoding
- Turning a category into separate 0/1 columns.
- Data leakage
- Information unavailable at prediction time sneaking into training.
- Training-serving skew
- Differences between how data is prepared for training and for live predictions.
- Feature store
- A shared system that serves identical features for training and production.
- Data drift
- Live data gradually changing so it no longer resembles the training data.
Great models are built on boring, reliable plumbing: validate, clean, transform, split and serve, all in one reproducible pipeline. Get that right and every model you build becomes easier to trust. Coming up next in the series: a complete RAG app with real embeddings and a vector database.
What's the messiest dataset you've had to clean? Share it in the comments.

Comments
Post a Comment