Salta el contingut

Projecte: pipeline amb validació completa i alertes

Aquesta pàgina tanca l'apartat de qualitat de dades lligant totes les peces anteriors en un únic pipeline end-to-end: ingesta → perfilat → validació → quarantena → informe de qualitat → alerta. És exactament l'arquitectura que qualsevol pipeline de producció hauria de tenir, i la que es demana implementar a la miniactivitat d'aquesta pàgina.

flowchart LR
    A["1. Ingesta\n(CSV / API / BD)"] --> B["2. Perfilat\n(data profiling)"]
    B --> C["3. Validació\n(Great Expectations\no Pandera)"]
    C -->|"vàlid"| D["4. Càrrega al\nData Warehouse"]
    C -->|"invàlid"| E["4b. Quarantena"]
    D --> F["5. Informe\nde qualitat"]
    E --> F
    F -->|"si hi ha\ndegradació"| G["6. Alerta\n(Slack / email)"]

    style A fill:#1d4ed8,color:#ffffff,stroke:#3b82f6
    style B fill:#7c3aed,color:#ffffff,stroke:#a78bfa
    style C fill:#b45309,color:#ffffff,stroke:#f59e0b
    style D fill:#166534,color:#ffffff,stroke:#22c55e
    style E fill:#7f1d1d,color:#ffffff,stroke:#ef4444
    style F fill:#0e7490,color:#ffffff,stroke:#06b6d4
    style G fill:#9d174d,color:#ffffff,stroke:#ec4899

Esquelet del pipeline complet

import logging
from dataclasses import dataclass, field
from datetime import datetime

import pandas as pd
import pandera as pa
from pandera import Column, DataFrameSchema, Check

logging.basicConfig(level=logging.INFO)


@dataclass
class InformeQualitat:
    dataset: str
    timestamp: datetime = field(default_factory=datetime.now)
    total_registres: int = 0
    registres_valids: int = 0
    registres_quarantena: int = 0
    dimensions: dict = field(default_factory=dict)

    @property
    def pct_error(self) -> float:
        if self.total_registres == 0:
            return 0.0
        return self.registres_quarantena / self.total_registres * 100

    def to_markdown(self) -> str:
        linies = [
            f"# Informe de qualitat — {self.dataset}",
            f"Generat: {self.timestamp:%Y-%m-%d %H:%M}",
            "",
            f"- Total de registres processats: {self.total_registres}",
            f"- Registres vàlids: {self.registres_valids}",
            f"- Registres en quarantena: {self.registres_quarantena} ({self.pct_error:.1f}%)",
            "",
            "## Dimensions mesurades",
        ]
        for dim, valor in self.dimensions.items():
            linies.append(f"- **{dim}**: {valor}")
        return "\n".join(linies)


def pas_1_ingesta(ruta: str) -> pd.DataFrame:
    logging.info("Ingerint dades de %s", ruta)
    return pd.read_csv(ruta)


def pas_2_perfilat(df: pd.DataFrame) -> dict:
    return {
        "completesa_mitjana": f"{(1 - df.isnull().mean().mean()) * 100:.1f}%",
        "files": len(df),
        "columnes": len(df.columns),
    }


def pas_3_valida(df: pd.DataFrame, esquema: DataFrameSchema) -> tuple[pd.DataFrame, pd.DataFrame]:
    try:
        df_valid = esquema.validate(df, lazy=True)
        return df_valid, pd.DataFrame()
    except pa.errors.SchemaErrors as e:
        indexs_invalids = e.failure_cases["index"].dropna().unique().astype(int)
        df_invalid = df.loc[df.index.isin(indexs_invalids)].copy()
        df_valid = df.drop(index=indexs_invalids, errors="ignore")
        return df_valid, df_invalid


def pas_4b_quarantena(df_invalid: pd.DataFrame, ruta_sortida: str) -> None:
    if not df_invalid.empty:
        df_invalid.to_csv(ruta_sortida, index=False)
        logging.warning("Registres en quarantena guardats a %s", ruta_sortida)


