Data pipelines and ETL
Be able to build a reproducible pipeline that extracts, transforms and loads data with a run history.
Prerequisites
- DPandas — tables in Pythonrequired
- DSQL — the basicsrequired
Intuition
A pipeline is a chain of steps that takes raw data to something useful. The classic division:
| Letter | Step | Example |
|---|---|---|
| E | Extract | fetch from an API, a database, files |
| T | Transform | clean, join, compute |
| L | Load | write 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:
- Idempotent — run it twice and the result is the same. No duplicates, no double insertions.
- Re-runnable per partition — you can re-run yesterday without touching the rest.
- 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:
| Pattern | How |
|---|---|
| Upsert | INSERT ... ON CONFLICT DO UPDATE on a natural key |
| Overwrite the partition | delete today's partition, write it again |
| Append-only with a run id | always 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:
| Check | Example |
|---|---|
| The row count within the expected range | 50 000 ± 20 % |
| No missing keys | id IS NULL = 0 |
| A unique key | no duplicates on id |
| The value range | age between 0 and 120 |
| Freshness | the latest row is from today |
| The schema unchanged | the 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
- pandas — User Guide (BSD-3) — BSD-3-Clause
- Apache Airflow — dokumentation (Apache-2.0) — Apache-2.0
- Great Expectations (Apache-2.0) — Apache-2.0