analitica tiene ya todas las piezas de cálculo: el lago en HDFS que subir_eventos_hdfs.py alimenta (04-02), ventas_diarias.py en Spark (05-03), las recomendaciones con ALS, y los flujos de Flink que no necesitan que nadie los lance porque nunca terminan (05-04). Lo que falta es lo que hace que los lotes ocurran cada día sin que nadie los lance a mano: esperar a que el fichero del día esté completo, comprobar que no está corrupto, ejecutar el spark-submit, cargar el resultado en la base de datos que consultan los paneles, avisar si algo falla, reintentar sin duplicar, y volver a hacerlo para siete días cuando se corrige un precio. La versión inicial de eso en Kilómetro Cero era un crontab con cuatro líneas y un script lanzar_todo.sh, y falló de todas las formas posibles: el fichero llegó tarde y Spark procesó medio día, un reintento cargó las ventas dos veces, y nadie se enteró hasta que un productor preguntó por qué su gráfica se había doblado. Esta lección trata el pipeline de datos como lo que es, un sistema distribuido más: un DAG de tareas con dependencias, planificación por tiempo y por dato, reintentos, idempotencia, backfill, SLAs y calidad de datos; presenta Apache Airflow como planificador de referencia (y nombra Prefect, Dagster y Argo Workflows); y construye dags/ventas_diarias.py, el pipeline diario de analitica, que cierra el módulo.
Contenido
- De scripts sueltos a pipelines
- Los conceptos: DAG de tareas, planificación, reintentos, idempotencia, backfill, SLAs
- cron y sus límites
- Apache Airflow: arquitectura y modelo de programación
- Alternativas: Prefect, Dagster y Argo Workflows
- Calidad de datos, linaje y catálogo
- Práctica:
dags/ventas_diarias.py - Backfill de la Semana de la Vendimia
- Errores Comunes y Consejos
- Ejercicios
- Conclusión
- De scripts sueltos a pipelines
El crontab original de analitica era este:
# Subir eventos del día anterior al lago, calcular ventas, cargar en PostgreSQL
0 2 * * * cd /opt/km0 && python servicios/analitica/subir_eventos_hdfs.py $(date -d yesterday +\%F)
0 3 * * * cd /opt/km0 && spark-submit servicios/analitica/ventas_diarias.py hdfs://namenode:8020/km0/eventos/$(date -d yesterday +\%F)/pedidos.jsonl ...
0 4 * * * cd /opt/km0 && python servicios/analitica/cargar_postgres.py $(date -d yesterday +\%F)Cada línea es correcta, y el conjunto es frágil por razones que ya conocemos del resto del curso:
- Las dependencias son implícitas y temporales. Spark arranca a las 3:00 suponiendo que la subida de las 2:00 terminó. El día que la subida tarda 70 minutos (un pico de campaña), Spark procesa un fichero a medias y la carga de las 4:00 publica ventas incompletas. Es la falacia de la latencia cero (01-04) aplicada al tiempo de un trabajo.
- No hay reintentos, o los hay sin idempotencia. Si
cargar_postgres.pyfalla a mitad, alguien lo relanza a mano y las filas se duplican, porque haceINSERTsin clave. - Reprocesar es manual. Corregir siete días es editar siete fechas a mano en tres comandos, en orden, sin equivocarse.
- No hay observabilidad. cron envía un correo con la salida estándar si el proceso devuelve un código distinto de cero, y nada más. ¿Cuánto tardó ayer? ¿Se ejecutó el día 12? ¿Qué tarea falla más?
- El estado vive en el
crontabde una máquina. Si esa máquina muere (falacia 1), no hay pipeline; si dos personas lo editan, nadie sabe qué versión corre.
Un pipeline de datos resuelve esto haciendo explícito lo implícito: las tareas y sus dependencias forman un grafo, cada tarea declara cuándo puede ejecutarse (por tiempo, o cuando el dato que necesita existe), qué hacer si falla y cuánto puede tardar, y el sistema que lo ejecuta registra cada ejecución y permite relanzar cualquier tramo. Es paralelismo de tareas (05-01, apartado 2) con un planificador que conoce el DAG completo.
- Los conceptos: DAG de tareas, planificación, reintentos, idempotencia, backfill, SLAs
El DAG de tareas. Las tareas son nodos y las dependencias aristas dirigidas; acíclico porque una tarea no puede depender de sí misma. A diferencia del DAG de operadores de Spark (05-03), aquí cada nodo es un trabajo completo (un spark-submit, una carga en base de datos, una llamada HTTP) y el planificador no mueve datos entre nodos: los nodos se comunican por el almacenamiento (HDFS, PostgreSQL) y solo intercambian metadatos pequeños (cuántas filas, qué ruta). Las tareas sin dependencia mutua se ejecutan en paralelo; el tiempo total lo marca el camino crítico.
Planificación por tiempo y por dato. "A las 3:00" es planificación por tiempo, y es insuficiente cuando la entrada la produce otro sistema. La planificación por dato (o por evento) añade sensores: tareas que no hacen nada salvo esperar a que una condición sea cierta (existe el fichero /km0/eventos/2026-09-14/pedidos.jsonl; hay una fila en una tabla; otro DAG ha terminado; un mensaje llegó a una cola) y que fallan si la condición no se cumple en un plazo. El DAG arranca a las 3:00 pero Spark no empieza hasta que el sensor confirma que el fichero está.
Reintentos con backoff. Muchos fallos son transitorios (un NameNode en failover, PostgreSQL saturado, un timeout). Cada tarea declara cuántas veces reintentar y con qué espera, creciente (backoff exponencial con tope y jitter, la misma política que 02-05 daba a los consumidores). Los fallos deterministas (un error en el código) agotan los reintentos y entonces sí fallan.
Idempotencia de cada tarea. Los reintentos y los reprocesamientos solo son seguros si ejecutar una tarea dos veces para el mismo día deja el mismo resultado que una. La técnica universal es escribir por partición con sobreescritura: la tarea del día D produce exactamente la partición D de su salida (el directorio dia=2026-09-14/ en Parquet, las filas con dia = '2026-09-14' en PostgreSQL) y la reemplaza entera, dentro de una transacción o con un renombrado atómico. Nunca append, nunca INSERT sin borrar antes o sin ON CONFLICT. ventas_diarias.py ya lo hace con partitionOverwriteMode=dynamic (05-03); la carga en PostgreSQL lo hará con DELETE ... WHERE dia = %s seguido de INSERT en la misma transacción.
Backfill. Ejecutar el pipeline para un rango de fechas pasadas: porque se corrigió una entrada (los precios de la Semana de la Vendimia), porque se cambió la lógica (una columna nueva que hay que rellenar históricamente), o porque el pipeline estuvo parado. Solo es posible si cada ejecución está parametrizada por su fecha lógica (no por "ayer") y si las tareas son idempotentes. Un pipeline que usa date -d yesterday no admite backfill.
SLAs y alertas. Un SLA (service level agreement) del pipeline es "las ventas del día D están en PostgreSQL antes de las 6:00 del día D+1". El planificador vigila el plazo y avisa cuando una tarea o el DAG lo incumplen, además de avisar en cada fallo definitivo. Las alertas van al canal del equipo; los fallos, a quien esté de guardia (07-01 tratará la monitorización en general).
Versionado del pipeline. El DAG es código en el repositorio (km0/dags/), revisado y desplegado como el resto: se sabe qué versión corrió cada día, y un cambio de lógica es un commit, no una edición del crontab. Idealmente, la versión del código queda registrada junto a cada ejecución.
- cron y sus límites
cron sigue siendo la herramienta correcta para "ejecutar este comando a esta hora en esta máquina" cuando no hay dependencias ni estado que gestionar: rotar logs, un backup, un recordatorio. Como orquestador de pipelines, la comparación es esta:
| cron | Airflow (y similares) | |
|---|---|---|
| Dependencias entre trabajos | Ninguna: se simulan con horas separadas | Explícitas, como DAG; una tarea arranca cuando sus predecesoras han terminado bien |
| Espera por datos | No (hay que programar un bucle en el script) | Sensores, con timeout y reprogramación |
| Reintentos | No | Por tarea, con backoff |
| Parametrización por fecha | A mano (date -d yesterday) |
Fecha lógica de cada ejecución (ds), disponible en todas las tareas |
| Backfill | Manual, comando a comando | airflow dags backfill -s ... -e ... |
| Historial y estado | El correo de cron, si acaso | Base de metadatos: cada ejecución, cada intento, duración, logs |
| Interfaz | crontab -e |
Web con el DAG, el estado por día, logs, relanzar tareas |
| Alertas y SLAs | No | Callbacks de fallo, SLA miss, integraciones |
| Alta disponibilidad | La máquina del crontab | Scheduler replicable, workers distribuidos |
| Concurrencia y cuotas | No (dos crons pueden solaparse) | max_active_runs, pools, depends_on_past |
| Coste | Cero | Un servicio más que operar (base de datos, scheduler, workers) |
- Apache Airflow: arquitectura y modelo de programación
Airflow (Airbnb, 2014; Apache desde 2016) define los pipelines como código Python y los ejecuta con esta arquitectura:
flowchart LR
DEV[Repositorio<br/>km0/dags/*.py] --> S
S[Scheduler<br/>parsea los DAGs, decide qué<br/>tareas toca ejecutar] --> EX[Executor<br/>Local · Celery · Kubernetes]
EX --> W1[Worker 1<br/>ejecuta tareas]
EX --> W2[Worker 2]
S <--> DB[(Base de metadatos<br/>PostgreSQL: DAG runs,<br/>task instances, XCom)]
W1 <--> DB
W2 <--> DB
WEB[Webserver<br/>interfaz, API] <--> DB
W1 --> HDFS[(HDFS)]
W1 --> SP[Spark]
W2 --> PG[(km0_analitica)]
- Scheduler. El corazón. Parsea periódicamente los ficheros de
dags/, calcula para cada DAG qué ejecuciones (DAG runs) tocan según suscheduley sustart_date, y para cada ejecución qué tareas tienen las dependencias cumplidas; las encola en el executor. Se puede ejecutar más de uno para alta disponibilidad. - Executor. Cómo se ejecutan las tareas:
LocalExecutor(procesos en la máquina del scheduler: suficiente para desarrollo y pipelines modestos),CeleryExecutor(una cola, Redis o RabbitMQ de 02-04, y workers en varias máquinas),KubernetesExecutor(un pod por tarea, 07-05). - Workers. Ejecutan el código de cada tarea. Para tareas que lanzan trabajo en otro sistema (Spark, una consulta), el worker solo espera y supervisa; el cálculo pesado no ocurre en Airflow.
- Base de metadatos. PostgreSQL con el estado de todo: qué DAG runs existen, en qué estado está cada instancia de tarea, cuántos intentos, los XCom. Es la fuente de verdad, y por eso Airflow no pierde el pipeline si un worker muere.
- Webserver. La interfaz: el grafo, la rejilla de ejecuciones por día, logs por intento, y los botones de relanzar, marcar como éxito o limpiar.
El modelo de programación:
- DAG. Un objeto
DAGcon id,schedule(una expresión cron, untimedelta, o un dataset para planificación por dato),start_date,catchupydefault_argspara las tareas. - Operadores y tareas. Cada tarea es una instancia de un operador:
BashOperator,PythonOperator,SparkSubmitOperator,SQLExecuteQueryOperator, cientos más en los providers. Los sensores son operadores que esperan:FileSensor,WebHdfsSensor,ExternalTaskSensor,SqlSensor. Las dependencias se declaran con>>. - Fecha lógica y plantillas. Cada DAG run tiene una fecha lógica (
logical_date, antesexecution_date) y un intervalo de datos (data_interval_start/end). Con unschedulediario, el run que procesa los datos del 14 de septiembre tiene fecha lógica 2026-09-14 y se ejecuta al terminar el intervalo, es decir, el 15 a la hora programada. La macro{{ ds }}en cualquier campo de plantilla vale2026-09-14para ese run: es lo que hace posible el backfill, porque una ejecución para el día 8 tendráds = 2026-09-08aunque se lance en octubre. catchup. Si esTrue, al activar un DAG constart_dateen el pasado, el scheduler crea y ejecuta un run por cada intervalo no ejecutado desde entonces. Es útil para poblar histórico y peligroso si no se espera (cientos de runs de golpe). ConFalse, solo se ejecuta desde el intervalo actual, y el histórico se hace con un backfill explícito.depends_on_past. La tarea de un run no arranca hasta que la misma tarea del run anterior haya terminado bien. Necesario cuando cada día se construye sobre el anterior (un acumulado); innecesario y perjudicial cuando los días son independientes, porque un fallo del día 12 bloquea el 13, el 14...retries,retry_delay,retry_exponential_backoff,max_retry_delay. La política de reintentos por tarea.- XCom. Un mecanismo para que una tarea deje un valor pequeño (un número de filas, una ruta) y otra lo lea, a través de la base de metadatos. Para metadatos, no para datos: cualquier cosa mayor de unos KB va al almacenamiento y por XCom viaja su ruta.
- Pools. Cuotas de concurrencia con nombre: un pool
sparkcon 2 huecos garantiza que nunca haya más de dosspark-submita la vez aunque veinte runs de backfill los pidan. max_active_runs,concurrency,trigger_rule(por defectoall_success: la tarea arranca si todas las anteriores tuvieron éxito;all_done,one_failed... para tareas de limpieza o notificación).
- Alternativas: Prefect, Dagster y Argo Workflows
Airflow es el estándar de hecho, con su peso: una base de datos, un scheduler, una interfaz que envejece, y un modelo (DAG runs por intervalo de tiempo) pensado para lotes diarios. Las alternativas atacan sus puntos débiles:
| Herramienta | Idea central | Cuándo encaja |
|---|---|---|
| Prefect | Flujos en Python normal (decoradores @flow, @task), ejecución dinámica, sin la rigidez de intervalo; despliegue híbrido (orquestación en la nube, ejecución en tu infraestructura) |
Equipos Python que quieren menos ceremonia; pipelines con lógica dinámica |
| Dagster | Orientado a activos de datos (software-defined assets): se declara qué tablas y ficheros existen y de qué dependen, y el orquestador deriva las tareas; tipado, pruebas y linaje integrados | Plataformas de datos con muchos conjuntos derivados; se quiere el catálogo y el linaje desde el principio |
| Argo Workflows | DAGs de contenedores en Kubernetes, definidos en YAML; cada paso es un pod | Todo ya corre en Kubernetes; pipelines de ML y CI con contenedores; sin Python obligatorio |
| Cron + scripts | Nada | Un trabajo sin dependencias en una máquina |
La elección para Kilómetro Cero es Airflow por madurez, por los operadores de Spark, HDFS y PostgreSQL que ya existen, y porque el equipo lo conoce; Dagster sería la alternativa seria si la plataforma de datos creciera hasta decenas de conjuntos derivados. Los conceptos (DAG, fecha lógica, idempotencia por partición, sensores, backfill) son los mismos en los cuatro.
- Calidad de datos, linaje y catálogo
Un pipeline que ejecuta bien un cálculo sobre datos malos produce resultados malos a tiempo. La calidad de datos se verifica dentro del pipeline, como tareas:
- Validaciones de entrada, antes de gastar cómputo: el fichero existe y está cerrado (no lo está escribiendo nadie), tiene un tamaño plausible (un día de campaña con 200 líneas es sospechoso), un muestreo parsea como JSON, los campos obligatorios están, la fracción de líneas corruptas es menor que un umbral, las fechas de los eventos caen en el día esperado.
- Validaciones de salida, antes de publicar: el total de importe coincide con el control del job (05-03 imprimía un
Control:por esa razón), no hay productores desconocidos, no hay importes negativos, el número de filas está en el rango histórico (±50 % del día equivalente de la semana anterior). - Cuarentena. Un fichero de entrada que no pasa la validación no se procesa a medias ni se descarta: se mueve a
/km0/cuarentena/<día>/con un informe de por qué, la tarea falla con un mensaje claro, y alguien decide. Los datos válidos de ese día se pueden reprocesar con backfill cuando se corrija la entrada.
Herramientas como Great Expectations o Soda expresan esas comprobaciones de forma declarativa y se integran en Airflow; para el pipeline de esta lección bastará una tarea Python. Dos conceptos más que solo nombramos: el linaje (de qué entradas y con qué código se produjo cada salida: ventas_diarias/dia=2026-09-14 viene de eventos/2026-09-14/pedidos.jsonl y de catalogo.csv con la versión a3f9 de ventas_diarias.py; OpenLineage lo estandariza y Airflow lo emite) y el catálogo de datos (el inventario de qué conjuntos existen, su esquema, su dueño y su frescura: el Hive Metastore de 05-02, DataHub, Amundsen). Ambos son lo que permite responder "¿de dónde sale este número?" sin leer código.
- Práctica:
dags/ventas_diarias.py
dags/ventas_diarias.py7.1 Airflow en docker-compose.yml
Airflow se añade al docker-compose.yml de km0/ con su propia base de metadatos y LocalExecutor (suficiente aquí; en producción, Celery o Kubernetes):
# km0/docker-compose.yml (fragmento)
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
# Conexiones a los sistemas de Kilómetro Cero, como URIs (evita crearlas a mano en la interfaz)
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
- ./servicios/analitica:/app/analitica
command: webserver
ports: ["8090:8080"]
airflow-scheduler:
<<: *airflow
command: scheduler
ports: []Tras docker compose run airflow airflow db migrate y crear un usuario, la interfaz está en localhost:8090. La tabla de destino en km0_analitica, con la clave que hace idempotente la carga:
-- km0/sql/analitica/ventas_diarias.sql
CREATE TABLE IF NOT EXISTS ventas_diarias (
dia date NOT NULL,
mercado text NOT NULL,
productor text NOT NULL,
nombre_productor text,
provincia text,
importe numeric(12,2) NOT NULL,
unidades integer NOT NULL,
pedidos integer NOT NULL,
cargado_en timestamptz NOT NULL DEFAULT now(),
PRIMARY KEY (dia, mercado, productor)
);7.2 El DAG
flowchart LR
S[esperar_eventos<br/>WebHdfsSensor<br/>/km0/eventos/ds/pedidos.jsonl] --> V[validar_entrada<br/>PythonOperator<br/>muestreo, tamaño, fechas]
V --> SP[calcular_ventas<br/>SparkSubmitOperator<br/>ventas_diarias.py]
SP --> C[cargar_postgres<br/>PythonOperator<br/>DELETE + INSERT por dia]
C --> Q[validar_salida<br/>PythonOperator<br/>total = control]
Q --> N[notificar<br/>trigger_rule = all_done]
V -. fallo: cuarentena .-> N
# km0/dags/ventas_diarias.py
"""Pipeline diario de analítica: eventos del lago -> ventas por productor/mercado/día -> km0_analitica.
Cada run procesa el día {{ ds }} (la fecha lógica) y es idempotente: se puede relanzar o hacer 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_EVENTOS = "/km0/eventos/{ds}/pedidos.jsonl"
RUTA_SALIDA = "/km0/agregados/ventas_diarias"
MIN_LINEAS = 1000 # por debajo de esto, el día es sospechoso
MAX_CORRUPTAS = 0.01 # 1 % de líneas ilegibles como máximo
def cliente_hdfs():
from hdfs import InsecureClient # WebHDFS, como subir_eventos_hdfs.py de 04-02
return InsecureClient("http://namenode:9870", user="analitica")
def validar_entrada(ds, ti, **_):
"""Comprueba el fichero del día antes de gastar un job de Spark. Mueve a cuarentena si no pasa."""
ruta = RUTA_EVENTOS.format(ds=ds)
hdfs = cliente_hdfs()
estado = hdfs.status(ruta)
total, corruptas, fuera_de_dia = 0, 0, 0
with hdfs.read(ruta, encoding="utf-8") as f:
for linea in f:
total += 1
try:
ev = json.loads(linea)
dia_ev = datetime.utcfromtimestamp(ev["fecha_ms"] / 1000).strftime("%Y-%m-%d")
if dia_ev != ds:
fuera_de_dia += 1
except (json.JSONDecodeError, KeyError, TypeError):
corruptas += 1
problemas = []
if total < MIN_LINEAS:
problemas.append(f"solo {total} líneas (mínimo {MIN_LINEAS})")
if total and corruptas / total > MAX_CORRUPTAS:
problemas.append(f"{corruptas} líneas corruptas de {total}")
if problemas:
destino = f"/km0/cuarentena/{ds}/pedidos.jsonl"
hdfs.makedirs(f"/km0/cuarentena/{ds}")
hdfs.rename(ruta, destino) # el fichero no se pierde: alguien lo revisará
hdfs.write(f"/km0/cuarentena/{ds}/informe.txt", "\n".join(problemas), overwrite=True)
raise AirflowFailException(f"Entrada de {ds} en cuarentena: " + "; ".join(problemas)) # sin reintentos
ti.xcom_push(key="lineas", value=total) # metadato pequeño para las tareas siguientes
ti.xcom_push(key="fuera_de_dia", value=fuera_de_dia)
print(f"{ds}: {total} líneas, {corruptas} corruptas, {fuera_de_dia} de otro día, {estado['length'] / 1e6:.1f} MB")
def cargar_postgres(ds, ti, **_):
"""Carga la partición dia=ds del Parquet en km0_analitica reemplazando el día completo: idempotente."""
import pyarrow.parquet as pq
from pyarrow import fs
hdfs_fs = fs.HadoopFileSystem("namenode", 8020)
tabla = pq.read_table(f"{RUTA_SALIDA}/dia={ds}", filesystem=hdfs_fs).to_pylist()
filas = [(ds, r["mercado"], r["productor"], r["nombre_productor"], r["provincia"],
r["importe"], r["unidades"], r["pedidos"]) for r in tabla]
pg = PostgresHook(postgres_conn_id="km0_analitica")
with pg.get_conn() as conn, conn.cursor() as cur: # UNA transacción: borrar + insertar, o nada
cur.execute("DELETE FROM ventas_diarias WHERE dia = %s", (ds,))
cur.executemany("""INSERT INTO ventas_diarias
(dia, mercado, productor, nombre_productor, provincia, importe, unidades, pedidos)
VALUES (%s, %s, %s, %s, %s, %s, %s, %s)""", filas)
conn.commit()
ti.xcom_push(key="filas", value=len(filas))
ti.xcom_push(key="total", value=float(sum(r["importe"] for r in tabla)))
def validar_salida(ds, ti, **_):
"""La suma cargada debe coincidir con la que Spark calculó, y el día debe tener un tamaño plausible."""
pg = PostgresHook(postgres_conn_id="km0_analitica")
total_pg, filas_pg = pg.get_first("SELECT COALESCE(SUM(importe), 0), COUNT(*) FROM ventas_diarias WHERE dia = %s", (ds,))
if abs(float(total_pg) - ti.xcom_pull(task_ids="cargar_postgres", key="total")) > 0.01:
raise AirflowFailException(f"Total en PostgreSQL {total_pg} != total cargado")
semana_antes = pg.get_first("SELECT COALESCE(SUM(importe), 0) FROM ventas_diarias WHERE dia = %s::date - 7", (ds,))[0]
if semana_antes and not 0.5 <= float(total_pg) / float(semana_antes) <= 2.0:
print(f"AVISO: el total {total_pg} se desvía más del 50 % del de hace una semana ({semana_antes})")
print(f"{ds}: {filas_pg} filas, total {total_pg} €")
def notificar(ds, dag_run, **_):
"""Resumen del run al canal del equipo. Se ejecuta siempre (trigger_rule=all_done)."""
estados = {ti.task_id: ti.state for ti in dag_run.get_task_instances()}
fallidas = [t for t, s in estados.items() if s == "failed"]
mensaje = f"ventas_diarias {ds}: " + ("OK" if not fallidas else f"FALLO en {', '.join(fallidas)}")
print(mensaje) # aquí iría el webhook de Slack/Teams o un correo
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, # los días son independientes
"email_on_failure": False, # las alertas van por 'notificar' y el callback
"sla": timedelta(hours=3), # cada tarea debe acabar 3 h después del inicio del run
}
with DAG(
dag_id="km0_ventas_diarias",
description="Ventas por productor, mercado y día a partir del lago de eventos",
schedule="0 3 * * *", # a las 3:00, para los datos del día anterior
start_date=datetime(2026, 9, 1),
catchup=False, # el histórico se hace con backfill explícito
max_active_runs=1, # un día a la vez en ejecución normal
default_args=default_args,
tags=["km0", "analitica", "lotes"],
) as dag:
esperar_eventos = WebHdfsSensor(
task_id="esperar_eventos",
webhdfs_conn_id="hdfs_km0",
filepath=RUTA_EVENTOS.format(ds="{{ ds }}"), # plantilla: la fecha lógica del run
poke_interval=300, # comprobar cada 5 minutos
timeout=6 * 3600, # rendirse (y fallar) a las 9:00
mode="reschedule", # libera el worker entre comprobaciones
)
validar = PythonOperator(task_id="validar_entrada", python_callable=validar_entrada)
calcular_ventas = SparkSubmitOperator(
task_id="calcular_ventas",
conn_id="spark_km0",
application="/app/analitica/ventas_diarias.py", # el programa de 05-03, sin cambios
application_args=[
f"{HDFS}{RUTA_EVENTOS.format(ds='{{ ds }}')}",
"/app/analitica/catalogo.csv",
f"{HDFS}{RUTA_SALIDA}",
],
conf={"spark.sql.shuffle.partitions": "8",
"spark.sql.sources.partitionOverwriteMode": "dynamic"}, # solo la partición dia={{ ds }}
executor_memory="1G", executor_cores=2, num_executors=2,
name="km0-ventas-{{ ds }}",
pool="spark", # como mucho 2 jobs Spark a la vez (pool creado en la UI/CLI)
)
cargar = PythonOperator(task_id="cargar_postgres", python_callable=cargar_postgres)
comprobar = PythonOperator(task_id="validar_salida", python_callable=validar_salida)
aviso = PythonOperator(task_id="notificar", python_callable=notificar, trigger_rule="all_done", retries=0)
esperar_eventos >> validar >> calcular_ventas >> cargar >> comprobar >> avisoLos puntos que conviene fijar:
- La fecha lógica gobierna todo.
{{ ds }}en el sensor y en los argumentos de Spark,dscomo parámetro de las funciones Python. El run del 15 de septiembre a las 3:00 tieneds = 2026-09-14y procesa el fichero del 14: es el intervalo de datos que acaba de cerrarse. Nada en el DAG dice "ayer". - El sensor en modo
rescheduleno ocupa un hueco del executor mientras espera: se reprograma cada 5 minutos. Conmode="poke"(por defecto) el worker quedaría bloqueado seis horas. Eltimeoutconvierte "el fichero no llegó" en un fallo visible a las 9:00 en lugar de un pipeline colgado. validar_entradafalla conAirflowFailException, que no reintenta: un fichero corrupto no se arregla esperando cinco minutos, y el fichero ya se ha movido a cuarentena. Los demás fallos (una excepción de red al leer HDFS) sí reintentan segúndefault_args.- Idempotencia en las tres escrituras. Spark sobrescribe solo
dia={{ ds }}(05-03);cargar_postgresborra e inserta el día en una transacción; la cuarentena usarename, atómico en HDFS. Relanzar cualquier tarea, o el run entero, deja el mismo estado. - XCom para metadatos.
lineas,filas,total: números, no datos. Los datos viajan por HDFS y PostgreSQL. pool="spark"limita a dos los jobs Spark simultáneos aunque un backfill cree siete runs a la vez (apartado 8), ymax_active_runs=1mantiene la ejecución normal en un día cada vez.notificarcontrigger_rule="all_done"se ejecuta tanto si el run fue bien como si alguna tarea falló, y por eso puede informar del estado; sin esa regla, un fallo envalidar_entradala dejaría enupstream_failedy nadie se enteraría. El SLA dedefault_argsañade una alerta si alguna tarea sigue sin terminar tres horas después de las 3:00.
Al activar el DAG en la interfaz, el scheduler crea el primer run en la siguiente ejecución programada (con catchup=False), y la rejilla muestra un cuadrado por día y tarea, verde, rojo o amarillo (reintentando). Un clic en un cuadrado da los logs de ese intento, incluida la salida de spark-submit con el explain() de 05-03.
- Backfill de la Semana de la Vendimia
El caso del ejercicio 3 de 05-03, ahora con el pipeline: pedidos corrige el precio de vino-crianza y regenera los ficheros pedidos.jsonl del 8 al 14 de septiembre en HDFS. Hay que recalcular esos siete días, y solo esos, sin tocar el resto y sin interferir con el run diario. Con el DAG parametrizado por fecha lógica e idempotente, es un comando:
# --reset-dagruns: los runs de esos días ya existen (en éxito); limpiarlos y reejecutar
docker compose exec airflow-scheduler airflow dags backfill km0_ventas_diarias \
--start-date 2026-09-08 --end-date 2026-09-14 \
--reset-dagruns --rerun-failed-tasksLo que ocurre: el scheduler crea (o resetea) siete DAG runs con ds del 08 al 14 y los ejecuta respetando las dependencias de cada uno y los límites globales: el pool spark permite dos calcular_ventas a la vez, así que los siete jobs Spark se ejecutan en cuatro tandas; max_active_runs no aplica al backfill (tiene su propio límite, --max-active-runs en versiones recientes, o el max_active_runs del DAG según la versión: conviene comprobarlo, porque un backfill de un año con siete runs a la vez puede saturar HDFS). Los sensores pasan de inmediato (los ficheros existen), las validaciones se repiten sobre los ficheros corregidos, Spark sobrescribe dia=2026-09-08/ ... dia=2026-09-14/ y deja los demás días intactos, y cada carga borra e inserta su día. Al terminar, la rejilla muestra los siete días con un nuevo intento en verde, y SELECT dia, SUM(importe) FROM ventas_diarias WHERE dia BETWEEN '2026-09-08' AND '2026-09-14' GROUP BY 1 refleja los precios corregidos.
Dos variantes útiles: para reejecutar solo desde la carga (Spark ya corrió bien, falló PostgreSQL), airflow tasks clear km0_ventas_diarias -t cargar_postgres --downstream -s 2026-09-08 -e 2026-09-14 limpia esa tarea y las siguientes de esos runs, y el scheduler las reejecuta; y para un run puntual con una fecha concreta, airflow dags trigger km0_ventas_diarias --logical-date 2026-09-14. Ninguna de las tres cosas era posible con el crontab del apartado 1 sin editar comandos a mano.
Errores Comunes y Consejos
- Usar "ayer" en lugar de la fecha lógica.
datetime.now()odate -d yesterdaydentro de una tarea hacen imposible el backfill y producen resultados distintos según cuándo se ejecute. Siempre{{ ds }}/data_interval_start. - Confundir cuándo se ejecuta un run con qué datos procesa. El run con fecha lógica 14 se ejecuta el 15. Es la fuente de más confusión de Airflow; pensar en "el intervalo que acaba de cerrarse" lo aclara.
- Tareas no idempotentes.
INSERTsin borrar antes,appenda un fichero, un contador. Reintentar duplica; backfill acumula. Partición completa con sobreescritura, siempre. - Cálculo pesado dentro del worker de Airflow. Un
PythonOperatorque lee 150 MB de JSON y agrega en pandas convierte al worker en un nodo de cómputo mal dimensionado. Airflow orquesta; Spark, la base de datos o un contenedor calculan. - Datos por XCom. Un DataFrame serializado en la base de metadatos. XCom es para rutas y contadores.
catchup=Truesin querer. Activar un DAG constart_datede hace dos años ycatchuppor defecto lanza 730 runs.catchup=Falsey backfill explícito.depends_on_past=Truepor precaución. Un día fallido bloquea todos los siguientes hasta que alguien lo arregla. Solo cuando cada día depende de verdad del anterior.- Sensores en modo
pokecon timeouts largos. Cada sensor ocupa un hueco del executor durante horas; diez sensores esperando bloquean el pipeline entero.mode="reschedule", o planificación por dataset. - Reintentar errores deterministas. Un fichero corrupto o un bug se reintentan tres veces con backoff y fallan una hora después.
AirflowFailExceptionpara lo que no se arregla esperando. - Sin validación de salida. El pipeline en verde no significa datos correctos. Un control de totales y un rango plausible frente al histórico cuestan veinte líneas y evitan gráficas dobladas.
- Secretos en el DAG. Contraseñas de PostgreSQL en el código. Conexiones de Airflow, variables de entorno o un gestor de secretos (06-04).
Ejercicios
Ejercicio 1: Añadir las recomendaciones al pipeline
Amplía el DAG para que, después de validar_salida, entrene las recomendaciones con recomendaciones_als.py (05-03) sobre los clics de los últimos 7 días y cargue el resultado en la tabla recomendaciones(cliente, producto, afinidad, generado_en) de km0_catalogo. Decide: ¿qué sensor necesita? ¿De qué tareas depende? ¿Cómo haces idempotente la carga si la tabla no tiene una "partición por día" natural? ¿Debe usar depends_on_past? ¿Qué pasa con las recomendaciones durante un backfill de siete días?
Ejercicio 2: Un día malo
El 20 de septiembre a las 3:00 ocurre lo siguiente: subir_eventos_hdfs.py tuvo un fallo y el fichero /km0/eventos/2026-09-19/pedidos.jsonl llega a HDFS a las 7:40 con 180 000 líneas, de las que 4 200 no son JSON válido. Describe, tarea a tarea, qué hace el DAG del apartado 7 entre las 3:00 y las 9:00: estados, reintentos, dónde acaba el fichero, qué ve el equipo. Después, alguien corrige el fichero y lo vuelve a subir a las 11:30. ¿Qué comando ejecuta el equipo para completar el día 19, y qué ocurre con el run del día 20 esa noche?
Ejercicio 3: cron o Airflow
Para cada uno de estos trabajos de Kilómetro Cero, decide si lo dejarías en cron, lo pondrías en Airflow, o lo sacarías a un flujo de 05-04, y justifica en una frase: (a) borrar los checkpoints de Flink de más de 30 días en HDFS; (b) generar cada lunes el informe semanal PDF por productor a partir de ventas_diarias, y enviarlo por correo; (c) recalcular el stock mínimo por producto y mercado cada 5 minutos; (d) exportar a MinIO una copia de km0_inventario cada noche; (e) cargar cada día en PostgreSQL las posiciones tardías de reparto.posiciones que el flujo desvió a la salida lateral.
Soluciones
Ejercicio 1.
Tareas nuevas: esperar_clics (un WebHdfsSensor sobre /km0/clics/{{ ds }}/ o, mejor, sobre un fichero de cierre _CERRADO que el proceso de subida por horas escriba al acabar el día; sin él, el directorio existe desde la primera hora y el sensor pasaría con datos incompletos), entrenar_als (SparkSubmitOperator con recomendaciones_als.py y un argumento con el rango {{ macros.ds_add(ds, -6) }} a {{ ds }}), cargar_recomendaciones (PythonOperator). Dependencias: esperar_clics >> entrenar_als >> cargar_recomendaciones, y entrenar_als también aguas abajo de validar_salida solo si las recomendaciones usan las ventas (no es el caso: usan clics), así que en rigor son dos ramas paralelas que confluyen en notificar. Idempotencia sin partición por día: la tabla se reemplaza entera en una transacción (DELETE FROM recomendaciones; INSERT ...; COMMIT), o se escribe en recomendaciones_nueva y se hace ALTER TABLE ... RENAME intercambiando las dos (atómico en PostgreSQL), o se guarda una columna generado_en = ds y la web lee WHERE generado_en = (SELECT MAX(generado_en) ...): la tercera opción es la que permite backfill sin pisar la versión actual. depends_on_past: no; cada entrenamiento es independiente. Durante un backfill de siete días, las siete ejecuciones entrenarían siete modelos con ventanas móviles de clics y cargarían siete veces: con la columna generado_en es inocuo pero inútil; lo razonable es que entrenar_als viva en otro DAG con su propio ciclo (diario, sin backfill salvo que se pida), o que en el backfill se excluya con --task-regex la rama de recomendaciones.
Ejercicio 2.
3:00: el scheduler crea el run ds=2026-09-19. esperar_eventos comprueba cada 5 minutos (reschedule, sin ocupar worker) y no encuentra el fichero: estado up_for_reschedule, cuadrado amarillo. 6:00: el SLA de 3 h se incumple y Airflow registra un SLA miss con aviso al equipo (el pipeline aún no ha fallado, pero va tarde). 7:40: el fichero aparece; en la comprobación de las 7:45 el sensor pasa a success. validar_entrada arranca: 180 000 líneas (> 1 000), pero 4 200 corruptas son el 2,3 % (> 1 %): mueve el fichero a /km0/cuarentena/2026-09-19/pedidos.jsonl, escribe informe.txt y lanza AirflowFailException: estado failed sin reintentos. calcular_ventas, cargar_postgres y validar_salida quedan en upstream_failed. notificar se ejecuta (all_done) y publica "ventas_diarias 2026-09-19: FALLO en validar_entrada"; el callback de fallo también avisa. A las 9:00 no pasa nada más: el timeout del sensor ya no aplica porque el sensor terminó. El equipo ve el cuadrado rojo en validar_entrada, lee el log con "4200 líneas corruptas de 180000" y encuentra el fichero en cuarentena.
11:30: el fichero corregido se sube a /km0/eventos/2026-09-19/pedidos.jsonl. El equipo ejecuta airflow tasks clear km0_ventas_diarias -t validar_entrada --downstream -s 2026-09-19 -e 2026-09-19 (o pulsa Clear en la tarea desde la interfaz): validar_entrada y las siguientes vuelven a None y el scheduler las reejecuta; el sensor no se repite porque no se limpió. Con max_active_runs=1, ese run ocupa el hueco hasta que termina (unos minutos). Esa noche, a las 3:00 del día 21, el run ds=2026-09-20 se crea con normalidad: depends_on_past=False, así que no le afecta lo que pasó con el 19, y el 19 ya está en verde.
Ejercicio 3.
(a) cron (o el propio Flink con retención de checkpoints configurada): un comando sin dependencias ni datos, hdfs dfs -rm con una fecha; si se salta un día, no pasa nada. (b) Airflow: depende de que ventas_diarias haya cargado los siete días (un ExternalTaskSensor sobre el DAG diario), tiene salida que hay que poder regenerar (backfill de una semana si se corrigen datos) y un envío que no debe duplicarse. (c) Flujo (05-04): cada 5 minutos con latencia de segundos y estado por clave es una ventana tumbling sobre stock.actualizado, no un lote lanzado 288 veces al día. (d) Airflow, aunque sea una sola tarea: se quiere historial, alerta si falla y un sensor o validación de que la copia se completó (tamaño, _SUCCESS); en cron sería aceptable si se añadiera monitorización externa. (e) Airflow: es un lote diario parametrizado por fecha que lee /km0/reparto/tardias/{{ ds }}/ y carga por partición, y encaja como una tarea más del pipeline diario: es la reconciliación entre el flujo y el lote de la que hablaba el ejercicio 1 de 05-04.
Conclusión
Un pipeline de datos es la parte del sistema distribuido que convierte trabajos sueltos en una plataforma: un DAG de tareas con dependencias explícitas, planificado por tiempo y por dato (sensores), con reintentos con backoff para lo transitorio y fallo inmediato para lo determinista, tareas idempotentes que escriben particiones completas con sobreescritura, backfill parametrizado por la fecha lógica y no por "ayer", SLAs y notificaciones para saber cuándo algo va tarde, validaciones de entrada y salida con cuarentena para no publicar datos malos a tiempo, y todo en el repositorio, versionado. cron no ofrece nada de eso; Airflow lo ofrece con un scheduler, un executor, workers, una base de metadatos y DAGs en Python, y Prefect, Dagster y Argo Workflows lo reformulan con distintos énfasis. dags/ventas_diarias.py encadena el sensor sobre /km0/eventos/{{ ds }}/pedidos.jsonl, la validación, el SparkSubmitOperator con el programa de 05-03, la carga transaccional en km0_analitica y la comprobación de totales, y el backfill de la Semana de la Vendimia se reduce a un comando con dos fechas.
Con esto se cierra el Módulo 5, y la plataforma de datos de Kilómetro Cero queda completa:
| Pieza | Qué es | Lección |
|---|---|---|
| Lago de datos | HDFS con un directorio por día: /km0/eventos/<día>/pedidos.jsonl, /km0/clics/<día>/, convertidos a Parquet particionado por dia en /km0/agregados/ |
04-02, 05-03 |
| Modelos de cómputo | Llevar el cálculo a los datos; scatter/gather, colas de trabajo, BSP, dataflow; shuffle, sesgo, rezagados, reejecución determinista | 05-01 |
| Lotes | MapReduce como base histórica y vocabulario; Spark con DataFrames, Catalyst y Parquet para ventas_diarias.py; MLlib/ALS para recomendaciones |
05-02, 05-03 |
| Flujos | Flink (y Structured Streaming) sobre pedidos.eventos y reparto.posiciones: tiempo de evento, marcas de agua, ventanas, checkpoints, sinks idempotentes; alertas de stock y panel de reparto |
05-04 |
| Pipelines | Airflow: km0_ventas_diarias con sensores, validación, Spark, carga idempotente, backfill |
05-05 |
Los datos están repartidos (Módulo 4) y se procesan en masa, por lotes y en tiempo real, sin que nadie lance nada a mano. Pero hay algo que hemos dado por supuesto en todos los módulos hasta ahora: cualquier proceso podía hablar con cualquier servicio y leer cualquier dato. ventas_diarias.py lee todo el lago; cargar_postgres entra en km0_analitica con una contraseña en una variable de entorno; el consumidor de Flink lee pedidos.eventos sin identificarse; un hdfs dfs -rm desde cualquier contenedor borraría un año de eventos; y las URLs prefirmadas de MinIO (04-03) fueron el único mecanismo de acceso que hemos diseñado con cuidado. En un sistema con decenas de servicios, cientos de trabajos y datos personales de Ana, Marc y Lucía, eso no puede seguir así. El Módulo 6 trata la seguridad en sistemas distribuidos, y empieza por las dos preguntas que ningún sistema puede eludir: quién es quién (autenticación) y quién puede qué (autorización).
Curso de Arquitecturas Distribuidas
Módulo 1: Introducción a los Sistemas Distribuidos
- Conceptos Básicos de Sistemas Distribuidos
- Modelos de Sistemas Distribuidos
- Ventajas y Desafíos de los Sistemas Distribuidos
- Las Falacias de la Computación Distribuida
- Tiempo, Relojes y Ordenación de Eventos
- Del Monolito a la Plataforma Distribuida: el Caso Kilómetro Cero
Módulo 2: Comunicación en Sistemas Distribuidos
- Protocolos de Comunicación
- RPC y RMI
- gRPC y Serialización de Datos
- Mensajería y Colas de Mensajes
- Patrones de Comunicación Asíncrona
Módulo 3: Consistencia y Replicación
- Modelos de Consistencia
- El Teorema CAP y PACELC
- Algoritmos de Consenso
- Replicación de Datos
- Transacciones Distribuidas y Sagas
Módulo 4: Almacenamiento Distribuido
- Particionado de Datos y Hashing Consistente
- Sistemas de Archivos Distribuidos
- Almacenamiento de Objetos
- Bases de Datos Distribuidas
- Cachés Distribuidos
Módulo 5: Computación Distribuida
- Modelos de Computación Distribuida
- MapReduce y Hadoop
- Spark y Computación en Memoria
- Procesamiento de Flujos de Datos
- Planificación de Trabajos y Pipelines de Datos
Módulo 6: Seguridad en Sistemas Distribuidos
- Autenticación y Autorización
- Cifrado y Protección de Datos
- Gestión de Identidades
- Seguridad entre Servicios: mTLS y Gestión de Secretos
- Puertas de Enlace, Limitación de Tasa y Auditoría
Módulo 7: Monitoreo y Mantenimiento
- Monitoreo de Sistemas Distribuidos
- Logs Centralizados y Trazabilidad Distribuida
- Gestión de Fallos y Recuperación
- Patrones de Resiliencia: Timeouts, Reintentos y Circuit Breaker
- Automatización y Orquestación
- Pruebas en Sistemas Distribuidos e Ingeniería del Caos
