Saltar al contenido

lección 8

Proyecto: pipeline con validación completa y alertas

Un pipeline end-to-end que ingiere datos, valida con GE, alerta si falla y genera un reporte de calidad. Todo lo aprendido en acción.

55 min

### El encargo: pipeline de calidad para RetailMax

RetailMax es una cadena de tiendas con 50 locales en España. Cada noche, a las 2 AM, un sistema legacy exporta un CSV con todas las ventas del día. Tu equipo de datos ingiere ese CSV, lo transforma y lo carga en el data warehouse para que el equipo de BI tenga los dashboards actualizados a las 7 AM cuando llegan a la oficina.

El problema: la semana pasada el CSV llegó con 40% de campos vacíos y nadie se enteró hasta que el director comercial vio el dashboard con cifras absurdas a las 9 AM. La reunión de seguimiento se basó en datos incorrectos. El CEO está furioso. Te piden: "que esto NO vuelva a pasar. Nunca."

Tu misión: construir un pipeline que no solo ingiera y transforme, sino que VALIDE en cada paso y ALERTE si algo no cumple los estándares. Si los datos son malos, el pipeline se detiene y alerta al equipo ANTES de que los datos malos lleguen al dashboard.

### Arquitectura del pipeline con calidad

Dos puntos de validación: entrada (raw) y salida (clean). Fallo en cualquiera = pipeline detenido.

### Paso 1: Simular los datos de RetailMax

En un proyecto real leerías un CSV de un bucket S3 o un servidor SFTP. Aquí vamos a simular los datos con variantes: un día "normal" y un día "malo" para probar que nuestras validaciones funcionan.

1import pandas as pd
2import numpy as np
3from datetime import datetime, timedelta
4
5def generar_ventas_dia(fecha: str, escenario: str = "normal") -> pd.DataFrame:
6 """Genera datos de ventas simulados para un día."""
7 np.random.seed(42)
8
9 if escenario == "normal":
10 n = 2500 # Volumen normal de un día
11 datos = {
12 'venta_id': [f"V-{fecha.replace('-','')}-{i:05d}" for i in range(n)],
13 'tienda_id': np.random.randint(1, 51, n),
14 'producto': np.random.choice(['Camiseta', 'Pantalón', 'Zapatos', 'Chaqueta', 'Accesorios'], n),
15 'cantidad': np.random.randint(1, 5, n),
16 'precio_unitario': np.random.choice([19.99, 39.99, 59.99, 89.99, 14.99], n),
17 'metodo_pago': np.random.choice(['tarjeta', 'efectivo', 'bizum'], n, p=[0.6, 0.25, 0.15]),
18 'fecha_venta': [fecha] * n,
19 'hora_venta': [f"{h:02d}:{m:02d}" for h, m in zip(np.random.randint(9, 22, n), np.random.randint(0, 60, n))],
20 }
21 df = pd.DataFrame(datos)
22 df['importe_total'] = df['cantidad'] * df['precio_unitario']
23 # Insertar mínimo ruido realista (1%)
24 mask_nulos = np.random.choice(n, int(n * 0.01), replace=False)
25 df.loc[mask_nulos[:5], 'metodo_pago'] = None
26 return df
27
28 elif escenario == "problematico":
29 n = 800 # Volumen muy bajo (problema!)
30 datos = {
31 'venta_id': [f"V-{fecha.replace('-','')}-{i:05d}" for i in range(n)] + ['V-ERROR', 'V-ERROR'], # Duplicados
32 'tienda_id': list(np.random.randint(1, 51, n)) + [99, 100], # Tiendas que no existen
33 'producto': list(np.random.choice(['Camiseta', 'Pantalón', 'Zapatos'], n)) + [None, None],
34 'cantidad': list(np.random.randint(1, 5, n)) + [-1, 0],
35 'precio_unitario': list(np.random.choice([19.99, 39.99, 59.99], n)) + [0.0, -10.0],
36 'metodo_pago': list(np.random.choice(['tarjeta', 'efectivo', 'bizum'], n)) + ['bitcoin', None],
37 'fecha_venta': [fecha] * n + ['2020-01-01', fecha], # Fecha antigua
38 'hora_venta': [f"{h:02d}:{m:02d}" for h, m in zip(np.random.randint(9, 22, n), np.random.randint(0, 60, n))] + ['25:00', '10:30'],
39 }
40 df = pd.DataFrame(datos)
41 df['importe_total'] = df['cantidad'].astype(float) * df['precio_unitario'].astype(float)
42 # 15% de nulos en método de pago
43 mask_nulos = np.random.choice(len(df), int(len(df) * 0.15), replace=False)
44 df.loc[mask_nulos, 'metodo_pago'] = None
45 return df
46
47# Generar ambos escenarios
48print("Día normal:", len(generar_ventas_dia("2024-01-15", "normal")), "registros")
49print("Día problemático:", len(generar_ventas_dia("2024-01-16", "problematico")), "registros")

