Batch vs. streaming
El procesamiento batch trabaja sobre conjuntos de datos acotados (bounded): un archivo, una tabla, el log de ayer. Se ejecuta periódicamente (cada hora, cada noche) y produce resultados completos sobre un conjunto cerrado.
El procesamiento streaming trabaja sobre datos no acotados (unbounded): un flujo continuo de eventos que nunca “termina”. El sistema procesa cada evento a medida que llega y emite resultados de forma incremental, con latencias de segundos o milisegundos.
| Aspecto | Batch | Streaming |
|---|---|---|
| Datos | Acotados, históricos | Continuos, en movimiento |
| Latencia | Minutos a horas | Milisegundos a segundos |
| Disparo | Programado (cron) | Por cada evento |
| Ejemplo | Informe diario de ventas | Detección de fraude al instante |
Streaming no reemplaza a batch: son complementarios. De hecho, la arquitectura Lambda (que veremos en la siguiente lección) los combina.
El problema del tiempo: event time vs. processing time
En un stream hay dos relojes distintos:
- Event time: el momento en que el evento ocurrió en el mundo real (timestamp que lleva el propio evento).
- Processing time: el momento en que el sistema procesa el evento (reloj del servidor).
¿Por qué importa la diferencia? Porque los eventos pueden llegar tarde y desordenados: un móvil sin cobertura guarda sus eventos y los envía una hora después; la red introduce retrasos variables. Si agrupamos por processing time, los resultados dependen del azar de la red; si agrupamos por event time, el resultado refleja la realidad — pero tenemos que decidir cuánto esperar a los rezagados.
Ventanas: trocear lo infinito
Para calcular agregados (contar, sumar, promediar) sobre un flujo infinito necesitamos acotarlo en ventanas:
| Tipo | Definición | Ejemplo |
|---|---|---|
| Tumbling (fija) | Ventanas de tamaño fijo, contiguas y sin solapamiento | “Clics por cada minuto”: [10:00–10:01), [10:01–10:02) |
| Sliding (deslizante) | Tamaño fijo, pero avanzan con un paso menor: se solapan | “Media de los últimos 5 minutos, recalculada cada minuto” |
| Session (sesión) | Agrupan eventos separados por menos de un gap de inactividad | “Sesiones de usuario: actividad seguida de 30 min sin eventos cierra la sesión” |
graph LR subgraph S["Stream de eventos (event time)"] direction LR E1["e1"] --- E2["e2"] --- E3["e3"] --- E4["e4"] --- E5["e5"] --- E6["e6"] --- E7["e7"] end subgraph W1["Ventana 1: 10:00-10:01"] E1 E2 E3 end subgraph W2["Ventana 2: 10:01-10:02"] E4 E5 end subgraph W3["Ventana 3: 10:02-10:03"] E6 E7 end E1 --> W1 E2 --> W1 E3 --> W1 E4 --> W2 E5 --> W2 E6 --> W3 E7 --> W3 W1 --> R1["count = 3"] W2 --> R2["count = 2"] W3 --> R3["count = 2"]
Watermarks: ¿cuándo damos por cerrada una ventana?
Si esperamos a los eventos tardíos, ¿cuánto hay que esperar? Los watermarks responden a esa pregunta: son una marca de progreso del event time que afirma “no espero más eventos con timestamp anterior a T”.
- Cuando el watermark supera el final de una ventana, esta se cierra y emite su resultado.
- Los eventos que lleguen después del watermark se consideran late events y se descartan, se envían a un canal aparte (side output) o se permiten dentro de un allowed lateness adicional.
Es un equilibrio: watermark muy conservador → más completitud pero más latencia; watermark agresivo → baja latencia pero más eventos descartados.
Garantías de entrega
¿Qué pasa si un nodo falla a mitad del procesamiento? Los sistemas de streaming ofrecen distintas semánticas de entrega:
| Semántica | Garantía | Consecuencia |
|---|---|---|
| At-most-once | Cada evento se procesa como máximo una vez | Puede haber pérdida de eventos |
| At-least-once | Ningún evento se pierde; se reintenta | Puede haber duplicados (hay que hacer el procesamiento idempotente) |
| Exactly-once | Cada evento afecta al resultado exactamente una vez | El ideal, pero exige coordinación (transacciones, checkpoints, sinks transaccionales) |
Exactly-once no significa que el evento “se lea una sola vez”, sino que su efecto sobre el estado y la salida se aplica una sola vez, incluso con reintentos y fallos intermedios.
Motores de stream processing
- Spark Structured Streaming: trata el stream como una tabla infinita a la que se le aplican consultas SQL/DataFrame. Usa un modelo de micro-batch (pequeños lotes cada N segundos) con checkpoints que permiten exactly-once. Ideal si ya usas el ecosistema Spark.
- Apache Flink: motor de streaming puro, evento a evento, con soporte de primer nivel para event time, watermarks, estado con checkpoints y exactly-once. Referencia para casos de baja latencia y lógica de ventanas compleja.
Ambos se integran con Kafka como fuente y destino: Kafka actúa como el “hub” de eventos y el motor como el cerebro que los procesa.
Ejercicio: ventana tumbling por minuto
🧪 Ejercicio
Contador de eventos por minuto (ventana tumbling)
Recibes una lista de eventos `(timestamp_unix, valor)` ya ordenados por event time. Implementa `ventana_tumbling(eventos, tam=60)` que agrupe los eventos en ventanas fijas de 60 segundos alineadas al minuto y devuelva un diccionario {inicio_ventana: cuenta}. Aplícalo a los eventos de ejemplo e imprime el resultado formateado.
🔍 El inicio de la ventana de un timestamp t es (t // tam) * tam.
from datetime import datetime, timezone
def ventana_tumbling(eventos, tam=60):
ventanas = {}
for ts, _valor in eventos:
inicio = (ts // tam) * tam
ventanas[inicio] = ventanas.get(inicio, 0) + 1
return ventanas
# Eventos: (timestamp unix, valor)
eventos = [
(1_700_000_005, "clic"),
(1_700_000_012, "clic"),
(1_700_000_041, "clic"),
(1_700_000_067, "clic"),
(1_700_000_099, "clic"),
(1_700_000_130, "clic"),
]
conteo = ventana_tumbling(eventos)
for inicio, n in sorted(conteo.items()):
hora = datetime.fromtimestamp(inicio, tz=timezone.utc).strftime("%H:%M:%S")
print(f"Ventana [{hora}, +60s): {n} eventos")Comprueba lo aprendido
Comprueba que lo pillaste
Un sistema de recomendación calcula la media de compras de los últimos 10 minutos, actualizándola cada 2 minutos. ¿Qué tipo de ventana usa?
La ventana tiene tamaño fijo (10 min) pero avanza con un paso menor (2 min), por lo que las ventanas se solapan: es una ventana sliding o deslizante.
Comprueba que lo pillaste
Un evento ocurrió a las 10:00:30 pero llega al sistema a las 10:02:10. ¿Cuál es su event time?
El event time es el momento en que ocurrió el evento en el mundo real (10:00:30), independientemente de cuándo lo procese el sistema (processing time, 10:02:10).
Comprueba que lo pillaste
¿Qué problema puede aparecer con semántica at-least-once si el procesamiento no es idempotente?
At-least-once reintenta ante fallos, así que un mismo evento puede procesarse más de una vez. Si la operación no es idempotente (por ejemplo, un contador que no deduplica), el resultado queda inflado por duplicados.
Resumen
- Batch procesa datos acotados con alta latencia; streaming procesa flujos continuos con baja latencia.
- Las ventanas acotan el flujo: tumbling (fijas sin solape), sliding (con solape) y session (por inactividad).
- Event time refleja cuándo ocurrió el evento; processing time, cuándo se procesó. Los eventos llegan tarde y desordenados.
- Los watermarks marcan hasta qué punto del event time esperamos eventos, y deciden cuándo cerrar ventanas.
- Las garantías van de at-most-once (pérdidas posibles) a exactly-once (efecto único), pasando por at-least-once (duplicados posibles).
- Spark Structured Streaming (micro-batch) y Flink (streaming puro) son los motores de referencia, usualmente sobre Kafka.
📚 Lecturas y fuentes
| Recurso | Tipo | Por qué leerlo |
|---|---|---|
| Procesamiento por lotes (Wikipedia en español) | Artículo | Para fijar el contraste con streaming: lee la introducción y “Ventajas”; el apartado final compara con los sistemas de tiempo real. |
| Streaming 101: The World Beyond Batch (Tyler Akidau) | Artículo | La mejor explicación de event time vs. processing time. Lee hasta el final de “Windowing” y deja “Streaming 102” para otro día. (en inglés) |
| Structured Streaming Programming Guide (Spark) | Docs | Ver los conceptos en código. Solo “Basic Concepts” y “Handling Late Data and Watermarking”; ignora las secciones de sources y sinks. (en inglés) |
| Designing Data-Intensive Applications (Kleppmann) | Libro | Capítulo 11 “Stream Processing”: la parte “Reasoning About Time” y las garantías de entrega (at-least-once, exactly-once) con más rigor. (en inglés) |