Most teams' first ML pipeline is a set of scripts that run in order: pull data, build features, train, evaluate, push a model. It works until the questions get harder. Which data trained the model that is serving now? Why did last night's run take six hours when only one feature changed? Why did the model get worse when no one touched the code? A v1 pipeline answers these with archaeology.
This article describes the redesign that most teams arrive at, here called v2, and how to get there without a risky cutover. The core change is small to state: every step becomes a function of declared inputs, and every output becomes an immutable artifact identified by a fingerprint of its content. Caching, lineage, reproducibility and safe backfills all follow from that. The stages themselves, data validation, feature computation, training and promotion, are covered in AI/ML pipeline architecture; this page is about the structure that holds them together.
What fails in a v1 pipeline
A v1 pipeline usually has four habits that cause most of its incidents. Steps communicate through mutable locations, so features/latest.parquet is overwritten each night and no one can say which version a model saw. Steps read undeclared inputs: a query for yesterday's partition computed from the clock, an environment variable, a table someone else owns. Reruns are all-or-nothing, so a change to the evaluation script retrains the model. And the record of what happened lives in logs, if anywhere.
Each habit is reasonable on its own. Together they mean that the pipeline's output is a function of things nobody wrote down, which is why it cannot be reproduced, why it cannot be partially rerun, and why a silent upstream change shows up only as a worse model. The v2 design removes the habits rather than adding monitoring on top of them.
The v2 principles
- Artifacts are immutable and typed. A dataset, feature table, model or evaluation report is written once to a location derived from its fingerprint and never modified. A new version is a new artifact.
- Steps are functions of declared inputs. A step reads only the artifacts and parameters it is given. Time, if it matters, is a parameter such as
as_of=2026-09-30, not a call to the clock. - Identity comes from content and recipe. A step's cache key hashes its code, environment, parameters and the fingerprints of its inputs. Same key, same output, so the step can be skipped.
- Contracts sit at boundaries. Each step states what it expects of its inputs and promises about its output, and the runner checks both.
- Lineage is recorded by the runner, not by step authors, so it is complete by construction.
Artifacts and fingerprints
An artifact needs an identity that changes exactly when its content changes. Hashing every byte of a terabyte dataset on every run is too slow, so fingerprint a manifest instead: the sorted list of files with their sizes and content hashes, which object stores usually provide or which you compute once when the file is written. A model's fingerprint is the hash of its weight files; an evaluation report's is the hash of its JSON.
from dataclasses import dataclass, field
import hashlib, json
@dataclass(frozen=True)
class Artifact:
kind: str # "dataset", "feature_table", "model", "eval_report"
uri: str # immutable location, e.g. .../artifacts/<fingerprint>/
fingerprint: str # identity of the content, not of the run
schema: dict = field(default_factory=dict, hash=False, compare=False)
def fingerprint_manifest(files: list[tuple[str, int, str]]) -> str:
"""Fingerprint a dataset from its manifest of (relative_path, size, content_hash).
Cheap to compute, stable across copies, and changes if any file changes."""
h = hashlib.sha256()
for path, size, digest in sorted(files):
h.update(f"{path}\t{size}\t{digest}\n".encode())
return h.hexdigest()
def cache_key(step_name: str, code_digest: str, params: dict,
inputs: dict[str, Artifact], env_digest: str) -> str:
payload = {
"step": step_name,
"code": code_digest, # hash of the step's source or container image digest
"env": env_digest, # pinned dependency set
"params": params, # every knob, including seed and any "as of" date
"inputs": {k: a.fingerprint for k, a in sorted(inputs.items())},
}
return hashlib.sha256(json.dumps(payload, sort_keys=True).encode()).hexdigest()The cache key deliberately includes the environment. A changed library version can change numerical results, and a cache that ignores it will hand back an artifact that the current code would not produce. Use the container image digest or a lock-file hash. Parameters must include everything that affects the output, including the random seed and any date that bounds the data.
The step runner
With those pieces, the runner is short. It computes the key, reuses an existing artifact if one is recorded and present, otherwise checks input contracts, runs the step into a staging directory, seals the output by fingerprinting it and making it read-only, checks the output contract and records the run.
def run_step(step, params, inputs, store, meta):
key = cache_key(step.name, step.code_digest, params, inputs, step.env_digest)
hit = meta.lookup(key)
if hit is not None and store.exists(hit.uri):
return hit # reuse: nothing recomputed
for name, art in inputs.items(): # contract on the way in
step.input_contracts[name].check(art)
staging = store.new_staging_dir()
step.fn(params=params, inputs=inputs, out_dir=staging) # pure: reads only inputs
out = store.seal(staging, kind=step.output_kind) # compute fingerprint, make immutable
step.output_contract.check(out) # contract on the way out
meta.record(key=key, step=step.name, params=params,
inputs={k: a.fingerprint for k, a in inputs.items()},
output=out) # lineage is a by-product
return outTwo properties matter. A step that crashes leaves only a staging directory, never a half-written artifact under a real fingerprint, so retries are safe. And because the record ties a key to its inputs' fingerprints, walking backwards from any model to the raw snapshots that produced it is a graph query, not an investigation. The model registry should store the model's fingerprint so that this walk starts from what is actually serving.
Data contracts at step boundaries
Caching makes a pipeline fast; contracts make it trustworthy. A contract states the schema, which columns may not be null, plausible value ranges, and how much the row count may move compared with the last accepted artifact. The runner checks input contracts before spending compute and output contracts before anything downstream can use the result.
@dataclass
class Contract:
columns: dict[str, str] # name -> type
not_null: set[str]
ranges: dict[str, tuple[float, float]]
max_row_delta: float = 0.2 # vs the previous accepted artifact
severity: str = "fail" # "fail", "quarantine" or "warn"
def check(self, art, previous=None, stats=None):
stats = stats or profile(art) # one pass: types, null counts, min, max, rows
errors = []
for col, typ in self.columns.items():
if stats.types.get(col) != typ:
errors.append(f"{col}: expected {typ}, got {stats.types.get(col)}")
errors += [f"{c}: {stats.nulls[c]} nulls" for c in self.not_null if stats.nulls.get(c)]
for col, (lo, hi) in self.ranges.items():
if stats.min[col] < lo or stats.max[col] > hi:
errors.append(f"{col}: outside [{lo}, {hi}]")
if previous is not None:
delta = abs(stats.rows - previous.rows) / max(previous.rows, 1)
if delta > self.max_row_delta:
errors.append(f"row count moved {delta:.0%}")
if errors:
raise ContractViolation(self.severity, errors)Give each check a severity. Schema breaks should fail the run. A suspicious row-count drop might quarantine the artifact for review while the previous model keeps serving. Mild distribution shifts can warn. Contracts that fail too often get ignored, so start with a few that encode real past incidents and add more as you learn. Contracts complement, rather than replace, the feature-level checks in a feature store.
Reproducibility, honestly
Content addressing guarantees that the same inputs and recipe map to the same recorded artifact. It does not guarantee that re-running training produces identical weights: GPU kernels with non-deterministic reductions, data-loader ordering and mixed precision all introduce run-to-run variance. Treat training as reproducible in distribution. Record seeds and the environment, and when you need to compare, run the recipe a few times and compare metric ranges, not single numbers. The cache still pays off, because the default is to reuse the recorded model rather than retrain it.
Worked example: changing one feature
A churn model's pipeline has six steps: snapshot raw events, validate, clean, build features, train, evaluate. A data scientist changes the window of one feature from 30 to 28 days. The snapshot, validate and clean steps have unchanged keys, because their code, parameters and inputs are identical, so they are skipped. The feature step's parameters changed, so it runs and produces a new feature table with a new fingerprint. That changes the train step's input fingerprint, so training runs, and evaluation follows.
Now suppose instead that someone edits only the evaluation script to add a metric. Only the evaluation step's code digest changes, so v2 reruns evaluation against the existing model in minutes, whereas v1 would have rebuilt features and retrained. Finally, suppose the upstream events table silently loses a day of data. The snapshot fingerprint changes as it should, but the validate step's row-count contract quarantines the artifact before any training compute is spent, and the serving model is untouched. Retraining triggers that sit on top of this are covered in automated retraining.
Migrating from v1 without a big bang
- Inventory. For each v1 script, list what it actually reads and writes, including hidden inputs such as the clock, environment variables and tables it queries by name. This list is usually longer than anyone expects.
- Wrap first, rewrite later. Turn each script into a v2 step with its inputs passed in explicitly and its output written to staging and sealed. The logic stays the same, so any difference in output points at a hidden input.
- Dual-run. Run v1 and v2 side by side on the same schedule, with v1 still feeding production.
- Diff for parity. Compare row counts, schemas and per-column summary statistics of every intermediate artifact, and model metrics within a tolerance you set in advance. Exact equality is expected for deterministic steps; training needs the in-distribution comparison above.
- Cut over one consumer at a time, starting with the least critical model, and keep v1 runnable and frozen for a fixed rollback period.
- Delete v1 once the rollback period passes with no reverts, so there is only one source of truth.
Failure modes
- Cache poisoning by undeclared inputs. A step that still reads the clock or a mutable table has a cache key that does not change when its output should. Stale artifacts are reused silently. Sandbox steps so they can read only staged inputs where you can.
- Over-broad keys. Including a timestamp or run identifier in parameters makes every key unique and the cache useless.
- Under-broad keys. Omitting the environment or seed returns artifacts the current code would not produce.
- Storage growth. Immutable artifacts accumulate. Garbage-collect by reachability: keep anything reachable from a registered or serving model and anything recent, and delete the rest.
- Contract fatigue. Too many warn-level checks get muted. Review the alert rate and delete checks that never catch anything real.
- Metadata store as a single point of failure. If it is down, nothing can run. Back it up and keep the artifact store self-describing, with a small metadata file next to each artifact, so lineage can be rebuilt.
Trade-offs
v2 costs discipline: every step must declare its inputs, and wrapping legacy code takes time. It costs storage, since nothing is overwritten. It moves complexity into the runner and metadata store, which become platform components someone must own. In return, partial reruns become cheap, backfills become a loop over parameter values, incidents become graph queries, and a broken upstream stops at a contract instead of in production. For a single model retrained monthly by one person, v1 with good discipline may be enough; once several models share data and people, v2 usually pays for itself. How it fits a wider release process is described in MLOps architecture.
What to do next
- Pick one production pipeline and write down every input each step reads, including the clock, environment variables and tables queried by name.
- Introduce immutable artifact locations keyed by fingerprint, starting with the model and the final feature table.
- Implement the cache key from code digest, environment digest, parameters and input fingerprints, and make the runner skip steps whose key is recorded.
- Add three contracts that encode past incidents, with explicit severities.
- Store the model fingerprint in the registry and test that you can walk from it to the raw snapshot.
- Dual-run v1 and v2 for at least one full retraining cycle, diff every artifact, then cut over one consumer.
- Add reachability-based garbage collection before storage becomes a reason to abandon immutability.