Skip to content
AI-grafen
EUniversityData handling· about 60 min· evolving, reviewed regularly· verified 2026-09-20· EN

Data pipelines and ETL

Be able to build a reproducible pipeline that extracts, transforms and loads data with a run history.

Prerequisites

Intuition

A pipeline is a chain of steps that takes raw data to something useful. The classic division:

LetterStepExample
EExtractfetch from an API, a database, files
TTransformclean, join, compute
LLoadwrite to a warehouse or a database

ELT instead of ETL is more common today: load the raw data in first, transform it in the warehouse afterwards. The advantage is that the raw data remains — change the transformation and you do not have to fetch everything again, and you can go back and see what actually came in.

Three properties that separate a pipeline that lasts from a script that worked once:

  1. Idempotent — run it twice and the result is the same. No duplicates, no double insertions.
  2. Re-runnable per partition — you can re-run yesterday without touching the rest.
  3. Observable — every run logs what went in, what came out and how long it took.

Formal

Idempotence in practice. Never write a plain INSERT in a pipeline. Three patterns that work:

PatternHow
UpsertINSERT ... ON CONFLICT DO UPDATE on a natural key
Overwrite the partitiondelete today's partition, write it again
Append-only with a run idalways append; always read the latest run per key

The second is the easiest to reason about: a run owns its partition entirely, and re-running it is always safe.

Quality gates between the steps. A pipeline without checks discovers errors only when somebody looks at a report a week later:

CheckExample
The row count within the expected range50 000 ± 20 %
No missing keysid IS NULL = 0
A unique keyno duplicates on id
The value rangeage between 0 and 120
Freshnessthe latest row is from today
The schema unchangedthe same columns and types as the previous run

The rule: fail loudly and early rather than letting rubbish through. A pipeline that stops and raises an alert is better than one that quietly writes incorrect data into a thousand downstream reports.

Orchestration. When the steps get more numerous you need something that runs them in the right order, handles retries and shows the history — Airflow, Dagster or Prefect. The basis of all three is the same: a DAG (a directed acyclic graph) of tasks, exactly like the dependency graph in this platform.

A backfill is what decides whether the design holds: if you have to re-run six months of history, one day at a time, does it work? If the pipeline reads «everything since the last run» instead of «a given period» the answer is no. Always make the period a parameter.

Code

import hashlib, json, time
from dataclasses import dataclass, field
from datetime import date
from pathlib import Path
import pandas as pd

@dataclass
class Check:
    name: str
    test: callable
    blocking: bool = True

CHECKS = [
    Check("a reasonable row count", lambda df: 100 <= len(df) <= 1_000_000),
    Check("no missing ids", lambda df: df["id"].notna().all()),
    Check("unique ids", lambda df: not df["id"].duplicated().any()),
    Check("reasonable ages", lambda df: df["age"].between(0, 120).all()),
    Check("few missing scores", lambda df: df["score"].isna().mean() < 0.2, blocking=False),
]

class PipelineError(Exception):
    pass

def check(df, step):
    outcome = []
    for c in CHECKS:
        ok = bool(c.test(df))
        outcome.append({"check": c.name, "ok": ok})
        if not ok and c.blocking:
            raise PipelineError(f"{step}: the check «{c.name}» failed")
    return outcome

def run(day: date, out: Path):
    """A run owns its partition entirely — re-running is always safe."""
    t0 = time.perf_counter()
    partition = out / f"day={day.isoformat()}"

    raw = fetch(day)                         # E — always for a GIVEN day
    checks = check(raw, "extract")

    clean = (raw.drop_duplicates(subset=["id"])   # T
                .assign(city=lambda d: d["city"].str.strip().str.lower())
                .query("0 <= age <= 120"))

    partition.mkdir(parents=True, exist_ok=True)
    for f in partition.glob("*.parquet"):
        f.unlink()                            # delete the partition first → idempotent
    clean.to_parquet(partition / "data.parquet", index=False)   # L

    manifest = {
        "day": day.isoformat(),
        "rows_in": len(raw), "rows_out": len(clean),
        "dropped": len(raw) - len(clean),
        "sha256": hashlib.sha256(
            (partition / "data.parquet").read_bytes()).hexdigest()[:16],
        "checks": checks,
        "seconds": round(time.perf_counter() - t0, 2),
    }
    (partition / "_manifest.json").write_text(
        json.dumps(manifest, indent=2, ensure_ascii=False), encoding="utf-8")
    return manifest

# A backfill: the period is a parameter, so history can be re-run day by day
def backfill(start: date, end: date, out: Path):
    d = start
    while d <= end:
        print(run(d, out))
        d = date.fromordinal(d.toordinal() + 1)

A _manifest.json per partition is what makes the pipeline auditable: the rows in and out, what was dropped, a checksum and which checks were run. The question «why are 3 000 rows missing on 14 March?» then becomes answerable.

Mastery means

  • Builds a pipeline with clear steps
  • Makes the steps idempotent and re-runnable
  • Logs the run history and the data quality

Sign in to do the exercises and build your mastery up.

Sources

All the sources and licences