lección 7
Streaming real: procesamiento continuo de eventos
Procesa flujos de datos en tiempo real: transformaciones, agregaciones por ventana, enriquecimiento y escritura a múltiples destinos.
⏱ 50 min
### De consumidor simple a procesador de streams
Hasta ahora hemos consumido mensajes uno a uno y hecho operaciones simples. Pero el streaming real va mucho más allá: transformar datos en vuelo, agregar por ventanas de tiempo, enriquecer con datos externos, detectar patrones y escribir resultados en múltiples destinos. Es la diferencia entre leer un periódico y operar una sala de control que monitoriza mil señales simultáneas.
En producción real, el streaming se usa para: dashboards en tiempo real, detección de fraude, recomendaciones instantáneas, alertas operacionales, sincronización de datos entre sistemas y ETL continuo (a diferencia del batch de siempre). Lo que antes tardaba horas en procesarse con Airflow, ahora fluye en segundos.
### Patrón 1: Transformación en vuelo (map)
El patrón más simple de streaming: leer de un topic, transformar cada mensaje y escribir el resultado en OTRO topic. Es como una fábrica: materia prima entra por un extremo, producto terminado sale por el otro. El topic de entrada no se modifica — es inmutable. El resultado va a un topic nuevo.
1from kafka import KafkaConsumer, KafkaProducer2import json34consumer = KafkaConsumer(5 "pedidos-raw",6 bootstrap_servers=["localhost:9092"],7 group_id="transformador-pedidos",8 auto_offset_reset="earliest",9 enable_auto_commit=False,10 value_deserializer=lambda m: json.loads(m.decode("utf-8")),11)1213producer = KafkaProducer(14 bootstrap_servers=["localhost:9092"],15 value_serializer=lambda v: json.dumps(v).encode("utf-8"),16)1718def transformar_pedido(raw: dict) -> dict:19 """Enriquece y normaliza un pedido raw."""20 return {21 "pedido_id": raw["pedido_id"],22 "cliente_id": raw["cliente_id"],23 "total_con_iva": round(raw["total"] * 1.21, 2),24 "moneda": "EUR",25 "categoria": clasificar_total(raw["total"]),26 "procesado_at": time.time(),27 }2829def clasificar_total(total: float) -> str:30 if total > 500: return "premium"31 if total > 100: return "standard"32 return "micro"3334print("Stream processor: pedidos-raw → pedidos-enriched")35for msg in consumer:36 enriched = transformar_pedido(msg.value)37 producer.send("pedidos-enriched", value=enriched)38 consumer.commit()39 print(f" {enriched['pedido_id']}: {enriched['total_con_iva']}EUR ({enriched['categoria']})")
Transformación en vuelo: leer de un topic, transformar, escribir en otro
### Patrón 2: Filtrado (filter)
A veces solo necesitas un subconjunto de los eventos. El patrón filter lee de un topic y solo escribe en el destino los mensajes que cumplen una condición. Ejemplo: de todos los pedidos, enviar a un topic "pedidos-vip" solo los mayores de 500€ para tratamiento prioritario.
1# Filtrado: solo pedidos VIP (> 500 EUR) van al topic destino2for msg in consumer:3 pedido = msg.value4 if pedido["total"] > 500:5 producer.send("pedidos-vip", value=pedido)6 print(f" VIP: {pedido['pedido_id']} -> {pedido['total']}EUR")7 else:8 print(f" Skip: {pedido['pedido_id']} ({pedido['total']}EUR)")9 consumer.commit()
Filtrado: solo los eventos que cumplen la condición pasan al topic destino
### Patrón 3: Agregación por ventana de tiempo
Este es el patrón más poderoso del streaming: agregar datos en ventanas de tiempo. "Cuántas ventas en los últimos 5 minutos". "Total facturado por hora". "Promedio de tiempo de respuesta cada 30 segundos". En batch calcularías esto al final del día; en streaming lo tienes CONSTANTEMENTE actualizado.
La implementación naive: acumulas en memoria y "flush" cada N segundos. Es simple pero no tolerante a fallos (si mueres, pierdes el estado). En producción usarías Kafka Streams o Flink con state stores persistentes. Pero para entender el concepto, la versión en Python puro es perfecta.
1import time2from collections import defaultdict34class VentanaTemporal:5 """Agregación por ventana de tiempo sliding window."""67 def __init__(self, window_seconds: int = 60):8 self.window_seconds = window_seconds9 self.eventos = [] # (timestamp, valor)1011 def agregar(self, valor: float, timestamp: float = None):12 ts = timestamp or time.time()13 self.eventos.append((ts, valor))14 self._limpiar_expirados()1516 def _limpiar_expirados(self):17 ahora = time.time()18 self.eventos = [(ts, v) for ts, v in self.eventos19 if ahora - ts <= self.window_seconds]2021 def total(self) -> float:22 self._limpiar_expirados()23 return sum(v for _, v in self.eventos)2425 def count(self) -> int:26 self._limpiar_expirados()27 return len(self.eventos)2829 def promedio(self) -> float:30 c = self.count()31 return self.total() / c if c > 0 else 03233# Uso en un stream processor34ventas_por_tienda = defaultdict(lambda: VentanaTemporal(window_seconds=300))3536for msg in consumer:37 venta = msg.value38 tienda = venta["tienda_id"]39 ventas_por_tienda[tienda].agregar(venta["precio"])4041 # Emitir métrica cada mensaje (o cada N mensajes)42 print(f" [{tienda}] Últimos 5min: {ventas_por_tienda[tienda].count()} ventas, "43 f"total: {ventas_por_tienda[tienda].total():.2f}EUR")44 consumer.commit()
Ventana temporal: agrega datos de los últimos N segundos/minutos
### Patrón 4: Enriquecimiento con datos externos
Un evento llega con un cliente_id pero necesitas su nombre, su segmento y su historial. El patrón enrichment hace un lookup a una base de datos o caché para ENRIQUECER el evento antes de procesarlo. Es como un formulario pre-rellenado: el evento trae lo mínimo y tú completas los datos.
1# Simulamos una "base de datos" de clientes (en producción sería Redis o PostgreSQL)2clientes_db = {3 "CLI-1": {"nombre": "María García", "segmento": "premium", "antiguedad_meses": 24},4 "CLI-2": {"nombre": "Carlos López", "segmento": "standard", "antiguedad_meses": 6},5 "CLI-3": {"nombre": "Ana Martínez", "segmento": "nuevo", "antiguedad_meses": 1},6}78def enriquecer_pedido(pedido: dict) -> dict:9 """Enriquece un pedido con datos del cliente."""10 cliente_info = clientes_db.get(pedido["cliente_id"], {})11 return {12 **pedido,13 "cliente_nombre": cliente_info.get("nombre", "Desconocido"),14 "cliente_segmento": cliente_info.get("segmento", "unknown"),15 "cliente_antiguedad": cliente_info.get("antiguedad_meses", 0),16 "es_cliente_leal": cliente_info.get("antiguedad_meses", 0) > 12,17 }1819# En el stream processor:20for msg in consumer:21 pedido_enriched = enriquecer_pedido(msg.value)22 producer.send("pedidos-enriched", value=pedido_enriched)23 print(f" {pedido_enriched['pedido_id']}: {pedido_enriched['cliente_nombre']} "24 f"({pedido_enriched['cliente_segmento']})")
Enriquecimiento: completar el evento con datos de una fuente externa
### Patrón 5: Detección de anomalías en tiempo real
Uno de los casos de uso más potentes del streaming: detectar patrones sospechosos EN EL MOMENTO. Ejemplo: un cliente hace 5 compras en 2 minutos (posible fraude). O el tiempo de respuesta de una API supera los 5 segundos durante 1 minuto (posible degradación). Sin streaming, esto lo detectarías en el reporte del día siguiente — demasiado tarde.
1from collections import defaultdict2import time34class DetectorAnomalias:5 """Detecta anomalías basadas en frecuencia en ventana de tiempo."""67 def __init__(self, max_eventos: int = 5, ventana_seg: int = 120):8 self.max_eventos = max_eventos9 self.ventana_seg = ventana_seg10 self.historial = defaultdict(list) # key -> [timestamps]1112 def registrar(self, key: str) -> bool:13 """Registra un evento y retorna True si es anómalo."""14 ahora = time.time()15 # Limpiar eventos fuera de la ventana16 self.historial[key] = [17 ts for ts in self.historial[key]18 if ahora - ts <= self.ventana_seg19 ]20 # Agregar nuevo evento21 self.historial[key].append(ahora)22 # ¿Supera el umbral?23 return len(self.historial[key]) > self.max_eventos2425# Uso en el stream processor26detector = DetectorAnomalias(max_eventos=5, ventana_seg=120)2728for msg in consumer:29 pedido = msg.value30 cliente = pedido["cliente_id"]3132 if detector.registrar(cliente):33 # ALERTA: posible fraude34 alerta = {35 "tipo": "frecuencia_alta",36 "cliente_id": cliente,37 "pedidos_en_ventana": len(detector.historial[cliente]),38 "timestamp": time.time(),39 }40 producer.send("alertas-fraude", value=alerta)41 print(f" 🚨 ALERTA FRAUDE: {cliente} ({len(detector.historial[cliente])} pedidos en 2min)")4243 consumer.commit()
Detección de anomalías: alertar cuando la frecuencia supera un umbral
### Escribir a múltiples destinos (fan-out de stream)
En un stream processor real, el resultado no siempre va a UN solo topic. Puede ir a múltiples destinos: un topic para analytics, otro para alertas, otro para un data lake. El stream processor decide a dónde va cada resultado basándose en su contenido.
Consejo de senior: mantén tus stream processors SIMPLES. Cada uno hace UNA cosa bien. Si necesitas transformar Y agregar Y alertar, son 3 processors leyendo del mismo topic — no uno gigante que hace todo. La composabilidad es la clave del streaming mantenible.
Lo que le diría a mi yo de hace 5 años: el estado en memoria (como nuestra VentanaTemporal) es frágil. Si el proceso muere, pierdes el estado. Para producción, usa Kafka Streams con RocksDB (state stores persistentes) o una base de datos externa como Redis. El streaming sin estado es fácil; el streaming CON estado es donde está la complejidad real.
Cuidado con el enriquecimiento síncrono: si tu lookup a la DB tarda 50ms por mensaje y tienes 10.000 msg/s, necesitas 500 segundos de procesamiento por segundo (imposible). Usa caché agresiva, pre-carga datos en memoria o haz lookups asíncronos con batching.
## ejercicios
Stream processor: transformar y enriquecer pedidos
Escribe un stream processor que lea del topic "pedidos-raw", calcule el IVA (21%), clasifique el pedido (micro/standard/premium), enriquezca con datos del cliente y escriba en "pedidos-processed". Qué deberías ver: una línea por pedido procesado con pedido_id, nombre del cliente, total con IVA y categoría. Al final, "Procesados: N pedidos".
Agregación por ventana: ventas por minuto
Implementa un stream processor que calcule el total de ventas por tienda en ventanas de 60 segundos. Cada vez que llega un mensaje, imprime el total acumulado de la ventana actual para esa tienda.
💡 Resultado esperado
[MAD] Ventana 60s: 1 ventas, total: 120.00EUR [BCN] Ventana 60s: 1 ventas, total: 85.50EUR [MAD] Ventana 60s: 2 ventas, total: 320.00EUR [MAD] Ventana 60s: 3 ventas, total: 370.00EUR
Detector de fraude: más de 3 compras en 2 minutos
Implementa un detector que alerte cuando un cliente hace más de 3 compras en 2 minutos. Lee del topic "pedidos" y publica alertas en "alertas-fraude".
💡 Resultado esperado
OK: CLI-A - pedido O1 OK: CLI-A - pedido O2 OK: CLI-B - pedido O3 OK: CLI-A - pedido O4
Routing: enviar a diferentes topics según contenido
Escribe un stream processor que lea del topic "eventos-raw" y enrute cada evento a un topic diferente según su tipo: "PedidoCreado" → "topic-pedidos", "PagoConfirmado" → "topic-pagos", cualquier otro → "topic-otros".
💡 Resultado esperado
PedidoCreado -> topic-pedidos PagoConfirmado -> topic-pagos UsuarioLogout -> topic-otros PedidoCreado -> topic-pedidos
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...