Dos escenarios: datos normales (pasan validación) y problemáticos (fallan validación)

### Paso 2: Definir el contrato de datos

1# Contrato para las ventas diarias de RetailMax
2CONTRATO_VENTAS = {
3 "nombre": "ventas_diarias_retailmax",
4 "version": "1.0.0",
5 "productor": "sistema-legacy-tiendas",
6 "consumidores": ["pipeline-datos", "dashboard-bi"],
7 "schema": {
8 "venta_id": {"obligatorio": True, "unico": True},
9 "tienda_id": {"obligatorio": True, "min": 1, "max": 50},
10 "producto": {"obligatorio": True, "max_nulos_pct": 2.0},
11 "cantidad": {"obligatorio": True, "min": 1, "max": 99},
12 "precio_unitario": {"obligatorio": True, "min": 0.01, "max": 999.99},
13 "importe_total": {"obligatorio": True, "min": 0.01},
14 "metodo_pago": {"obligatorio": True, "max_nulos_pct": 3.0,
15 "valores_validos": ["tarjeta", "efectivo", "bizum"]},
16 "fecha_venta": {"obligatorio": True},
17 "hora_venta": {"obligatorio": True},
18 },
19 "volumen": {
20 "min_diario": 1500,
21 "max_diario": 10000,
22 },
23 "sla": {
24 "hora_disponibilidad": "03:00",
25 "max_retraso_horas": 2,
26 },
27}

Contrato formal: schema, rangos, volumen esperado y SLA de frescura

### Paso 3: Implementar las validaciones