def pas_6_alerta_si_cal(informe: InformeQualitat, llindar_pct: float = 5.0) -> None:
    if informe.pct_error > llindar_pct:
        logging.error(
            "ALERTA: %.1f%% de registres en quarantena (llindar: %.1f%%). Cal revisió.",
            informe.pct_error, llindar_pct,
        )
    else:
        logging.info("Qualitat dins del llindar acceptable (%.1f%%).", informe.pct_error)


def executa_pipeline(ruta_dades: str, esquema: DataFrameSchema, nom_dataset: str) -> InformeQualitat:
    df = pas_1_ingesta(ruta_dades)
    dimensions = pas_2_perfilat(df)
    df_valid, df_invalid = pas_3_valida(df, esquema)
    pas_4b_quarantena(df_invalid, f"quarantena_{nom_dataset}.csv")

    informe = InformeQualitat(
        dataset=nom_dataset,
        total_registres=len(df),
        registres_valids=len(df_valid),
        registres_quarantena=len(df_invalid),
        dimensions=dimensions,
    )
    pas_6_alerta_si_cal(informe)

    with open(f"informe_qualitat_{nom_dataset}.md", "w", encoding="utf-8") as f:
        f.write(informe.to_markdown())

    return informe

Aquest esquelet és deliberadament simple (una funció per pas, sense orquestador), perquè el propòsit és entendre l'arquitectura lògica. En un entorn real, cada pas seria una tasca d'un DAG d'Airflow (vegeu la pàgina d'Airflow d'aquest bloc), amb els seus propis reintents i dependències.


Checklist d'un pipeline de qualitat complet

  • Ingesta: es registra d'on venen les dades i quan s'han rebut.
  • Perfilat: es genera com a mínim un resum de completesa i unicitat abans de validar (útil per detectar sorpreses abans de definir regles massa estrictes).
  • Validació: hi ha un esquema explícit (Pandera o Great Expectations) que cobreix com a mínim completesa, unicitat, validesa i rangs plausibles.
  • Quarantena: els registres invàlids es guarden per a revisió, no es descarten silenciosament ni aturen tot el pipeline per un sol registre dolent.
  • Informe: es genera un resum llegible (Markdown, HTML o un registre a una taula de mètriques) amb el percentatge d'error i les dimensions mesurades.
  • Alerta: si el percentatge d'error supera un llindar (idealment dinàmic, vegeu la pàgina d'alertes), s'envia una notificació abans que ningú de negoci detecti el problema.
  • Contracte: si el dataset es comparteix amb un altre equip, hi ha un contracte de dades versionat que documenta l'schema i les garanties de servei.

AC5074/05/04 — Miniactivitat: implementació de checks de qualitat sobre un dataset real

Tria un dataset públic real (per exemple, un dataset de Kaggle o dades obertes de l'Ajuntament/Generalitat amb com a mínim 500 files i 6 columnes) i implementa el pipeline complet d'aquesta pàgina:

  1. Perfilat inicial: genera un informe de completesa i unicitat per columna (pots fer servir perfil_basic() de la pàgina de dimensions, o ydata-profiling).
  2. Defineix un esquema de validació (Pandera o Great Expectations, a triar) que cobreixi com a mínim 3 de les 5 dimensions de qualitat vistes en aquest bloc (completesa, unicitat, validesa, consistència, frescor), amb almenys 6 regles diferents en total.
  3. Implementa la quarantena: separa els registres invàlids en un fitxer a part, amb el motiu de l'error.
  4. Genera un informe de qualitat en Markdown o HTML amb: total de registres, percentatge vàlid/invàlid, i el detall de quines regles han fallat més sovint.
  5. Implementa una alerta (pot ser simulada: escriure a un fitxer de log amb nivell ERROR, no cal Slack real) que s'activi si el percentatge d'error supera un 10%.
  6. Redacta un contracte de dades breu (com el de la pàgina de contractes) per al dataset triat: schema, propietari fictici, i almenys una regla de "canvi trencador" específica per a aquest dataset.

Lliurament: codi font, els dos fitxers de sortida (quarantena i informe), i el contracte en YAML. Documenta al codi, amb comentaris, quina dimensió de qualitat cobreix cada regla de validació.


Mòdul M5074 Sistemes de Big Data | Institut Sa Palomera (Blanes) | Curs CEIABD 2026-2027