analitica ja té totes les peces de càlcul: el llac a HDFS que pujar_esdeveniments_hdfs.py alimenta (04-02), vendes_diaries.py a Spark (05-03), les recomanacions amb ALS, i els fluxos de Flink que no necessiten que ningú els llanci perquè no acaben mai (05-04). El que falta és allò que fa que els lots passin cada dia sense que ningú els llanci a mà: esperar que el fitxer del dia sigui complet, comprovar que no està corromput, executar l'spark-submit, carregar el resultat a la base de dades que consulten els panells, avisar si alguna cosa falla, reintentar sense duplicar, i tornar-ho a fer per a set dies quan es corregeix un preu. La versió inicial d'això a Quilòmetre Zero era un crontab amb quatre línies i un script llancar_tot.sh, i va fallar de totes les maneres possibles: el fitxer va arribar tard i Spark va processar mig dia, un reintent va carregar les vendes dues vegades, i ningú no se'n va adonar fins que un productor va preguntar per què la seva gràfica s'havia doblat. Aquesta lliçó tracta el pipeline de dades com el que és, un sistema distribuït més: un DAG de tasques amb dependències, planificació per temps i per dada, reintents, idempotència, backfill, SLA i qualitat de dades; presenta Apache Airflow com a planificador de referència (i esmenta Prefect, Dagster i Argo Workflows); i construeix dags/vendes_diaries.py, el pipeline diari d'analitica, que tanca el mòdul.

Contingut

  1. D'scripts solts a pipelines
  2. Els conceptes: DAG de tasques, planificació, reintents, idempotència, backfill, SLA
  3. cron i els seus límits
  4. Apache Airflow: arquitectura i model de programació
  5. Alternatives: Prefect, Dagster i Argo Workflows
  6. Qualitat de dades, llinatge i catàleg
  7. Pràctica: dags/vendes_diaries.py
  8. Backfill de la Setmana de la Verema
  9. Errors Comuns i Consells
  10. Exercicis
  11. Conclusió

  1. D'scripts solts a pipelines

El crontab original d'analitica era aquest:

# Pujar els esdeveniments del dia anterior al llac, calcular vendes, carregar a PostgreSQL
0 2 * * *  cd /opt/km0 && python serveis/analitica/pujar_esdeveniments_hdfs.py $(date -d yesterday +\%F)
0 3 * * *  cd /opt/km0 && spark-submit serveis/analitica/vendes_diaries.py hdfs://namenode:8020/km0/esdeveniments/$(date -d yesterday +\%F)/comandes.jsonl ...
0 4 * * *  cd /opt/km0 && python serveis/analitica/carregar_postgres.py $(date -d yesterday +\%F)

Cada línia és correcta, i el conjunt és fràgil per raons que ja coneixem de la resta del curs:

  • Les dependències són implícites i temporals. Spark arrenca a les 3:00 suposant que la pujada de les 2:00 ha acabat. El dia que la pujada triga 70 minuts (un pic de campanya), Spark processa un fitxer a mitges i la càrrega de les 4:00 publica vendes incompletes. És la fal·làcia de la latència zero (01-04) aplicada al temps d'un treball.
  • No hi ha reintents, o n'hi ha sense idempotència. Si carregar_postgres.py falla a mig camí, algú el rellança a mà i les files es dupliquen, perquè fa INSERT sense clau.
  • Reprocessar és manual. Corregir set dies és editar set dates a mà en tres ordres, en ordre, sense equivocar-se.
  • No hi ha observabilitat. cron envia un correu amb la sortida estàndard si el procés retorna un codi diferent de zero, i res més. Quant va trigar ahir? Es va executar el dia 12? Quina tasca falla més?
  • L'estat viu al crontab d'una màquina. Si aquesta màquina mor (fal·làcia 1), no hi ha pipeline; si dues persones l'editen, ningú no sap quina versió corre.

Un pipeline de dades resol això fent explícit l'implícit: les tasques i les seves dependències formen un graf, cada tasca declara quan es pot executar (per temps, o quan la dada que necessita existeix), què fer si falla i quant pot trigar, i el sistema que l'executa registra cada execució i permet rellançar qualsevol tram. És paral·lelisme de tasques (05-01, apartat 2) amb un planificador que coneix el DAG complet.

  1. Els conceptes: DAG de tasques, planificació, reintents, idempotència, backfill, SLA

El DAG de tasques. Les tasques són nodes i les dependències, arestes dirigides; acíclic perquè una tasca no pot dependre de si mateixa. A diferència del DAG d'operadors de Spark (05-03), aquí cada node és un treball complet (un spark-submit, una càrrega a base de dades, una crida HTTP) i el planificador no mou dades entre nodes: els nodes es comuniquen per l'emmagatzematge (HDFS, PostgreSQL) i només intercanvien metadades petites (quantes files, quina ruta). Les tasques sense dependència mútua s'executen en paral·lel; el temps total el marca el camí crític.

Planificació per temps i per dada. "A les 3:00" és planificació per temps, i és insuficient quan l'entrada la produeix un altre sistema. La planificació per dada (o per esdeveniment) hi afegeix sensors: tasques que no fan res més que esperar que una condició sigui certa (existeix el fitxer /km0/esdeveniments/2026-09-14/comandes.jsonl; hi ha una fila en una taula; un altre DAG ha acabat; un missatge ha arribat a una cua) i que fallen si la condició no es compleix dins d'un termini. El DAG arrenca a les 3:00 però Spark no comença fins que el sensor confirma que el fitxer hi és.

Reintents amb backoff. Moltes fallades són transitòries (un NameNode en failover, PostgreSQL saturat, un timeout). Cada tasca declara quantes vegades reintentar i amb quina espera, creixent (backoff exponencial amb sostre i jitter, la mateixa política que 02-05 donava als consumidors). Les fallades deterministes (un error al codi) esgoten els reintents i aleshores sí que fallen.