1from dataclasses import dataclass, field
2from typing import Any
3
4@dataclass
5class ResultadoCheck:
6 nombre: str
7 paso: bool
8 detalle: str
9 severidad: str = "ERROR"
10
11@dataclass
12class ResultadoValidacion:
13 etapa: str
14 checks: list[ResultadoCheck] = field(default_factory=list)
15
16 @property
17 def exito(self) -> bool:
18 return all(c.paso for c in self.checks if c.severidad == "ERROR")
19
20 def resumen(self) -> str:
21 ok = sum(1 for c in self.checks if c.paso)
22 fail = sum(1 for c in self.checks if not c.paso)
23 return f"{self.etapa}: {ok} OK, {fail} FAIL → {'✅' if self.exito else '🔴'}"
24
25def validar_raw(df: pd.DataFrame, contrato: dict) -> ResultadoValidacion:
26 """Validación de datos RAW contra el contrato."""
27 resultado = ResultadoValidacion(etapa="RAW")
28
29 # Check 1: Volumen
30 vol = contrato["volumen"]
31 n = len(df)
32 resultado.checks.append(ResultadoCheck(
33 nombre="volumen_minimo",
34 paso=n >= vol["min_diario"],
35 detalle=f"{n} registros (mín: {vol['min_diario']})",
36 ))
37
38 # Check 2: Columnas esperadas
39 cols_esperadas = set(contrato["schema"].keys())
40 cols_presentes = set(df.columns)
41 faltantes = cols_esperadas - cols_presentes
42 resultado.checks.append(ResultadoCheck(
43 nombre="schema_columnas",
44 paso=len(faltantes) == 0,
45 detalle=f"Faltantes: {faltantes}" if faltantes else "Todas las columnas presentes",
46 ))
47
48 # Check 3: Duplicados en clave
49 if "venta_id" in df.columns:
50 n_dup = df["venta_id"].duplicated().sum()
51 resultado.checks.append(ResultadoCheck(
52 nombre="unicidad_venta_id",
53 paso=n_dup == 0,
54 detalle=f"{n_dup} duplicados en venta_id",
55 ))
56
57 # Check 4: Nulos en campos obligatorios
58 for campo, reglas in contrato["schema"].items():
59 if campo in df.columns and reglas.get("obligatorio"):
60 max_nulos = reglas.get("max_nulos_pct", 0)
61 pct_nulos = df[campo].isna().mean() * 100
62 resultado.checks.append(ResultadoCheck(
63 nombre=f"nulos_{campo}",
64 paso=pct_nulos <= max_nulos,
65 detalle=f"{pct_nulos:.1f}% nulos (máx: {max_nulos}%)",
66 severidad="ERROR" if pct_nulos > max_nulos * 2 else "WARNING",
67 ))
68
69 return resultado
70
71def validar_clean(df: pd.DataFrame, contrato: dict) -> ResultadoValidacion:
72 """Validación de datos limpios (post-transformación)."""
73 resultado = ResultadoValidacion(etapa="CLEAN")
74
75 # Check: rangos
76 for campo, reglas in contrato["schema"].items():
77 if campo not in df.columns:
78 continue
79 if "min" in reglas:
80 fuera = (df[campo] < reglas["min"]).sum()
81 resultado.checks.append(ResultadoCheck(
82 nombre=f"rango_min_{campo}",
83 paso=fuera == 0,
84 detalle=f"{fuera} valores < {reglas['min']}",
85 ))
86 if "max" in reglas:
87 fuera = (df[campo] > reglas["max"]).sum()
88 resultado.checks.append(ResultadoCheck(
89 nombre=f"rango_max_{campo}",
90 paso=fuera == 0,
91 detalle=f"{fuera} valores > {reglas['max']}",
92 ))
93 if "valores_validos" in reglas:
94 invalidos = (~df[campo].dropna().isin(reglas["valores_validos"])).sum()
95 resultado.checks.append(ResultadoCheck(
96 nombre=f"valores_{campo}",
97 paso=invalidos == 0,
98 detalle=f"{invalidos} valores no permitidos",
99 ))
100
101 return resultado

Funciones de validación separadas por etapa: raw (schema + volumen) y clean (rangos + reglas)

### Paso 4: Sistema de alertas y reporte

1def enviar_alerta(resultado: ResultadoValidacion, canal: str = "slack"):
2 """Envía alerta cuando la validación falla."""
3 if resultado.exito:
4 return # No alertar si todo OK
5
6 errores = [c for c in resultado.checks if not c.paso]
7
8 print(f"\n{'='*60}")
9 print(f"🚨 ALERTA — Pipeline RetailMax — {resultado.etapa}")
10 print(f"{'='*60}")
11 print(f"Canal: #{canal}")
12 print(f"Timestamp: {datetime.now().isoformat()}")
13 print(f"Checks fallidos: {len(errores)}/{len(resultado.checks)}")
14 print(f"\nDetalle:")
15 for error in errores:
16 icono = "🔴" if error.severidad == "ERROR" else "⚠️"
17 print(f" {icono} {error.nombre}: {error.detalle}")
18 print(f"\n📋 Acción: Investigar datos de entrada. NO se cargaron al warehouse.")
19 print(f"📖 Runbook: https://wiki.retailmax.internal/runbooks/pipeline-ventas")
20 print(f"{'='*60}")
21
22def generar_reporte_calidad(df: pd.DataFrame, resultado_raw: ResultadoValidacion,
23 resultado_clean: ResultadoValidacion) -> str:
24 """Genera un reporte de calidad del día."""
25 reporte = []
26 reporte.append("=" * 60)
27 reporte.append("📊 REPORTE DE CALIDAD — RetailMax Ventas Diarias")
28 reporte.append(f" Fecha: {datetime.now().strftime('%Y-%m-%d %H:%M')}")
29 reporte.append(f" Registros procesados: {len(df)}")
30 reporte.append("=" * 60)
31
32 # Métricas de las 5 dimensiones
33 completitud = (1 - df.isna().sum().sum() / (df.shape[0] * df.shape[1])) * 100
34 unicidad = (1 - df['venta_id'].duplicated().mean()) * 100 if 'venta_id' in df.columns else 100
35
36 reporte.append(f"\n Completitud: {completitud:.1f}%")
37 reporte.append(f" Unicidad: {unicidad:.1f}%")
38 reporte.append(f" Validación RAW: {resultado_raw.resumen()}")
39 reporte.append(f" Validación CLEAN: {resultado_clean.resumen()}")
40
41 # Estado global
42 todo_ok = resultado_raw.exito and resultado_clean.exito
43 reporte.append(f"\n Estado: {'✅ CARGA EXITOSA' if todo_ok else '🔴 CARGA BLOQUEADA'}")
44 reporte.append("=" * 60)
45
46 return "\n".join(reporte)

