Datapipelines och ETL
Kunna bygga en reproducerbar pipeline som hämtar, transformerar och laddar data med körhistorik.
Förkunskaper
- DPandas — tabeller i Pythonkrävs
- DSQL — grundernakrävs
Intuition
En pipeline är en kedja av steg som tar rådata till något användbart. Den klassiska uppdelningen:
| Bokstav | Steg | Exempel |
|---|---|---|
| E | Extract | hämta från API, databas, filer |
| T | Transform | städa, slå ihop, beräkna |
| L | Load | skriv till lager eller databas |
ELT i stället för ETL är vanligare i dag: ladda in rådatan först, transformera sedan i lagret. Fördelen är att rådatan finns kvar — ändrar du transformationen behöver du inte hämta allt igen, och du kan gå tillbaka och se vad som faktiskt kom in.
Tre egenskaper som skiljer en pipeline som håller från ett skript som fungerade en gång:
- Idempotent — kör den två gånger och resultatet är detsamma. Inga dubbletter, inga dubbla insättningar.
- Omkörbar per partition — kan köra om gårdagen utan att röra resten.
- Observerbar — varje körning loggar vad som gick in, vad som kom ut och hur lång tid det tog.
Formellt
Idempotens i praktiken. Skriv aldrig INSERT rakt av i en pipeline. Tre mönster som fungerar:
| Mönster | Hur |
|---|---|
| Upsert | INSERT ... ON CONFLICT DO UPDATE på en naturlig nyckel |
| Överskriv partition | radera dagens partition, skriv om den |
| Append-only med körnings-id | lägg alltid till; läs alltid senaste körningen per nyckel |
Den andra är enklast att resonera om: en körning äger sin partition helt, och att köra om den är alltid säkert.
Kvalitetsgrindar mellan stegen. En pipeline utan kontroller upptäcker fel först när någon tittar på en rapport en vecka senare:
| Kontroll | Exempel |
|---|---|
| Radantal inom förväntat spann | 50 000 ± 20 % |
| Inga nycklar saknas | id IS NULL = 0 |
| Unik nyckel | inga dubbletter på id |
| Värdeintervall | alder mellan 0 och 120 |
| Färskhet | senaste raden är från i dag |
| Schema oförändrat | samma kolumner och typer som förra körningen |
Regeln: faila hellre högt och tidigt än att släppa igenom skräp. En pipeline som stannar och larmar är bättre än en som tyst skriver in felaktig data i tusen nedströms-rapporter.
Orkestrering. När stegen blir fler behövs något som kör dem i rätt ordning, hanterar återförsök och visar historik — Airflow, Dagster eller Prefect. Grunden i alla tre är samma: en DAG (riktad acyklisk graf) av uppgifter, precis som beroendegrafen i den här plattformen.
Backfill är det som avgör om designen håller: ska du köra om sex månader historik, en dag i taget, fungerar det då? Om pipelinen läser «allt sedan förra körningen» i stället för «en angiven period» är svaret nej. Gör alltid perioden till en parameter.
Kod
import hashlib, json, time
from dataclasses import dataclass, field
from datetime import date
from pathlib import Path
import pandas as pd
@dataclass
class Kontroll:
namn: str
test: callable
blockerande: bool = True
KONTROLLER = [
Kontroll("radantal rimligt", lambda df: 100 <= len(df) <= 1_000_000),
Kontroll("inga saknade id", lambda df: df["id"].notna().all()),
Kontroll("unika id", lambda df: not df["id"].duplicated().any()),
Kontroll("rimliga åldrar", lambda df: df["alder"].between(0, 120).all()),
Kontroll("få saknade poäng", lambda df: df["poang"].isna().mean() < 0.2, blockerande=False),
]
class Pipelinefel(Exception):
pass
def kontrollera(df, steg):
utfall = []
for k in KONTROLLER:
ok = bool(k.test(df))
utfall.append({"kontroll": k.namn, "ok": ok})
if not ok and k.blockerande:
raise Pipelinefel(f"{steg}: kontrollen «{k.namn}» misslyckades")
return utfall
def kor(dag: date, ut: Path):
"""En körning äger sin partition helt — omkörning är alltid säker."""
t0 = time.perf_counter()
partition = ut / f"dag={dag.isoformat()}"
ra = hamta(dag) # E — alltid för en ANGIVEN dag
kontroller = kontrollera(ra, "extract")
ren = (ra.drop_duplicates(subset=["id"]) # T
.assign(stad=lambda d: d["stad"].str.strip().str.lower())
.query("0 <= alder <= 120"))
partition.mkdir(parents=True, exist_ok=True)
for f in partition.glob("*.parquet"):
f.unlink() # radera partitionen först → idempotent
ren.to_parquet(partition / "data.parquet", index=False) # L
manifest = {
"dag": dag.isoformat(),
"rader_in": len(ra), "rader_ut": len(ren),
"bortfall": len(ra) - len(ren),
"sha256": hashlib.sha256(
(partition / "data.parquet").read_bytes()).hexdigest()[:16],
"kontroller": kontroller,
"sekunder": round(time.perf_counter() - t0, 2),
}
(partition / "_manifest.json").write_text(
json.dumps(manifest, indent=2, ensure_ascii=False), encoding="utf-8")
return manifest
# Backfill: perioden är en parameter, så historik kan köras om dag för dag
def backfill(fran: date, till: date, ut: Path):
d = fran
while d <= till:
print(kor(d, ut))
d = date.fromordinal(d.toordinal() + 1)
_manifest.json per partition är det som gör pipelinen granskningsbar: radantal in och ut, bortfall, checksumma och vilka kontroller som kördes. Frågan «varför saknas 3 000 rader den 14 mars?» blir då möjlig att besvara.
Behärskning innebär
- Bygger en pipeline med tydliga steg
- Gör stegen idempotenta och omkörbara
- Loggar körhistorik och datakvalitet
Logga in för att göra övningarna och bygga upp din behärskning.
Källor
- 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