Saltar al contenido

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, KafkaProducer
2import json
3
4consumer = 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)
12
13producer = KafkaProducer(
14 bootstrap_servers=["localhost:9092"],
15 value_serializer=lambda v: json.dumps(v).encode("utf-8"),
16)
17
18def 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 }
28
29def clasificar_total(total: float) -> str:
30 if total > 500: return "premium"
31 if total > 100: return "standard"
32 return "micro"
33
34print("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 destino
2for msg in consumer:
3 pedido = msg.value
4 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 time
2from collections import defaultdict
3
4class VentanaTemporal:
5 """Agregación por ventana de tiempo sliding window."""
6
7 def __init__(self, window_seconds: int = 60):
8 self.window_seconds = window_seconds
9 self.eventos = [] # (timestamp, valor)
10
11 def agregar(self, valor: float, timestamp: float = None):
12 ts = timestamp or time.time()
13 self.eventos.append((ts, valor))
14 self._limpiar_expirados()
15
16 def _limpiar_expirados(self):
17 ahora = time.time()
18 self.eventos = [(ts, v) for ts, v in self.eventos
19 if ahora - ts <= self.window_seconds]
20
21 def total(self) -> float:
22 self._limpiar_expirados()
23 return sum(v for _, v in self.eventos)
24
25 def count(self) -> int:
26 self._limpiar_expirados()
27 return len(self.eventos)
28
29 def promedio(self) -> float:
30 c = self.count()
31 return self.total() / c if c > 0 else 0
32
33# Uso en un stream processor
34ventas_por_tienda = defaultdict(lambda: VentanaTemporal(window_seconds=300))
35
36for msg in consumer:
37 venta = msg.value
38 tienda = venta["tienda_id"]
39 ventas_por_tienda[tienda].agregar(venta["precio"])
40
41 # 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

Los patrones fundamentales del streaming se combinan como piezas de LEGO

### 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}
7
8def 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 }
18
19# 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 defaultdict
2import time
3
4class DetectorAnomalias:
5 """Detecta anomalías basadas en frecuencia en ventana de tiempo."""
6
7 def __init__(self, max_eventos: int = 5, ventana_seg: int = 120):
8 self.max_eventos = max_eventos
9 self.ventana_seg = ventana_seg
10 self.historial = defaultdict(list) # key -> [timestamps]
11
12 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 ventana
16 self.historial[key] = [
17 ts for ts in self.historial[key]
18 if ahora - ts <= self.ventana_seg
19 ]
20 # Agregar nuevo evento
21 self.historial[key].append(ahora)
22 # ¿Supera el umbral?
23 return len(self.historial[key]) > self.max_eventos
24
25# Uso en el stream processor
26detector = DetectorAnomalias(max_eventos=5, ventana_seg=120)
27
28for msg in consumer:
29 pedido = msg.value
30 cliente = pedido["cliente_id"]
31
32 if detector.registrar(cliente):
33 # ALERTA: posible fraude
34 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)")
42
43 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

[01]

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

Cargando editor...
[02]

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
Cargando editor...
[03]

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
Cargando editor...
[04]

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