Alertas con contexto completo + reporte de calidad resumido para stakeholders

### Paso 5: El pipeline completo

1def pipeline_retailmax(fecha: str, escenario: str = "normal"):
2 """Pipeline completo con validación y alertas."""
3 print(f"\n{'▓'*60}")
4 print(f"▓ PIPELINE RETAILMAX — {fecha} ({escenario})")
5 print(f"{'▓'*60}\n")
6
7 # === PASO 1: INGESTA ===
8 print("📥 [1/4] Ingesta de datos...")
9 df_raw = generar_ventas_dia(fecha, escenario)
10 print(f" Leídos {len(df_raw)} registros del CSV")
11
12 # === PASO 2: VALIDAR RAW ===
13 print("\n🔍 [2/4] Validando datos raw...")
14 resultado_raw = validar_raw(df_raw, CONTRATO_VENTAS)
15 print(f" {resultado_raw.resumen()}")
16
17 if not resultado_raw.exito:
18 enviar_alerta(resultado_raw, "data-alerts")
19 print("\n⛔ PIPELINE DETENIDO — Datos raw no cumplen contrato")
20 return None
21
22 # === PASO 3: TRANSFORMAR ===
23 print("\n⚙️ [3/4] Limpiando y transformando...")
24 df_clean = df_raw.copy()
25 # Eliminar duplicados
26 df_clean = df_clean.drop_duplicates(subset=['venta_id'])
27 # Eliminar filas con campos clave nulos
28 df_clean = df_clean.dropna(subset=['venta_id', 'tienda_id', 'importe_total'])
29 # Corregir tipos
30 df_clean['cantidad'] = pd.to_numeric(df_clean['cantidad'], errors='coerce')
31 df_clean['precio_unitario'] = pd.to_numeric(df_clean['precio_unitario'], errors='coerce')
32 print(f" {len(df_raw)}{len(df_clean)} registros (eliminados {len(df_raw)-len(df_clean)} problemáticos)")
33
34 # === PASO 4: VALIDAR CLEAN ===
35 print("\n🔍 [4/4] Validando datos limpios...")
36 resultado_clean = validar_clean(df_clean, CONTRATO_VENTAS)
37 print(f" {resultado_clean.resumen()}")
38
39 if not resultado_clean.exito:
40 enviar_alerta(resultado_clean, "data-alerts")
41 print("\n⛔ PIPELINE DETENIDO — Datos limpios no cumplen estándares")
42 return None
43
44 # === CARGA ===
45 print("\n📤 Cargando en warehouse...")
46 # df_clean.to_parquet(f"warehouse/ventas_{fecha}.parquet")
47 print(f" ✅ {len(df_clean)} registros cargados exitosamente")
48
49 # === REPORTE ===
50 reporte = generar_reporte_calidad(df_clean, resultado_raw, resultado_clean)
51 print(f"\n{reporte}")
52
53 return df_clean
54
55# Ejecutar con escenario normal (debe pasar)
56resultado_ok = pipeline_retailmax("2024-01-15", "normal")

El pipeline completo: ingesta → validar → transformar → validar → cargar (o alertar)

### Probando con datos problemáticos

1# Ejecutar con escenario problemático (debe fallar y alertar)
2resultado_malo = pipeline_retailmax("2024-01-16", "problematico")
3
4# El pipeline se detuvo ANTES de cargar datos malos al warehouse
5# El equipo recibió la alerta a las 3:05 AM
6# A las 7 AM cuando llegan a la oficina, ya saben qué pasó
7# El dashboard NO se actualizó con datos incorrectos

Con datos malos, el pipeline se detiene limpiamente y alerta — exactamente lo que queríamos

