Procesamiento en streaming Avanzado¶
Procesar datos a medida que llegan, en lugar de acumularlos y procesarlos por lotes. Permite alertas en segundos, vistas siempre actualizadas y reacciones automáticas.
1. Batch vs. streaming¶
| Batch | Streaming | |
|---|---|---|
| Datos | Conjunto acotado (el fichero de ayer) | Flujo infinito |
| Latencia | Minutos u horas | Milisegundos o segundos |
| Complejidad | Menor: todo está disponible, se puede reintentar entero | Mayor: orden, retrasos, estado, fallos a mitad |
| Ejemplos | Informe diario de costes, reentrenar un modelo | Detección de anomalías en telemetría, alertas de fraude |
Batch es un caso particular de streaming
Un lote es un stream acotado. Por eso los motores modernos (Flink, Spark Structured Streaming, Beam) usan la misma API para ambos.
2. Tiempo de evento vs. tiempo de procesamiento¶
- Tiempo de evento: cuándo ocurrió (el timestamp que pone el sensor del sitio edge).
- Tiempo de procesamiento: cuándo lo procesa tu sistema.
Pueden diferir mucho: un sitio sin conexión durante una hora envía sus eventos después, tarde y desordenados. Si agregas por tiempo de procesamiento, esos eventos caen en la ventana equivocada. Para resultados correctos, se agrega por tiempo de evento.
3. Ventanas¶
| Ventana | Forma | Ejemplo |
|---|---|---|
| Fija (tumbling) | Intervalos consecutivos sin solapar | Errores por sitio cada 5 minutos |
| Deslizante (sliding / hopping) | Intervalos que se solapan | Media de CPU de los últimos 10 min, calculada cada minuto |
| Sesión | Se cierra tras un periodo de inactividad | Sesiones de un operador en el dashboard |
| Global | Todo el flujo, con disparadores propios | Contadores acumulados |
4. Watermarks y datos tardíos¶
Un watermark es la afirmación del sistema: "ya no espero eventos con tiempo anterior a T". Cuando el watermark supera el final de una ventana, esta se cierra y emite su resultado.
flowchart LR
E1["evento t=10:02"] --> W
E2["evento t=10:04"] --> W
E3["evento t=09:58<br/>(llega tarde)"] --> W
W["Ventana 10:00-10:05<br/>watermark = máx. visto − 2 min"] --> R["Resultado al pasar<br/>el watermark de 10:05"]
- Watermark conservador (espera mucho) → resultados más completos pero más tardíos.
- Agresivo → resultados rápidos, pero más eventos tardíos descartados o que obligan a corregir.
- Opciones para los tardíos: descartarlos, enviarlos a una salida lateral para tratarlos aparte, o actualizar el resultado ya emitido (allowed lateness).
5. Estado y tolerancia a fallos¶
Contar, agregar o unir flujos requiere estado (el contador por sitio, el último valor por clave). Si el procesador cae, ese estado no puede perderse.
- Checkpoints (Flink): instantáneas periódicas y consistentes del estado y de las posiciones de lectura; al recuperarse, vuelve al último checkpoint y reprocesa desde ahí → resultados exactly-once sobre el estado.
- Changelog topics (Kafka Streams): el estado local (RocksDB) se respalda en un topic compactado; otra instancia puede reconstruirlo.
6. Motores¶
| Motor | Modelo | Cuándo |
|---|---|---|
| Kafka Streams | Librería Java dentro de tu aplicación; escala con instancias; estado en RocksDB | Procesamiento de Kafka a Kafka sin clúster adicional |
| Apache Flink | Clúster de procesamiento con estado, tiempo de evento y checkpoints de primera clase | Streaming complejo y a gran escala; también SQL sobre streams |
| Spark Structured Streaming | Micro-lotes (o continuo) sobre Spark | Equipos que ya usan Spark; latencias de segundos |
| Apache Beam | Modelo de programación que se ejecuta sobre Flink, Spark o Dataflow | Portabilidad entre motores |
Ejemplo con Kafka Streams¶
// Número de eventos de error por sitio en ventanas de 5 minutos (tiempo de evento)
StreamsBuilder builder = new StreamsBuilder();
builder.stream("node-events", Consumed.with(Serdes.String(), nodeEventSerde)
.withTimestampExtractor((rec, prev) -> ((NodeEvent) rec.value()).occurredAt().toEpochMilli()))
.filter((siteId, event) -> event.level() == Level.ERROR)
.groupByKey()
.windowedBy(TimeWindows.ofSizeAndGrace(Duration.ofMinutes(5), Duration.ofMinutes(2))) // 2 min de gracia para tardíos
.count()
.toStream()
.filter((window, count) -> count > 20)
.map((window, count) -> KeyValue.pair(window.key(), new SiteAlert(window.key(), count, window.window().endTime())))
.to("site-alerts", Produced.with(Serdes.String(), siteAlertSerde));
7. Change Data Capture (CDC)¶
Convertir cada cambio de una base de datos en un evento, leyendo su log de transacciones (WAL de PostgreSQL, binlog de MySQL). Herramienta de referencia: Debezium (sobre Kafka Connect).
- Sin modificar la aplicación y sin consultas periódicas costosas.
- Base del patrón outbox: la aplicación escribe en una tabla
outboxy Debezium la publica en Kafka. - Permite alimentar cachés, índices de búsqueda y el lakehouse casi en tiempo real.
8. Arquitecturas Lambda y Kappa¶
| Arquitectura | Idea | Problema |
|---|---|---|
| Lambda | Una capa batch (exacta, lenta) + una capa streaming (aproximada, rápida), combinadas | Dos implementaciones de la misma lógica |
| Kappa | Solo streaming; para recalcular, se reprocesa el log desde el principio | Requiere retener el log y un motor capaz |
La tendencia actual se acerca a Kappa, con el log (Kafka) y el lakehouse como almacenamiento duradero.
Preguntas de repaso¶
¿Por qué agregar por tiempo de evento y no de procesamiento?
Porque el tiempo de procesamiento depende de retrasos de red, desconexiones y reintentos. Agregar por tiempo de evento asigna cada evento a la ventana en la que realmente ocurrió, y el resultado es reproducible.
¿Qué es un watermark?
Una marca que indica que el sistema no espera más eventos anteriores a cierto tiempo; permite cerrar ventanas y emitir resultados sabiendo cuánto retraso se tolera.
¿Qué aporta CDC frente a consultar la base de datos periódicamente?
Captura todos los cambios (incluidos borrados e intermedios), con baja latencia y sin cargar la base de datos con consultas repetidas, y en el orden en que se confirmaron.
Ejercicios¶
Ejercicio 1 · Básico — Elegir la ventana
Indica qué tipo de ventana usarías: (a) informe de despliegues por hora; (b) alerta si la media de latencia de los últimos 10 minutos supera 500 ms, evaluada cada minuto; (c) agrupar las acciones de un operador hasta que pase 15 min inactivo.
Solución
(a) Fija de 1 hora. (b) Deslizante de 10 minutos con avance de 1 minuto. (c) Sesión con 15 minutos de inactividad.
Ejercicio 2 · Medio — Sitios que se reconectan tarde
Los sitios edge pueden estar desconectados hasta 6 horas y luego envían sus eventos. Quieres un recuento de errores por sitio y hora que sea correcto. ¿Cómo configuras ventanas, watermark y datos tardíos?
Solución
Ventanas fijas de 1 h por tiempo de evento. Un watermark de 6 h retrasaría todos los resultados 6 h; mejor:
watermark corto (p. ej. 5 min) para emitir resultados pronto + allowed lateness de 6 h para actualizar
la ventana cuando lleguen tardíos (el destino debe aceptar actualizaciones: upsert por siteId + hora). Los que
lleguen aún más tarde, a una salida lateral para una corrección batch. Alternativa: calcular en streaming una
vista provisional y recalcular la definitiva en batch sobre el lakehouse.
Ejercicio 3 · Medio — Outbox con Debezium
Describe el flujo completo para publicar RolloutStarted de forma fiable usando una tabla outbox y Debezium.
Solución
- En la misma transacción:
INSERT INTO rollout …yINSERT INTO outbox(id, aggregate_type, aggregate_id, type, payload). - Debezium lee el WAL de PostgreSQL, detecta la inserción en
outboxy, con el Outbox Event Router, publica en el topicrollout.eventscon claveaggregate_id. - Los consumidores procesan de forma idempotente usando
iddel evento. - Una tarea borra filas antiguas de
outbox(Debezium ya las leyó del WAL). Resultado: el evento se publica si y solo si la transacción se confirmó, sin dual write.
Ejercicio 4 · Avanzado — Kafka Streams o Flink
Debes: (a) enriquecer eventos de nodos con la configuración del sitio y (b) detectar patrones complejos (3 reinicios seguidos de un fallo de despliegue en 30 min) sobre 2 millones de eventos/s. ¿Qué motor eliges para cada uno?
Solución
(a) Kafka Streams: unir el stream de eventos con una KTable (el topic compactado de configuración) es su
caso de uso natural; se despliega como un servicio más, sin clúster aparte.
(b) Flink: volumen muy alto, estado grande, tiempo de evento y detección de patrones (librería CEP) son su
punto fuerte; los checkpoints dan tolerancia a fallos con estado grande.