Building a data pipeline
A pipeline is every step from raw records to a trained model written down as code, with checks between the steps that stop the run when the data is wrong.
- 15 min read
- 3 reading levels
- Updated
Read these first
On this page 8
One lesson, three depths. Pick the one that fits you today — you can switch any time.
Beginner — No maths. Plain English.
A pipeline is the whole journey from raw records to a trained model, written down as code that anybody can run.
Think about the assembly of a thali in a busy canteen. The plate moves along: rice here, dal there, sabzi, roti, pickle. Each station does one thing and passes it on.
If the dal station is out, the line stops. Nobody serves a plate with an empty bowl and hopes the customer will not notice.
A data pipeline is that line, with the stations being steps and the stopping being deliberate.
Why it beats a notebook
Most models start life in a notebook. Cells run in whatever order somebody clicked. It works on a Tuesday afternoon and nobody can reproduce it on Friday.
A pipeline fixes four things:
- It runs the same way every time, because the order is written down.
- It runs without a person, so it can run nightly.
- It fails loudly, instead of producing a quietly wrong model.
- It leaves a record of what happened, so a change can be traced.
The stations on the line
raw records
|
[ extract ] pull from the source
|
[ validate ] is this data fit to use? --- no ---> STOP
|
[ clean ] repeats, units, blanks
|
[ validate ] did cleaning break anything? --- no ---> STOP
|
[ features ] build the columns the model reads
|
[ split ] train and test, no leaks
|
[ train ] fit and score
|
report: what came in, what was dropped, what it scoredEvery step in this section of the site is one of those boxes. This lesson is where they get bolted together.
The gates are the point
Those validate boxes are the whole reason a pipeline is better than a script.
A gate asks plain questions. Are there enough rows? Are the answers present? Is anything impossible? If any answer is wrong, the run stops and nothing is trained.
Stopping is the correct behaviour. A pipeline that trains on broken data produces something worse than nothing. It produces a wrong model wearing the badge of an automated system.
The report at the end
Every run should print what it did. Rows in. Rows dropped, and why. The score. A fingerprint of the data.
That report is what somebody reads six months later when they ask why the model changed. Without it, every investigation starts from zero.
Somewhere you have seen this
The recommendations on a shopping app refresh overnight. Behind that is a pipeline that ran at 3am, pulled yesterday's activity, checked it, rebuilt features and retrained.
If it failed, an on-call engineer got a message, and yesterday's model kept serving. That fallback is part of the design, not an accident.
What is honestly hard here
Deciding where a gate should sit, and how strict it should be, takes experience you do not have on day one.
Too strict and the pipeline stops every night for harmless variation, and people start ignoring the alerts. Too loose and it lets through the batch that matters.
You will get this wrong at first. Everybody does. Start with a few basic checks, and add one every time something slips through.
Remember this
- A pipeline is every step written down as code, in a fixed order.
- Gates between the steps are the point. Stopping beats training on bad data.
- Print a report every run. It is what makes a future investigation possible.
What to learn next
- What is MLOps? — everything that happens after the pipeline produces a model.
- CI/CD for machine learning — testing and shipping the pipeline itself.
- Monitoring and model drift — the gates that run after deployment, not before it.
Developer — Code and libraries.
Everything in this section, assembled into one file you can run. Two runs: an ordinary day, and a day when the upstream label feed breaks.
Setup
pip install numpy pandas scikit-learnThe whole line, end to end
import hashlib
import numpy as np
import pandas as pd
from sklearn.linear_model import LogisticRegression
from sklearn.metrics import roc_auc_score
from sklearn.model_selection import train_test_split
from sklearn.pipeline import make_pipeline
from sklearn.preprocessing import StandardScaler
class GateFailed(Exception):
"""Raised when data is not fit to train on. Stopping is the correct behaviour."""
def extract(seed=0, break_it=False):
rng = np.random.default_rng(seed)
n = 1500
income = rng.gamma(4.0, 15.0, n) + 10
years_job = rng.integers(0, 20, n).astype(float)
repaid = rng.binomial(1, 1/(1+np.exp(-(0.05*income + 0.12*years_job - 4.0))))
df = pd.DataFrame({"customer_id": np.arange(n),
"city": rng.choice([" Pune", "mumbai", "MUMBAI ", "nagpur"], n),
"income": income, "years_job": years_job, "repaid": repaid})
df.loc[rng.random(n) < 0.05, "income"] = np.nan # some people decline
df.loc[rng.random(n) < 0.02, "income"] *= 1000 # some rows typed in rupees
df = pd.concat([df, df.sample(60, random_state=1)], ignore_index=True) # ingest retried
if break_it:
df.loc[df.sample(700, random_state=2).index, "repaid"] = np.nan # upstream label outage
return df
def clean(df):
out = df.drop_duplicates(subset="customer_id", keep="first").copy()
out["city"] = out["city"].str.strip().str.lower()
out["income"] = np.where(out["income"] > 1000, out["income"]/1000, out["income"])
out["income_missing"] = out["income"].isna().astype(int)
out["income"] = out["income"].fillna(out["income"].median())
return out
def gate(df, stage):
checks = {
"no duplicate customers": not df["customer_id"].duplicated().any(),
"labels present": df["repaid"].isna().mean() < 0.01,
"income in a sane range": df["income"].between(0, 1000).all(),
"enough rows": len(df) >= 1000,
"both classes present": df["repaid"].nunique() == 2,
}
failed = [name for name, ok in checks.items() if not ok]
print(f" gate after {stage}: " + ("PASS" if not failed else "FAIL -> " + ", ".join(failed)))
if failed:
raise GateFailed(f"{stage}: {failed}")
return df
def fingerprint(df):
return hashlib.sha256(df.to_csv(index=False).encode()).hexdigest()[:12]
def train(df):
cols = ["income", "years_job", "income_missing"]
tr, te = train_test_split(df, test_size=0.3, random_state=42, stratify=df["repaid"])
m = make_pipeline(StandardScaler(), LogisticRegression(max_iter=1000)).fit(tr[cols], tr["repaid"])
return round(roc_auc_score(te["repaid"], m.predict_proba(te[cols])[:, 1]), 3)
def run(name, **kw):
print(f"[{name}]")
try:
raw = extract(**kw)
print(f" extracted {len(raw)} rows, fingerprint {fingerprint(raw)}")
tidy = gate(clean(raw), "clean")
print(f" kept {len(tidy)} rows, fingerprint {fingerprint(tidy)}")
print(f" test AUC {train(tidy)}")
except GateFailed as e:
print(f" pipeline stopped, no model was trained -> {e}")
print()
run("normal day")
run("upstream label outage", break_it=True)[normal day] extracted 1560 rows, fingerprint ddfed40a185d gate after clean: PASS kept 1500 rows, fingerprint 210d6b011ede test AUC 0.803 [upstream label outage] extracted 1560 rows, fingerprint a23c91c42338 gate after clean: FAIL -> labels present pipeline stopped, no model was trained -> clean: ['labels present']
What the second run bought you
On the outage day, 674 of the 1500 cleaned rows arrived with no label at all, leaving 826 usable. Without the gate, that reaches train_test_split and the fitter, which either raise an error somewhere far less legible or train on a distorted class balance. Either way the failure surfaces late, in a stack trace that does not mention the real cause.
With the gate, it stopped in one line and told you which check failed. That message is the difference between a five-minute fix and a two-day investigation.
Note what did not happen: no model was written. Yesterday's model keeps serving. This is the behaviour you want, and it comes from raising an exception instead of logging a warning.
The five gate checks, and why each is there
| check | catches |
|---|---|
| no duplicate customers | a retried ingest, and the leakage it causes |
| labels present | the upstream outage in run two |
| income in a sane range | the rupees-versus-thousands bug, if the fix ever regresses |
| enough rows | a partial extract, where the source returned one page |
| both classes present | a filter that accidentally removed every positive |
Each of these has cost somebody a week somewhere. The list is short deliberately — five checks that run on every batch beat fifty that nobody maintains.
The income in a sane range check is worth a second look. The cleaning step already fixed the units, so this gate verifies the fix rather than the input. Gates after a transformation are how you find out that your own code broke, and they are the ones most often left out.
The fingerprints
extract produced ddfed40a185d on the normal day and a23c91c42338 on the outage day. Different data, different fingerprint, visible immediately.
The post-clean fingerprint 210d6b011ede is what you record next to the model. When somebody asks in March which rows trained this thing, that string is the answer. See data versioning.
Moving this to a scheduler
The structure above is what Airflow, Dagster and Prefect run for you. What they add:
- Retries with backoff for steps that fail on transient errors, which is why each step must be idempotent. See ETL for machine learning.
- Partitioned runs. Every run is parameterised by a date, so a backfill re-runs the past with the past's parameters.
- Dependency graphs, so independent steps run in parallel and a failure blocks only what depends on it.
- Alerting, so a
GateFailedreaches a person rather than a log file.
Do not reach for a scheduler on day one. A run() function and cron covers a single daily pipeline. Add the orchestrator when you have several pipelines that share steps, or when backfills become routine.
Common mistakes
Warning instead of raising. A warning is a message nobody reads. If the data is unfit, raise. If it is only unusual, log it with a number so you can plot the trend.
Gates only at the start. Half your bugs come from your own transformations. Gate after each significant step.
Checks with no thresholds you chose. df["income"].between(0, 1000) is a decision about your domain. Write the reasoning in a comment, or the next person will loosen it to make an alert stop.
Running steps in a notebook in a different order. The order in the code is the contract. If it only works when you run cell 7 before cell 3, it is not a pipeline.
No fallback when the pipeline stops. Stopping is right only if the previous model keeps serving. Decide what happens on failure before it happens.
A report nobody stores. Print the counts, and append them to a file. A single run's numbers are noise; the series is the signal, and it is how you spot a step slowly degrading.
Try it yourself
Add a sixth check: the fraction of rows with income_missing == 1 must stay between 0.02 and 0.10. Then change the extract's missing rate from 0.05 to 0.30 and confirm the gate fires.
That is a distribution check, not a correctness check. Nothing in the data is invalid — the shape changed. These are the checks that catch upstream changes nobody told you about, and they are the hardest to tune. Start with wide bounds and tighten them as you learn what a normal week looks like.
What to learn next
- What is MLOps? — everything that happens after the pipeline produces a model.
- CI/CD for machine learning — testing and shipping the pipeline itself.
- Monitoring and model drift — the gates that run after deployment, not before it.
Researcher — Mathematics and papers.
The pipeline as a directed acyclic graph
A pipeline is a DAG $G = (V, E)$ where vertices are transformations and edges are data dependencies. Two properties make it tractable:
Determinism. For a fixed input and fixed parameters, each node produces identical output. This is what makes caching, memoisation and incremental recomputation sound. Sources of non-determinism to eliminate deliberately: unseeded RNGs, dict iteration in older runtimes, parallel float reduction order, and any use of wall-clock now() inside a transformation rather than as an injected parameter.
Idempotency. Re-executing a node with the same inputs leaves the system in the same state, as in ETL for machine learning. This is what makes retries safe, and retries are the only reason distributed pipelines survive.
Given both, the correct rerun set after a change to node $v$ is the set of descendants of $v$ in $G$, computed by reachability in $O(|V| + |E|)$. Without both, the correct rerun set is $V$, and incremental computation is unsound.
Staleness as a first-class concept
An asset $a$ is stale if any upstream asset has a materialisation newer than $a$'s, or if the code producing $a$ has changed since. Formally, with $\tau(a)$ the materialisation time and $c(a)$ the code version:
$$ \text{stale}(a) \iff \exists\, u \in \text{parents}(a) : \tau(u) > \tau(a) \;\;\vee\;\; c(a) \neq c_{\text{recorded}}(a) $$
Task-centric orchestrators (Airflow's original model) track task runs, not asset states, so this question cannot be answered from the scheduler's own metadata. Asset-centric orchestrators (Dagster, and Airflow's later datasets feature) track it directly, which is why ML workloads have drifted towards them: the deliverable is a table or a model, and "is this model current with respect to its inputs" is the question people actually ask.
Validation as statistical hypothesis testing
Gates fall into two categories, and conflating them causes most alert fatigue.
Schema and integrity constraints are deterministic predicates. They are either satisfied or not; there is no false positive rate. Type checks, uniqueness, referential integrity, domain ranges.
Distributional checks are hypothesis tests, and every one has a false alarm rate. With $k$ checks at significance $\alpha$ running daily, the expected number of spurious alerts per year is approximately $365 \, k \, \alpha$. At $k = 50$ and $\alpha = 0.01$, that is 182 false alarms a year, which is precisely how teams learn to ignore the pipeline.
Mitigations that actually work:
- Control the false discovery rate across checks (Benjamini-Hochberg) rather than testing each independently.
- Prefer effect-size thresholds to p-values. A KS statistic exceeding 0.1 is a statement about magnitude; $p < 0.05$ on ten million rows is a statement about sample size.
- Learn bounds from a rolling baseline window rather than fixing them once, so seasonal variation does not trip the gate.
Breck et al. (2019), Data Validation for Machine Learning, SysML, describe Google's approach: infer a schema from a baseline, evolve it under human review, and treat every anomaly as a request for a schema amendment rather than an automatic block. The human-in-the-loop step is the part most reimplementations drop, and it is the part that makes the system survivable.
What to record per run
The minimum for post-hoc debugging:
- Input fingerprints for every source, and the output fingerprint per stage.
- Row counts in and out per stage, plus the reason for each drop, as counts by reason.
- All gate results including passes, with the observed value alongside the threshold. Recording only failures makes trends invisible.
- Code version, resolved dependency lockfile hash, and every random seed.
- Wall-clock duration and peak memory per stage — the earliest signal of a data volume change.
- The evaluation metrics, joined to the data fingerprint.
Items 2 and 3 turn a pipeline into a time series about itself. A stage whose drop rate moves from 2% to 6% over three weeks is the kind of failure no single-run check catches, and it is common.
The ML Test Score
Breck et al. (2017), The ML Test Score: A Rubric for ML Production Readiness and Technical Debt Reduction, IEEE Big Data, propose 28 tests across four areas: features and data, model development, infrastructure, and monitoring. Scoring is deliberately harsh — a point for a manual test, two for an automated one, with the final score the minimum across the four areas rather than the sum.
The minimum rule is the interesting design decision. A team with excellent model tests and no data tests scores zero, which matches the failure pattern observed in practice.
The data-and-features tests are: feature expectations captured in a schema, features benefit-cost tested, features free of disallowed content, a data pipeline with appropriate privacy controls, new features added quickly, and all input feature code tested. Most teams pass fewer than half.
Reproducibility in the strict sense
Bit-identical reproduction requires more than pinned data and code:
- Pinned library versions and the transitive dependency graph, since a patch release of a numerical library changes results.
- Fixed thread counts, because parallel floating-point reduction is not associative and results vary with the number of workers.
- Deterministic GPU kernels where applicable, which usually costs measurable throughput.
- A fixed hardware target, since instruction sets differ in fused-multiply-add availability.
Most teams need statistical reproducibility — the same conclusion within noise — rather than bit-identity, and should say which they mean. Claiming bit-identity and delivering statistical reproducibility is how a reproducibility claim becomes worthless. Run the same configuration with several seeds and report the spread; a difference smaller than the seed-to-seed spread is not a result.
Reading
- Breck et al., The ML Test Score, IEEE Big Data 2017.
- Breck et al., Data Validation for Machine Learning, SysML 2019.
- Sculley et al., Hidden Technical Debt in Machine Learning Systems, NeurIPS 2015.
- Polyzotis et al., Data Management Challenges in Production Machine Learning, SIGMOD 2017.
- Kleppmann, Designing Data-Intensive Applications, O'Reilly 2017 — chapter 10 on batch processing and dataflow.
What to learn next
- What is MLOps? — everything that happens after the pipeline produces a model.
- CI/CD for machine learning — testing and shipping the pipeline itself.
- Monitoring and model drift — the gates that run after deployment, not before it.