Consejo de senior: este patrón (validar entrada → transformar → validar salida → cargar o abortar) es el ESTÁNDAR en equipos de datos profesionales. No es complejo, no requiere herramientas caras. Solo requiere la DISCIPLINA de no saltarse la validación "porque tengo prisa". El día que tienes prisa es el día que necesitas la validación más que nunca.

En producción real, NUNCA elimines registros silenciosamente. Si tu transformación elimina el 30% de los registros, eso en sí mismo debería ser una alerta. Lo correcto es: si el % de registros eliminados supera un umbral (ej: 5%), alerta al equipo. Puede ser un problema upstream, no datos que debas filtrar.

### Recapitulación: todo lo que has aprendido

En estas 7 lecciones has construido un arsenal completo de calidad de datos. Empezaste entendiendo POR QUÉ importa (con horror stories reales), aprendiste las 5 dimensiones para MEDIR la calidad, creaste validadores con Python puro, adoptaste Great Expectations como framework profesional, definiste contratos de datos entre equipos, configuraste alertas multi-canal y finalmente integraste TODO en un pipeline de producción.

Ya no eres el ingeniero que dice "el pipeline funciona". Eres el ingeniero que dice "el pipeline funciona, los datos son correctos, y si dejan de serlo me entero en minutos". Esa es la diferencia que marca esta skill. Esa es la diferencia entre un junior y un senior. Bienvenido al otro lado.

### Has terminado cuando…

  • Tu pipeline PARA si la validación de entrada falla, y termina con código de salida ≠ 0 (compruébalo: `python pipeline.py; echo $?` debe imprimir 1)
  • La validación de salida DETECTA los defectos de negocio del dataset de reservas: precio negativo, cero noches, checkout anterior al checkin, estado desconocido y hotel inexistente
  • Con los datos limpios, el pipeline LLEGA AL FINAL e imprime el reporte de calidad
  • Tu contrato está en UN SOLO SITIO (una constante o un diccionario), no repartido por el código
  • Las alertas distinguen al menos dos severidades y dicen a quién avisar
  • El informe semanal mira LOS SIETE DÍAS, no solo el último — un incidente del día 4 no puede desaparecer del informe
  • Puedes ejecutarlo dos veces y obtener EXACTAMENTE la misma salida (nada de datetime.now())

Para tu portfolio: sube el pipeline a un repositorio con el CSV de ejemplo, el contrato en un YAML aparte, los tests con pytest, y un README que explique qué decide parar el pipeline y por qué. Eso último es lo que distingue a alguien que ha hecho un curso de alguien que entiende el oficio.

## ejercicios

[01]

Pipeline con validación para reservas de hotel

HotelPlus te contrata para construir un pipeline de reservas diarias con validación completa. El CSV llega cada noche con reservas del día. Implementa el pipeline con: validación de entrada (schema + volumen), transformación (limpieza) y validación de salida (reglas de negocio).

💡 Resultado esperado

🏨 PIPELINE HOTELPLUS — Reservas diarias

🔍 Validando entrada...
   ⚠️ 1 duplicado(s) en reserva_id — se limpiará en transformar
Cargando editor...
[02]

Generar reporte ejecutivo de calidad semanal

El VP de Datos quiere un reporte SEMANAL de calidad que pueda enseñar al board. Debe ser comprensible por no-técnicos: scores, tendencias, top problemas y recomendaciones. Genera el reporte para 5 datasets durante 7 días.

💡 Resultado esperado

╔══════════════════════════════════════════════════════════╗
║               📊 REPORTE SEMANAL DE CALIDAD               ║
║             Semana: 2024-01-15 → 2024-01-21              ║
╠══════════════════════════════════════════════════════════╣
Cargando editor...
[03]

Pipeline completo con GE + alertas + reporte

Integra TODO: un pipeline que usa Great Expectations para validar, envía alertas simuladas si falla, y genera un reporte de calidad al finalizar. Es tu "template" para cualquier pipeline futuro.

💡 Resultado esperado

💰 PIPELINE DE FACTURAS

🔍 Validando entrada...
   ✅ Entrada OK
Cargando editor...

Regístrate para guardar tu progreso.

## comentarios

Reporta erratas, ayuda a otros o comparte tu opinión. Sé constructivo.

Inicia sesión para comentar y responder.

cargando comentarios...