Data Engineering for AI

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.

On this page 8
  1. Why it beats a notebook
  2. The stations on the line
  3. The gates are the point
  4. The report at the end
  5. Somewhere you have seen this
  6. What is honestly hard here
  7. Remember this
  8. What to learn next

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 scored

Every 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

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

bash
pip install numpy pandas scikit-learn

The whole line, end to end

pipeline.py
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)
Output
[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

checkcatches
no duplicate customersa retried ingest, and the leakage it causes
labels presentthe upstream outage in run two
income in a sane rangethe rupees-versus-thousands bug, if the fix ever regresses
enough rowsa partial extract, where the source returned one page
both classes presenta 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 GateFailed reaches 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

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:

  1. Input fingerprints for every source, and the output fingerprint per stage.
  2. Row counts in and out per stage, plus the reason for each drop, as counts by reason.
  3. All gate results including passes, with the observed value alongside the threshold. Recording only failures makes trends invisible.
  4. Code version, resolved dependency lockfile hash, and every random seed.
  5. Wall-clock duration and peak memory per stage — the earliest signal of a data volume change.
  6. 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 to learn next

These follow on from what you just read.

  • Feature and Data Pipelines in Production

    Point-in-time correctness

    Point-in-time correctness means every number shown to a model was actually known at that exact moment in the past, never filled in using something from later.

  • Feature and Data Pipelines in Production

    Feature freshness

    Feature freshness is how old a stored number is right now, and whether that age is still safe to trust for the decision it is about to make.

  • Feature and Data Pipelines in Production

    Backfilling a new feature

    Backfilling means computing a brand-new feature for every past date it needs to exist, not only going forward from today.