Hoppa till innehållet
AI-grafen
E· Universitetdatahantering· ca 60 min· utvecklande· verifierad 2026-09-20

Datapipelines och ETL

Kunna bygga en reproducerbar pipeline som hämtar, transformerar och laddar data med körhistorik.

Förkunskaper

Intuition

En pipeline är en kedja av steg som tar rådata till något användbart. Den klassiska uppdelningen:

BokstavStegExempel
EExtracthämta från API, databas, filer
TTransformstäda, slå ihop, beräkna
LLoadskriv 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:

  1. Idempotent — kör den två gånger och resultatet är detsamma. Inga dubbletter, inga dubbla insättningar.
  2. Omkörbar per partition — kan köra om gårdagen utan att röra resten.
  3. 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önsterHur
UpsertINSERT ... ON CONFLICT DO UPDATE på en naturlig nyckel
Överskriv partitionradera dagens partition, skriv om den
Append-only med körnings-idlä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:

KontrollExempel
Radantal inom förväntat spann50 000 ± 20 %
Inga nycklar saknasid IS NULL = 0
Unik nyckelinga dubbletter på id
Värdeintervallalder mellan 0 och 120
Färskhetsenaste raden är från i dag
Schema oförändratsamma 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

Alla källor och licenser