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.
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...