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.

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