Idempotència de cada tasca. Els reintents i els reprocessaments només són segurs si executar una tasca dues vegades per al mateix dia deixa el mateix resultat que una. La tècnica universal és escriure per partició amb sobreescriptura: la tasca del dia D produeix exactament la partició D de la seva sortida (el directori dia=2026-09-14/ en Parquet, les files amb dia = '2026-09-14' a PostgreSQL) i la reemplaça sencera, dins d'una transacció o amb un reanomenament atòmic. Mai append, mai INSERT sense esborrar abans o sense ON CONFLICT. vendes_diaries.py ja ho fa amb partitionOverwriteMode=dynamic (05-03); la càrrega a PostgreSQL ho farà amb DELETE ... WHERE dia = %s seguit d'INSERT dins de la mateixa transacció.

Backfill. Executar el pipeline per a un rang de dates passades: perquè s'ha corregit una entrada (els preus de la Setmana de la Verema), perquè s'ha canviat la lògica (una columna nova que cal omplir històricament), o perquè el pipeline ha estat aturat. Només és possible si cada execució està parametritzada per la seva data lògica (no per "ahir") i si les tasques són idempotents. Un pipeline que fa servir date -d yesterday no admet backfill.

SLA i alertes. Un SLA (service level agreement) del pipeline és "les vendes del dia D són a PostgreSQL abans de les 6:00 del dia D+1". El planificador vigila el termini i avisa quan una tasca o el DAG l'incompleixen, a més d'avisar a cada fallada definitiva. Les alertes van al canal de l'equip; les fallades, a qui estigui de guàrdia (07-01 tractarà el monitoratge en general).

Versionat del pipeline. El DAG és codi al repositori (km0/dags/), revisat i desplegat com la resta: se sap quina versió va córrer cada dia, i un canvi de lògica és un commit, no una edició del crontab. Idealment, la versió del codi queda registrada al costat de cada execució.

  1. cron i els seus límits

cron continua sent l'eina correcta per a "executar aquesta ordre a aquesta hora en aquesta màquina" quan no hi ha dependències ni estat que gestionar: rotar logs, una còpia de seguretat, un recordatori. Com a orquestrador de pipelines, la comparació és aquesta:

cron Airflow (i similars)
Dependències entre treballs Cap: se simulen amb hores separades Explícites, com a DAG; una tasca arrenca quan les seves predecessores han acabat bé
Espera per dades No (cal programar un bucle a l'script) Sensors, amb timeout i reprogramació
Reintents No Per tasca, amb backoff
Parametrització per data A mà (date -d yesterday) Data lògica de cada execució (ds), disponible a totes les tasques
Backfill Manual, ordre a ordre airflow dags backfill -s ... -e ...
Historial i estat El correu de cron, si de cas Base de metadades: cada execució, cada intent, durada, logs
Interfície crontab -e Web amb el DAG, l'estat per dia, logs, rellançar tasques
Alertes i SLA No Callbacks de fallada, SLA miss, integracions
Alta disponibilitat La màquina del crontab Scheduler replicable, workers distribuïts
Concurrència i quotes No (dos crons es poden encavalcar) max_active_runs, pools, depends_on_past
Cost Zero Un servei més per operar (base de dades, scheduler, workers)

  1. Apache Airflow: arquitectura i model de programació

Airflow (Airbnb, 2014; Apache des del 2016) defineix els pipelines com a codi Python i els executa amb aquesta arquitectura:

flowchart LR
    DEV[Repositori<br/>km0/dags/*.py] --> S
    S[Scheduler<br/>parseja els DAG, decideix quines<br/>tasques toca executar] --> EX[Executor<br/>Local · Celery · Kubernetes]
    EX --> W1[Worker 1<br/>executa tasques]
    EX --> W2[Worker 2]
    S <--> DB[(Base de metadades<br/>PostgreSQL: DAG runs,<br/>task instances, XCom)]
    W1 <--> DB
    W2 <--> DB
    WEB[Webserver<br/>interficie, API] <--> DB
    W1 --> HDFS[(HDFS)]
    W1 --> SP[Spark]
    W2 --> PG[(km0_analitica)]
  • Scheduler. El cor. Parseja periòdicament els fitxers de dags/, calcula per a cada DAG quines execucions (DAG runs) toquen segons el seu schedule i el seu start_date, i per a cada execució quines tasques tenen les dependències complertes; les encua a l'executor. Se'n pot executar més d'un per a alta disponibilitat.
  • Executor. Com s'executen les tasques: LocalExecutor (processos a la màquina de l'scheduler: suficient per a desenvolupament i pipelines modestos), CeleryExecutor (una cua, Redis o RabbitMQ de 02-04, i workers en diverses màquines), KubernetesExecutor (un pod per tasca, 07-05).
  • Workers. Executen el codi de cada tasca. Per a tasques que llancen feina en un altre sistema (Spark, una consulta), el worker només espera i supervisa; el càlcul pesant no passa a Airflow.
  • Base de metadades. PostgreSQL amb l'estat de tot: quins DAG runs existeixen, en quin estat és cada instància de tasca, quants intents, els XCom. És la font de veritat, i per això Airflow no perd el pipeline si un worker mor.
  • Webserver. La interfície: el graf, la graella d'execucions per dia, logs per intent, i els botons de rellançar, marcar com a èxit o netejar.

El model de programació:

  • DAG. Un objecte DAG amb id, schedule (una expressió cron, un timedelta, o un dataset per a planificació per dada), start_date, catchup i default_args per a les tasques.
  • Operadors i tasques. Cada tasca és una instància d'un operador: BashOperator, PythonOperator, SparkSubmitOperator, SQLExecuteQueryOperator, centenars més als providers. Els sensors són operadors que esperen: FileSensor, WebHdfsSensor, ExternalTaskSensor, SqlSensor. Les dependències es declaren amb >>.
  • Data lògica i plantilles. Cada DAG run té una data lògica (logical_date, abans execution_date) i un interval de dades (data_interval_start/end). Amb un schedule diari, el run que processa les dades del 14 de setembre té data lògica 2026-09-14 i s'executa quan acaba l'interval, és a dir, el dia 15 a l'hora programada. La macro {{ ds }} a qualsevol camp de plantilla val 2026-09-14 per a aquest run: és el que fa possible el backfill, perquè una execució per al dia 8 tindrà ds = 2026-09-08 encara que es llanci a l'octubre.
  • catchup. Si és True, en activar un DAG amb start_date en el passat, l'scheduler crea i executa un run per cada interval no executat des d'aleshores. És útil per poblar l'històric i perillós si no s'espera (centenars de runs de cop). Amb False, només s'executa des de l'interval actual, i l'històric es fa amb un backfill explícit.
  • depends_on_past. La tasca d'un run no arrenca fins que la mateixa tasca del run anterior hagi acabat bé. Necessari quan cada dia es construeix sobre l'anterior (un acumulat); innecessari i perjudicial quan els dies són independents, perquè una fallada del dia 12 bloqueja el 13, el 14...
  • retries, retry_delay, retry_exponential_backoff, max_retry_delay. La política de reintents per tasca.
  • XCom. Un mecanisme perquè una tasca deixi un valor petit (un nombre de files, una ruta) i una altra el llegeixi, a través de la base de metadades. Per a metadades, no per a dades: qualsevol cosa més gran d'uns KB va a l'emmagatzematge i per XCom hi viatja la seva ruta.
  • Pools. Quotes de concurrència amb nom: un pool spark amb 2 places garanteix que mai no hi hagi més de dos spark-submit alhora encara que vint runs de backfill els demanin.
  • max_active_runs, concurrency, trigger_rule (per defecte all_success: la tasca arrenca si totes les anteriors han tingut èxit; all_done, one_failed... per a tasques de neteja o notificació).

  1. Alternatives: Prefect, Dagster i Argo Workflows

Airflow és l'estàndard de fet, amb el seu pes: una base de dades, un scheduler, una interfície que envelleix, i un model (DAG runs per interval de temps) pensat per a lots diaris. Les alternatives n'ataquen els punts febles:

Eina Idea central Quan encaixa
Prefect Fluxos en Python normal (decoradors @flow, @task), execució dinàmica, sense la rigidesa de l'interval; desplegament híbrid (orquestració al núvol, execució a la teva infraestructura) Equips Python que volen menys cerimònia; pipelines amb lògica dinàmica
Dagster Orientat a actius de dades (software-defined assets): es declara quines taules i fitxers existeixen i de què depenen, i l'orquestrador en deriva les tasques; tipatge, proves i llinatge integrats Plataformes de dades amb molts conjunts derivats; es vol el catàleg i el llinatge des del principi
Argo Workflows DAG de contenidors a Kubernetes, definits en YAML; cada pas és un pod Tot ja corre a Kubernetes; pipelines de ML i CI amb contenidors; sense Python obligatori
Cron + scripts Res Un treball sense dependències en una màquina

L'elecció per a Quilòmetre Zero és Airflow per maduresa, pels operadors de Spark, HDFS i PostgreSQL que ja existeixen, i perquè l'equip el coneix; Dagster seria l'alternativa seriosa si la plataforma de dades creixés fins a desenes de conjunts derivats. Els conceptes (DAG, data lògica, idempotència per partició, sensors, backfill) són els mateixos en tots quatre.

  1. Qualitat de dades, llinatge i catàleg

Un pipeline que executa bé un càlcul sobre dades dolentes produeix resultats dolents a temps. La qualitat de dades es verifica dins del pipeline, com a tasques:

  • Validacions d'entrada, abans de gastar còmput: el fitxer existeix i està tancat (no l'està escrivint ningú), té una mida plausible (un dia de campanya amb 200 línies és sospitós), un mostreig es parseja com a JSON, els camps obligatoris hi són, la fracció de línies corruptes és inferior a un llindar, les dates dels esdeveniments cauen en el dia esperat.
  • Validacions de sortida, abans de publicar: el total d'import coincideix amb el control del job (05-03 imprimia un Control: per aquesta raó), no hi ha productors desconeguts, no hi ha imports negatius, el nombre de files és dins del rang històric (±50 % del dia equivalent de la setmana anterior).
  • Quarantena. Un fitxer d'entrada que no passa la validació no es processa a mitges ni es descarta: es mou a /km0/quarantena/<dia>/ amb un informe del motiu, la tasca falla amb un missatge clar, i algú decideix. Les dades vàlides d'aquell dia es poden reprocessar amb backfill quan es corregeixi l'entrada.

Eines com Great Expectations o Soda expressen aquestes comprovacions de manera declarativa i s'integren a Airflow; per al pipeline d'aquesta lliçó n'hi haurà prou amb una tasca Python. Dos conceptes més que només esmentem: el llinatge (de quines entrades i amb quin codi es va produir cada sortida: vendes_diaries/dia=2026-09-14 ve d'esdeveniments/2026-09-14/comandes.jsonl i de cataleg.csv amb la versió a3f9 de vendes_diaries.py; OpenLineage ho estandarditza i Airflow ho emet) i el catàleg de dades (l'inventari de quins conjunts existeixen, el seu esquema, el seu propietari i la seva frescor: el Hive Metastore de 05-02, DataHub, Amundsen). Tots dos són el que permet respondre "d'on surt aquest número?" sense llegir codi.

  1. Pràctica: dags/vendes_diaries.py

7.1 Airflow a docker-compose.yml

Airflow s'afegeix al docker-compose.yml de km0/ amb la seva pròpia base de metadades i LocalExecutor (suficient aquí; en producció, Celery o Kubernetes):

# km0/docker-compose.yml (fragment)
services:
  airflow-db:
    image: postgres:16
    environment: { POSTGRES_USER: airflow, POSTGRES_PASSWORD: airflow, POSTGRES_DB: airflow }
  airflow: &airflow
    image: apache/airflow:2.9.3-python3.11
    environment:
      AIRFLOW__CORE__EXECUTOR: LocalExecutor
      AIRFLOW__DATABASE__SQL_ALCHEMY_CONN: postgresql+psycopg2://airflow:airflow@airflow-db/airflow
      AIRFLOW__CORE__LOAD_EXAMPLES: "false"
      AIRFLOW__CORE__DEFAULT_TIMEZONE: Europe/Madrid
      # Connexions als sistemes de Quilòmetre Zero, com a URI (evita crear-les a mà a la interfície)
      AIRFLOW_CONN_HDFS_KM0: http://namenode:9870
      AIRFLOW_CONN_SPARK_KM0: spark://spark-master:7077
      AIRFLOW_CONN_KM0_ANALITICA: postgresql://analitica:analitica@postgres-analitica:5432/km0_analitica
      _PIP_ADDITIONAL_REQUIREMENTS: apache-airflow-providers-apache-spark apache-airflow-providers-apache-hdfs hdfs pyarrow
    volumes:
      - ./dags:/opt/airflow/dags
      - ./serveis/analitica:/app/analitica
    command: webserver
    ports: ["8090:8080"]
  airflow-scheduler:
    <<: *airflow
    command: scheduler
    ports: []

Després de docker compose run airflow airflow db migrate i de crear un usuari, la interfície és a localhost:8090. La taula de destinació a km0_analitica, amb la clau que fa idempotent la càrrega:

-- km0/sql/analitica/vendes_diaries.sql
CREATE TABLE IF NOT EXISTS vendes_diaries (
    dia              date        NOT NULL,
    mercat           text        NOT NULL,
    productor        text        NOT NULL,
    nom_productor    text,
    provincia        text,
    import_eur       numeric(12,2) NOT NULL,
    unitats          integer     NOT NULL,
    comandes         integer     NOT NULL,
    carregat_en      timestamptz NOT NULL DEFAULT now(),
    PRIMARY KEY (dia, mercat, productor)
);

7.2 El DAG

flowchart LR
    S[esperar_esdeveniments<br/>WebHdfsSensor<br/>/km0/esdeveniments/ds/comandes.jsonl] --> V[validar_entrada<br/>PythonOperator<br/>mostreig, mida, dates]
    V --> SP[calcular_vendes<br/>SparkSubmitOperator<br/>vendes_diaries.py]
    SP --> C[carregar_postgres<br/>PythonOperator<br/>DELETE + INSERT per dia]
    C --> Q[validar_sortida<br/>PythonOperator<br/>total = control]
    Q --> N[notificar<br/>trigger_rule = all_done]
    V -. fallada: quarantena .-> N
# km0/dags/vendes_diaries.py
"""Pipeline diari d'analítica: esdeveniments del llac -> vendes per productor/mercat/dia -> km0_analitica.

Cada run processa el dia {{ ds }} (la data lògica) i és idempotent: es pot rellançar o fer-ne backfill.
"""
from datetime import datetime, timedelta
import json

from airflow import DAG
from airflow.operators.python import PythonOperator
from airflow.providers.apache.hdfs.sensors.web_hdfs import WebHdfsSensor
from airflow.providers.apache.spark.operators.spark_submit import SparkSubmitOperator
from airflow.providers.postgres.hooks.postgres import PostgresHook
from airflow.exceptions import AirflowFailException

HDFS = "hdfs://namenode:8020"
RUTA_ESDEVENIMENTS = "/km0/esdeveniments/{ds}/comandes.jsonl"
RUTA_SORTIDA = "/km0/agregats/vendes_diaries"
MIN_LINIES = 1000                 # per sota d'això, el dia és sospitós
MAX_CORRUPTES = 0.01              # 1 % de línies illegibles com a màxim


def client_hdfs():
    from hdfs import InsecureClient                       # WebHDFS, com pujar_esdeveniments_hdfs.py de 04-02
    return InsecureClient("http://namenode:9870", user="analitica")


def validar_entrada(ds, ti, **_):
    """Comprova el fitxer del dia abans de gastar un job de Spark. El mou a quarantena si no passa."""
    ruta = RUTA_ESDEVENIMENTS.format(ds=ds)
    hdfs = client_hdfs()
    estat = hdfs.status(ruta)
    total, corruptes, fora_de_dia = 0, 0, 0
    with hdfs.read(ruta, encoding="utf-8") as f:
        for linia in f:
            total += 1
            try:
                ev = json.loads(linia)
                dia_ev = datetime.utcfromtimestamp(ev["data_ms"] / 1000).strftime("%Y-%m-%d")
                if dia_ev != ds:
                    fora_de_dia += 1
            except (json.JSONDecodeError, KeyError, TypeError):
                corruptes += 1
    problemes = []
    if total < MIN_LINIES:
        problemes.append(f"només {total} línies (mínim {MIN_LINIES})")
    if total and corruptes / total > MAX_CORRUPTES:
        problemes.append(f"{corruptes} línies corruptes de {total}")
    if problemes:
        desti = f"/km0/quarantena/{ds}/comandes.jsonl"
        hdfs.makedirs(f"/km0/quarantena/{ds}")
        hdfs.rename(ruta, desti)                          # el fitxer no es perd: algú el revisarà
        hdfs.write(f"/km0/quarantena/{ds}/informe.txt", "\n".join(problemes), overwrite=True)
        raise AirflowFailException(f"Entrada de {ds} en quarantena: " + "; ".join(problemes))   # sense reintents
    ti.xcom_push(key="linies", value=total)               # metadada petita per a les tasques següents
    ti.xcom_push(key="fora_de_dia", value=fora_de_dia)
    print(f"{ds}: {total} línies, {corruptes} corruptes, {fora_de_dia} d'un altre dia, {estat['length'] / 1e6:.1f} MB")


def carregar_postgres(ds, ti, **_):
    """Carrega la partició dia=ds del Parquet a km0_analitica reemplaçant el dia complet: idempotent."""
    import pyarrow.parquet as pq
    from pyarrow import fs
    hdfs_fs = fs.HadoopFileSystem("namenode", 8020)
    taula = pq.read_table(f"{RUTA_SORTIDA}/dia={ds}", filesystem=hdfs_fs).to_pylist()
    files = [(ds, r["mercat"], r["productor"], r["nom_productor"], r["provincia"],
              r["import_eur"], r["unitats"], r["comandes"]) for r in taula]
    pg = PostgresHook(postgres_conn_id="km0_analitica")
    with pg.get_conn() as conn, conn.cursor() as cur:       # UNA transacció: esborrar + inserir, o res
        cur.execute("DELETE FROM vendes_diaries WHERE dia = %s", (ds,))
        cur.executemany("""INSERT INTO vendes_diaries
                           (dia, mercat, productor, nom_productor, provincia, import_eur, unitats, comandes)
                           VALUES (%s, %s, %s, %s, %s, %s, %s, %s)""", files)
        conn.commit()
    ti.xcom_push(key="files", value=len(files))
    ti.xcom_push(key="total", value=float(sum(r["import_eur"] for r in taula)))


def validar_sortida(ds, ti, **_):
    """La suma carregada ha de coincidir amb la que Spark va calcular, i el dia ha de tenir una mida plausible."""
    pg = PostgresHook(postgres_conn_id="km0_analitica")
    total_pg, files_pg = pg.get_first("SELECT COALESCE(SUM(import_eur), 0), COUNT(*) FROM vendes_diaries WHERE dia = %s", (ds,))
    if abs(float(total_pg) - ti.xcom_pull(task_ids="carregar_postgres", key="total")) > 0.01:
        raise AirflowFailException(f"Total a PostgreSQL {total_pg} != total carregat")
    setmana_abans = pg.get_first("SELECT COALESCE(SUM(import_eur), 0) FROM vendes_diaries WHERE dia = %s::date - 7", (ds,))[0]
    if setmana_abans and not 0.5 <= float(total_pg) / float(setmana_abans) <= 2.0:
        print(f"AVÍS: el total {total_pg} es desvia més del 50 % del de fa una setmana ({setmana_abans})")
    print(f"{ds}: {files_pg} files, total {total_pg} €")


def notificar(ds, dag_run, **_):
    """Resum del run al canal de l'equip. S'executa sempre (trigger_rule=all_done)."""
    estats = {ti.task_id: ti.state for ti in dag_run.get_task_instances()}
    fallides = [t for t, s in estats.items() if s == "failed"]
    missatge = f"vendes_diaries {ds}: " + ("OK" if not fallides else f"FALLADA a {', '.join(fallides)}")
    print(missatge)                                       # aquí hi aniria el webhook de Slack/Teams o un correu


default_args = {
    "owner": "analitica",
    "retries": 3,
    "retry_delay": timedelta(minutes=5),
    "retry_exponential_backoff": True,                    # 5, 10, 20 min...
    "max_retry_delay": timedelta(minutes=30),
    "depends_on_past": False,                             # els dies són independents
    "email_on_failure": False,                            # les alertes van per 'notificar' i el callback
    "sla": timedelta(hours=3),                            # cada tasca ha d'acabar 3 h després de l'inici del run
}

with DAG(
    dag_id="km0_vendes_diaries",
    description="Vendes per productor, mercat i dia a partir del llac d'esdeveniments",
    schedule="0 3 * * *",                                 # a les 3:00, per a les dades del dia anterior
    start_date=datetime(2026, 9, 1),
    catchup=False,                                        # l'històric es fa amb backfill explícit
    max_active_runs=1,                                    # un dia alhora en execució normal
    default_args=default_args,
    tags=["km0", "analitica", "lots"],
) as dag:

    esperar_esdeveniments = WebHdfsSensor(
        task_id="esperar_esdeveniments",
        webhdfs_conn_id="hdfs_km0",
        filepath=RUTA_ESDEVENIMENTS.format(ds="{{ ds }}"),  # plantilla: la data lògica del run
        poke_interval=300,                                # comprovar cada 5 minuts
        timeout=6 * 3600,                                 # rendir-se (i fallar) a les 9:00
        mode="reschedule",                                # allibera el worker entre comprovacions
    )

    validar = PythonOperator(task_id="validar_entrada", python_callable=validar_entrada)

    calcular_vendes = SparkSubmitOperator(
        task_id="calcular_vendes",
        conn_id="spark_km0",
        application="/app/analitica/vendes_diaries.py",   # el programa de 05-03, sense canvis
        application_args=[
            f"{HDFS}{RUTA_ESDEVENIMENTS.format(ds='{{ ds }}')}",
            "/app/analitica/cataleg.csv",
            f"{HDFS}{RUTA_SORTIDA}",
        ],
        conf={"spark.sql.shuffle.partitions": "8",
              "spark.sql.sources.partitionOverwriteMode": "dynamic"},   # només la partició dia={{ ds }}
        executor_memory="1G", executor_cores=2, num_executors=2,
        name="km0-vendes-{{ ds }}",
        pool="spark",                                     # com a màxim 2 jobs Spark alhora (pool creat a la UI/CLI)
    )

    carregar = PythonOperator(task_id="carregar_postgres", python_callable=carregar_postgres)
    comprovar = PythonOperator(task_id="validar_sortida", python_callable=validar_sortida)
    avis = PythonOperator(task_id="notificar", python_callable=notificar, trigger_rule="all_done", retries=0)

    esperar_esdeveniments >> validar >> calcular_vendes >> carregar >> comprovar >> avis

Els punts que convé fixar:

  • La data lògica ho governa tot. {{ ds }} al sensor i als arguments de Spark, ds com a paràmetre de les funcions Python. El run del 15 de setembre a les 3:00 té ds = 2026-09-14 i processa el fitxer del 14: és l'interval de dades que s'acaba de tancar. Res al DAG no diu "ahir".
  • El sensor en mode reschedule no ocupa una plaça de l'executor mentre espera: es reprograma cada 5 minuts. Amb mode="poke" (per defecte) el worker quedaria bloquejat sis hores. El timeout converteix "el fitxer no ha arribat" en una fallada visible a les 9:00 en comptes d'un pipeline penjat.
  • validar_entrada falla amb AirflowFailException, que no reintenta: un fitxer corromput no s'arregla esperant cinc minuts, i el fitxer ja s'ha mogut a quarantena. La resta de fallades (una excepció de xarxa en llegir HDFS) sí que reintenten segons default_args.
  • Idempotència a les tres escriptures. Spark sobreescriu només dia={{ ds }} (05-03); carregar_postgres esborra i insereix el dia dins d'una transacció; la quarantena fa servir rename, atòmic a HDFS. Rellançar qualsevol tasca, o el run sencer, deixa el mateix estat.
  • XCom per a metadades. linies, files, total: números, no dades. Les dades viatgen per HDFS i PostgreSQL.
  • pool="spark" limita a dos els jobs Spark simultanis encara que un backfill creï set runs alhora (apartat 8), i max_active_runs=1 manté l'execució normal en un dia cada vegada.
  • notificar amb trigger_rule="all_done" s'executa tant si el run ha anat bé com si alguna tasca ha fallat, i per això pot informar de l'estat; sense aquesta regla, una fallada a validar_entrada la deixaria en upstream_failed i ningú no se n'assabentaria. L'SLA de default_args hi afegeix una alerta si alguna tasca continua sense acabar tres hores després de les 3:00.

En activar el DAG a la interfície, l'scheduler crea el primer run a la següent execució programada (amb catchup=False), i la graella mostra un quadrat per dia i tasca, verd, vermell o groc (reintentant). Un clic en un quadrat dona els logs d'aquell intent, inclosa la sortida d'spark-submit amb l'explain() de 05-03.

  1. Backfill de la Setmana de la Verema

El cas de l'exercici 3 de 05-03, ara amb el pipeline: comandes corregeix el preu de vi-crianca i regenera els fitxers comandes.jsonl del 8 al 14 de setembre a HDFS. Cal recalcular aquests set dies, i només aquests, sense tocar la resta i sense interferir amb el run diari. Amb el DAG parametritzat per data lògica i idempotent, és una ordre:

# --reset-dagruns: els runs d'aquests dies ja existeixen (amb èxit); netejar-los i reexecutar
docker compose exec airflow-scheduler airflow dags backfill km0_vendes_diaries \
  --start-date 2026-09-08 --end-date 2026-09-14 \
  --reset-dagruns --rerun-failed-tasks

El que passa: l'scheduler crea (o reinicia) set DAG runs amb ds del 08 al 14 i els executa respectant les dependències de cadascun i els límits globals: el pool spark permet dos calcular_vendes alhora, de manera que els set jobs Spark s'executen en quatre tandes; max_active_runs no s'aplica al backfill (té el seu propi límit, --max-active-runs a les versions recents, o el max_active_runs del DAG segons la versió: convé comprovar-ho, perquè un backfill d'un any amb set runs alhora pot saturar HDFS). Els sensors passen immediatament (els fitxers existeixen), les validacions es repeteixen sobre els fitxers corregits, Spark sobreescriu dia=2026-09-08/ ... dia=2026-09-14/ i deixa la resta de dies intactes, i cada càrrega esborra i insereix el seu dia. En acabar, la graella mostra els set dies amb un intent nou en verd, i SELECT dia, SUM(import_eur) FROM vendes_diaries WHERE dia BETWEEN '2026-09-08' AND '2026-09-14' GROUP BY 1 reflecteix els preus corregits.

Dues variants útils: per reexecutar només des de la càrrega (Spark ja ha corregut bé, ha fallat PostgreSQL), airflow tasks clear km0_vendes_diaries -t carregar_postgres --downstream -s 2026-09-08 -e 2026-09-14 neteja aquesta tasca i les següents d'aquests runs, i l'scheduler les reexecuta; i per a un run puntual amb una data concreta, airflow dags trigger km0_vendes_diaries --logical-date 2026-09-14. Cap de les tres coses no era possible amb el crontab de l'apartat 1 sense editar ordres a mà.

Errors Comuns i Consells

  • Fer servir "ahir" en comptes de la data lògica. datetime.now() o date -d yesterday dins d'una tasca fan impossible el backfill i produeixen resultats diferents segons quan s'executi. Sempre {{ ds }} / data_interval_start.
  • Confondre quan s'executa un run amb quines dades processa. El run amb data lògica 14 s'executa el dia 15. És la font de més confusió d'Airflow; pensar en "l'interval que s'acaba de tancar" ho aclareix.
  • Tasques no idempotents. INSERT sense esborrar abans, append a un fitxer, un comptador. Reintentar duplica; el backfill acumula. Partició completa amb sobreescriptura, sempre.
  • Càlcul pesant dins del worker d'Airflow. Un PythonOperator que llegeix 150 MB de JSON i agrega amb pandas converteix el worker en un node de còmput mal dimensionat. Airflow orquestra; Spark, la base de dades o un contenidor calculen.
  • Dades per XCom. Un DataFrame serialitzat a la base de metadades. XCom és per a rutes i comptadors.
  • catchup=True sense voler. Activar un DAG amb start_date de fa dos anys i catchup per defecte llança 730 runs. catchup=False i backfill explícit.
  • depends_on_past=True per precaució. Un dia fallit bloqueja tots els següents fins que algú ho arregla. Només quan cada dia depèn de debò de l'anterior.
  • Sensors en mode poke amb timeouts llargs. Cada sensor ocupa una plaça de l'executor durant hores; deu sensors esperant bloquegen el pipeline sencer. mode="reschedule", o planificació per dataset.
  • Reintentar errors deterministes. Un fitxer corromput o un bug es reintenten tres vegades amb backoff i fallen una hora després. AirflowFailException per a allò que no s'arregla esperant.
  • Sense validació de sortida. El pipeline en verd no vol dir dades correctes. Un control de totals i un rang plausible respecte de l'històric costen vint línies i eviten gràfiques doblades.
  • Secrets al DAG. Contrasenyes de PostgreSQL al codi. Connexions d'Airflow, variables d'entorn o un gestor de secrets (06-04).

Exercicis

Exercici 1: Afegir les recomanacions al pipeline

Amplia el DAG perquè, després de validar_sortida, entreni les recomanacions amb recomanacions_als.py (05-03) sobre els clics dels últims 7 dies i carregui el resultat a la taula recomanacions(client, producte, afinitat, generat_en) de km0_cataleg. Decideix: quin sensor necessita? De quines tasques depèn? Com fas idempotent la càrrega si la taula no té una "partició per dia" natural? Ha de fer servir depends_on_past? Què passa amb les recomanacions durant un backfill de set dies?

Exercici 2: Un dia dolent

El 20 de setembre a les 3:00 passa això: pujar_esdeveniments_hdfs.py ha tingut una fallada i el fitxer /km0/esdeveniments/2026-09-19/comandes.jsonl arriba a HDFS a les 7:40 amb 180 000 línies, de les quals 4 200 no són JSON vàlid. Descriu, tasca a tasca, què fa el DAG de l'apartat 7 entre les 3:00 i les 9:00: estats, reintents, on acaba el fitxer, què veu l'equip. Després, algú corregeix el fitxer i el torna a pujar a les 11:30. Quina ordre executa l'equip per completar el dia 19, i què passa amb el run del dia 20 aquella nit?

Exercici 3: cron o Airflow

Per a cadascun d'aquests treballs de Quilòmetre Zero, decideix si el deixaries a cron, el posaries a Airflow, o el trauries a un flux de 05-04, i justifica-ho en una frase: (a) esborrar els checkpoints de Flink de més de 30 dies a HDFS; (b) generar cada dilluns l'informe setmanal PDF per productor a partir de vendes_diaries, i enviar-lo per correu; (c) recalcular l'estoc mínim per producte i mercat cada 5 minuts; (d) exportar a MinIO una còpia de km0_inventari cada nit; (e) carregar cada dia a PostgreSQL les posicions tardanes de repartiment.posicions que el flux va desviar a la sortida lateral.

Solucions

Exercici 1.

Tasques noves: esperar_clics (un WebHdfsSensor sobre /km0/clics/{{ ds }}/ o, millor, sobre un fitxer de tancament _TANCAT que el procés de pujada per hores escrigui en acabar el dia; sense ell, el directori existeix des de la primera hora i el sensor passaria amb dades incompletes), entrenar_als (SparkSubmitOperator amb recomanacions_als.py i un argument amb el rang {{ macros.ds_add(ds, -6) }} a {{ ds }}), carregar_recomanacions (PythonOperator). Dependències: esperar_clics >> entrenar_als >> carregar_recomanacions, i entrenar_als també aigües avall de validar_sortida només si les recomanacions fan servir les vendes (no és el cas: fan servir clics), així que en rigor són dues branques paral·leles que conflueixen a notificar. Idempotència sense partició per dia: la taula es reemplaça sencera dins d'una transacció (DELETE FROM recomanacions; INSERT ...; COMMIT), o s'escriu a recomanacions_nova i es fa ALTER TABLE ... RENAME intercanviant totes dues (atòmic a PostgreSQL), o es desa una columna generat_en = ds i la web llegeix WHERE generat_en = (SELECT MAX(generat_en) ...): la tercera opció és la que permet backfill sense trepitjar la versió actual. depends_on_past: no; cada entrenament és independent. Durant un backfill de set dies, les set execucions entrenarien set models amb finestres mòbils de clics i carregarien set vegades: amb la columna generat_en és innocu però inútil; el raonable és que entrenar_als visqui en un altre DAG amb el seu propi cicle (diari, sense backfill llevat que es demani), o que al backfill s'exclogui amb --task-regex la branca de recomanacions.

Exercici 2.

3:00: l'scheduler crea el run ds=2026-09-19. esperar_esdeveniments comprova cada 5 minuts (reschedule, sense ocupar worker) i no troba el fitxer: estat up_for_reschedule, quadrat groc. 6:00: l'SLA de 3 h s'incompleix i Airflow registra un SLA miss amb avís a l'equip (el pipeline encara no ha fallat, però va tard). 7:40: el fitxer apareix; a la comprovació de les 7:45 el sensor passa a success. validar_entrada arrenca: 180 000 línies (> 1 000), però 4 200 corruptes són el 2,3 % (> 1 %): mou el fitxer a /km0/quarantena/2026-09-19/comandes.jsonl, escriu informe.txt i llança AirflowFailException: estat failed sense reintents. calcular_vendes, carregar_postgres i validar_sortida queden en upstream_failed. notificar s'executa (all_done) i publica "vendes_diaries 2026-09-19: FALLADA a validar_entrada"; el callback de fallada també avisa. A les 9:00 no passa res més: el timeout del sensor ja no s'aplica perquè el sensor ha acabat. L'equip veu el quadrat vermell a validar_entrada, llegeix el log amb "4200 línies corruptes de 180000" i troba el fitxer a quarantena.

11:30: el fitxer corregit es puja a /km0/esdeveniments/2026-09-19/comandes.jsonl. L'equip executa airflow tasks clear km0_vendes_diaries -t validar_entrada --downstream -s 2026-09-19 -e 2026-09-19 (o prem Clear a la tasca des de la interfície): validar_entrada i les següents tornen a None i l'scheduler les reexecuta; el sensor no es repeteix perquè no s'ha netejat. Amb max_active_runs=1, aquest run ocupa la plaça fins que acaba (uns minuts). Aquella nit, a les 3:00 del dia 21, el run ds=2026-09-20 es crea amb normalitat: depends_on_past=False, així que no l'afecta el que ha passat amb el 19, i el 19 ja és en verd.

Exercici 3.

(a) cron (o el mateix Flink amb retenció de checkpoints configurada): una ordre sense dependències ni dades, hdfs dfs -rm amb una data; si se salta un dia, no passa res. (b) Airflow: depèn que vendes_diaries hagi carregat els set dies (un ExternalTaskSensor sobre el DAG diari), té una sortida que cal poder regenerar (backfill d'una setmana si es corregeixen dades) i un enviament que no s'ha de duplicar. (c) Flux (05-04): cada 5 minuts amb latència de segons i estat per clau és una finestra tumbling sobre estoc.actualitzat, no un lot llançat 288 vegades al dia. (d) Airflow, encara que sigui una sola tasca: es vol historial, alerta si falla i un sensor o validació que la còpia s'ha completat (mida, _SUCCESS); a cron seria acceptable si s'hi afegís monitoratge extern. (e) Airflow: és un lot diari parametritzat per data que llegeix /km0/repartiment/tardanes/{{ ds }}/ i carrega per partició, i encaixa com una tasca més del pipeline diari: és la reconciliació entre el flux i el lot de què parlava l'exercici 1 de 05-04.

Conclusió

Un pipeline de dades és la part del sistema distribuït que converteix treballs solts en una plataforma: un DAG de tasques amb dependències explícites, planificat per temps i per dada (sensors), amb reintents amb backoff per a allò transitori i fallada immediata per a allò determinista, tasques idempotents que escriuen particions completes amb sobreescriptura, backfill parametritzat per la data lògica i no per "ahir", SLA i notificacions per saber quan alguna cosa va tard, validacions d'entrada i de sortida amb quarantena per no publicar dades dolentes a temps, i tot al repositori, versionat. cron no ofereix res d'això; Airflow ho ofereix amb un scheduler, un executor, workers, una base de metadades i DAG en Python, i Prefect, Dagster i Argo Workflows ho reformulen amb èmfasis diferents. dags/vendes_diaries.py encadena el sensor sobre /km0/esdeveniments/{{ ds }}/comandes.jsonl, la validació, l'SparkSubmitOperator amb el programa de 05-03, la càrrega transaccional a km0_analitica i la comprovació de totals, i el backfill de la Setmana de la Verema es redueix a una ordre amb dues dates.

Amb això es tanca el Mòdul 5, i la plataforma de dades de Quilòmetre Zero queda completa:

Peça Què és Lliçó
Llac de dades HDFS amb un directori per dia: /km0/esdeveniments/<dia>/comandes.jsonl, /km0/clics/<dia>/, convertits a Parquet particionat per dia a /km0/agregats/ 04-02, 05-03
Models de còmput Portar el càlcul a les dades; scatter/gather, cues de treball, BSP, dataflow; shuffle, biaix, ressagats, reexecució determinista 05-01
Lots MapReduce com a base històrica i vocabulari; Spark amb DataFrames, Catalyst i Parquet per a vendes_diaries.py; MLlib/ALS per a recomanacions 05-02, 05-03
Fluxos Flink (i Structured Streaming) sobre comandes.esdeveniments i repartiment.posicions: temps d'esdeveniment, marques d'aigua, finestres, checkpoints, sinks idempotents; alertes d'estoc i panell de repartiment 05-04
Pipelines Airflow: km0_vendes_diaries amb sensors, validació, Spark, càrrega idempotent, backfill 05-05

Les dades estan repartides (Mòdul 4) i es processen en massa, per lots i en temps real, sense que ningú no llanci res a mà. Però hi ha una cosa que hem donat per feta en tots els mòduls fins ara: qualsevol procés podia parlar amb qualsevol servei i llegir qualsevol dada. vendes_diaries.py llegeix tot el llac; carregar_postgres entra a km0_analitica amb una contrasenya en una variable d'entorn; el consumidor de Flink llegeix comandes.esdeveniments sense identificar-se; un hdfs dfs -rm des de qualsevol contenidor esborraria un any d'esdeveniments; i les URL presignades de MinIO (04-03) han estat l'únic mecanisme d'accés que hem dissenyat amb cura. En un sistema amb desenes de serveis, centenars de treballs i dades personals de l'Anna, en Marc i la Llúcia, això no pot continuar així. El Mòdul 6 tracta la seguretat en sistemes distribuïts, i comença per les dues preguntes que cap sistema no pot eludir: qui és qui (autenticació) i qui pot fer què (autorització).

Curs d'Arquitectures Distribuïdes

Mòdul 1: Introducció als Sistemes Distribuïts

Mòdul 2: Comunicació en Sistemes Distribuïts

Mòdul 3: Consistència i Replicació

Mòdul 4: Emmagatzematge Distribuït

Mòdul 5: Computació Distribuïda

Mòdul 6: Seguretat en Sistemes Distribuïts

Mòdul 7: Monitoratge i Manteniment

Mòdul 8: Casos d'Estudi i Aplicacions

© Copyright 2026. Tots els drets reservats