Todo lo calculado hasta ahora en el módulo partía de un conjunto cerrado: el fichero del 14 de septiembre, los clics de siete días. Pero los datos de Kilómetro Cero no nacen cerrados. pedidos.eventos recibe cientos de eventos stock.actualizado por segundo durante una campaña, los 140 repartidores (como furgoneta-3) envían 2,4 millones de posiciones al día, y hay preguntas que no pueden esperar al lote nocturno: qué productos se están quedando sin stock en Girona ahora, dónde está cada repartidor ahora, qué repartidor lleva tres minutos sin dar señal. El procesamiento de flujos responde a esas preguntas ejecutando el cálculo de forma continua, evento a evento, sobre un conjunto que nunca termina. Eso rompe suposiciones que en los lotes eran gratuitas: no se puede "esperar a tener todo" antes de agrupar, los eventos llegan desordenados y con retraso respecto al momento en que ocurrieron, el estado del cálculo debe sobrevivir a los fallos sin poder reejecutar "desde el principio", y la entrada puede llegar más rápido de lo que se procesa. Esta lección da nombre a cada uno de esos problemas (tiempo de evento, marcas de agua, ventanas, estado y checkpoints, exactly-once, backpressure), compara las arquitecturas Lambda y Kappa y las tres herramientas habituales (Kafka Streams, Flink, Spark Structured Streaming), y resuelve los dos casos de analitica en tiempo real: la alerta de stock bajo y el panel de reparto. La entrega de esos resultados al navegador del operador queda para 08-02; aquí terminamos en un tópico de Kafka.
Contenido
- De los lotes a los flujos: qué cambia
- Tiempo de evento, tiempo de procesamiento y marcas de agua
- Ventanas: tumbling, sliding y session
- Estado y checkpoints
- Exactly-once en streaming
- Backpressure
- Arquitecturas Lambda y Kappa
- Herramientas: Kafka Streams, Flink y Spark Structured Streaming
- Los casos de Kilómetro Cero: stock bajo y panel de reparto
- Práctica: simulación en Python, PyFlink y Structured Streaming
- Errores Comunes y Consejos
- Ejercicios
- Conclusión
- De los lotes a los flujos: qué cambia
Recuperamos la tabla de 05-01 y la afinamos con lo aprendido desde entonces:
| Lote | Flujo | |
|---|---|---|
| Entrada | Acotada: un fichero, un directorio, "el día 14" | No acotada: un tópico de Kafka que no termina |
| Unidad | El conjunto entero | El evento: un hecho inmutable con marca de tiempo (pedido.creado, stock.actualizado, una posición) |
| Cuándo se produce la salida | Al terminar | Continuamente: cada evento, cada ventana, cada segundo |
| "Todos los datos" | Existe: se espera a leerlos | No existe: hay que decidir cuándo se ha visto suficiente (marcas de agua) |
| Orden | Se puede ordenar todo antes de calcular | Los eventos llegan desordenados y tarde |
| Estado | Vive en el job; si falla, se reejecuta desde la entrada | Vive para siempre en el operador; hay que persistirlo (checkpoints) |
| Tolerancia a fallos | Reejecutar el job (05-02, 05-03) | Restaurar el estado y reanudar desde el offset del checkpoint |
| Latencia | Minutos a horas | Milisegundos a segundos |
| Kilómetro Cero | ventas_diarias.py |
Alerta de stock bajo, panel de reparto |
Un motor de flujos ejecuta el mismo tipo de DAG de operadores que Spark (05-03), pero desplegado de forma permanente: cada operador es un proceso (o varios, en paralelo por clave) que espera eventos, los procesa y emite hacia el siguiente, sin que el DAG termine nunca. La fuente típica es Kafka (02-04), que aporta lo que un flujo necesita: particiones (paralelismo), offsets (posición reanudable) y retención (releer desde atrás). El sink es otro tópico, una base de datos o un panel.
Hay dos formas de mover eventos por el DAG: evento a evento (Flink, Kafka Streams: latencia de milisegundos) y por microlotes (Spark Structured Streaming: cada pocos cientos de milisegundos o segundos se toma lo que ha llegado y se ejecuta como un pequeño lote). La segunda reutiliza toda la maquinaria de lotes y simplifica exactly-once; la primera da latencias menores. Para el panel de reparto y la alerta de stock, segundos son suficientes, así que ambas sirven.
- Tiempo de evento, tiempo de procesamiento y marcas de agua
En 01-05 distinguimos el momento en que algo ocurrió del momento en que otro nodo se entera. En flujos esa distinción es la que más consecuencias tiene:
- Tiempo de evento (event time): cuándo ocurrió el hecho, según el reloj de quien lo generó. Es el
fecha_msde la envoltura de 02-05. La posición defurgoneta-3tomada a las 10:04:58 tiene ese tiempo aunque el móvil la envíe a las 10:07 por falta de cobertura. - Tiempo de procesamiento (processing time): cuándo el operador ve el evento, según su reloj. Depende de la cola en Kafka, de la red, de la carga.
La pregunta "¿cuántos stock.actualizado de queso-curado hubo entre las 10:00 y las 10:05?" solo tiene una respuesta correcta en tiempo de evento. Con tiempo de procesamiento, un reinicio del consumidor a las 10:03 metería en la ventana 10:05–10:10 eventos ocurridos a las 10:02; el resultado dependería del comportamiento del sistema, no de los hechos. Pero calcular en tiempo de evento crea un problema nuevo: cuando llega un evento con fecha_ms = 10:04:58 a las 10:07, ¿la ventana 10:00–10:05 ya se había emitido? ¿Hay que reabrirla? ¿Cuánto esperar antes de darla por cerrada?
La respuesta es la marca de agua (watermark): una afirmación que el sistema hace fluir por el DAG que dice "no espero más eventos con tiempo de evento anterior a T". Se genera en la fuente con una heurística, normalmente máximo tiempo de evento visto − retraso tolerado: con un retraso tolerado de 2 minutos, tras ver un evento de las 10:07:00 la marca de agua es 10:05:00, y en ese momento la ventana 10:00–10:05 se considera completa y se emite. El retraso tolerado es una decisión de negocio y de medición: cuánto tardan de verdad los eventos en llegar (el percentil 99 del desfase llegada − fecha_ms), contra cuánta latencia se acepta en el resultado. Poco retraso tolerado: resultados rápidos que descartan más eventos; mucho: resultados más completos pero tardíos.
Los eventos que llegan después de que la marca de agua haya pasado su ventana son eventos tardíos (late events). Cada motor ofrece tres destinos: descartarlos (el comportamiento por defecto), incorporarlos reemitiendo la ventana actualizada durante un margen adicional (allowed lateness en Flink, que obliga a mantener la ventana en el estado más tiempo), o desviarlos a una salida lateral (side output) para contarlos, registrarlos o reprocesarlos por lotes. Contar los tardíos es obligatorio: es la métrica que dice si el retraso tolerado está bien elegido.
flowchart LR
subgraph W1[Ventana 10:00–10:05]
e1[e1 10:01]
e2[e2 10:03]
e4[e4 10:04:58<br/>llega a las 10:08]
end
subgraph W2[Ventana 10:05–10:10]
e3[e3 10:06]
e5[e5 10:07]
end
e3 -. "marca de agua = 10:06 − 2 min = 10:04<br/>W1 sigue abierta" .-> WM1[ ]
e5 -. "marca de agua = 10:07 − 2 min = 10:05<br/>W1 se cierra y se emite" .-> WM2[ ]
e4 -. "llega tras el cierre: TARDÍO" .-> L[descartar / reemitir / salida lateral]
Con varias particiones de Kafka, cada partición tiene su propia marca de agua y la del operador es el mínimo de todas: una partición sin tráfico (un mercado cerrado por la noche) retiene la marca de agua global y bloquea la emisión de ventanas de todos. Los motores lo tratan con idle timeouts que excluyen las particiones inactivas del mínimo.
- Ventanas: tumbling, sliding y session
Sobre un flujo infinito, cualquier agregación ("cuántos", "el mínimo") necesita un límite: la ventana. Las tres formas básicas:
| Ventana | Definición | Un evento pertenece a | Ejemplo en Kilómetro Cero | Tamaño del estado |
|---|---|---|---|---|
| Tumbling (fija, saltos) | Intervalos consecutivos y disjuntos de tamaño fijo: 10:00–10:05, 10:05–10:10 | Exactamente una ventana | Stock mínimo por producto y mercado cada 5 minutos | Una ventana abierta por clave (más las que esperan la marca de agua) |
| Sliding (deslizante) | Tamaño fijo, avance menor: cada 1 minuto, los últimos 5 | Varias ventanas (tamaño / avance) | Distancia recorrida por furgoneta-3 en los últimos 5 minutos, refrescada cada minuto |
Tamaño/avance ventanas por clave: 5 |
| Session (sesión) | Sin tamaño fijo: se abre con un evento y se cierra tras un gap sin eventos | Una sesión, que crece | "Repartidor sin señal": sesión de posiciones que se cierra tras 3 minutos sin recibir ninguna | Una sesión abierta por clave; las sesiones se fusionan si un evento tardío las une |
| Global + trigger | Toda la historia, con disparadores explícitos | Una | Contador total de pedidos del día, emitido cada 10 s | Un acumulador por clave |
Las ventanas se combinan con una clave: "tumbling de 5 minutos por (producto, mercado)" mantiene una ventana por cada pareja, y el paralelismo del operador es por clave (todas las queso-curado/girona van a la misma instancia, como en el shuffle de 05-02). Y con el tiempo de evento: es el fecha_ms el que decide en qué ventana cae un evento, y la marca de agua la que decide cuándo se emite.
flowchart TB
subgraph T[Tumbling 5 min]
direction LR
t1[10:00–10:05] --- t2[10:05–10:10] --- t3[10:10–10:15]
end
subgraph S[Sliding 5 min cada 1 min]
direction LR
s1[10:00–10:05]
s2[10:01–10:06]
s3[10:02–10:07]
end
subgraph G[Session gap 3 min]
direction LR
g1[10:00:10 … 10:04:50] -- "gap > 3 min: sin señal" --- g2[10:09:30 … 10:21:00]
end
- Estado y checkpoints
Una ventana abierta, un contador por clave, la última posición conocida de cada repartidor: todo eso es estado, y en un flujo vive indefinidamente en el operador. Los motores lo guardan en un state backend local al operador (memoria, o RocksDB en disco local para estados de gigabytes) particionado por clave, exactamente como un actor por clave (05-01). El problema es la durabilidad: si el nodo muere, su estado local se pierde, y no se puede "reejecutar desde el principio" porque el principio fue hace meses.
La solución es el checkpoint: periódicamente (cada 10 s, cada minuto) el motor escribe una copia consistente del estado de todos los operadores, junto con los offsets de Kafka que ese estado refleja, en almacenamiento duradero (HDFS, MinIO). Ante un fallo, restaura el último checkpoint en nodos sanos y reanuda el consumo desde esos offsets: los eventos posteriores al checkpoint se vuelven a procesar, y el estado acaba siendo el mismo que si no hubiera habido fallo.
La palabra difícil es consistente: el estado de todos los operadores debe corresponder al mismo punto del flujo, aunque cada uno vaya por un evento distinto. Flink lo consigue con el algoritmo de Chandy-Lamport adaptado (las barreras de checkpoint): la fuente inyecta en el flujo un marcador con el número de checkpoint; cada operador, al recibirlo por todas sus entradas, guarda su estado y reenvía el marcador; el estado guardado refleja exactamente los eventos anteriores al marcador. Es una instantánea distribuida sin detener el flujo. Spark Structured Streaming lo tiene más fácil: cada microlote es una unidad, y el checkpoint registra qué microlotes se han completado con sus rangos de offsets. Kafka Streams guarda el estado en tópicos de Kafka (changelog topics) y lo reconstruye releyéndolos.
Un savepoint es un checkpoint disparado a mano, con formato estable, que se usa para parar el trabajo, cambiar el código o el paralelismo, y reanudar desde el mismo estado: el equivalente de un despliegue sin perder "los últimos cinco minutos". Y el tamaño del estado importa: una session window por repartidor es pequeña, pero "todos los pedidos de las últimas 24 horas por cliente" son gigabytes que hay que checkpointear cada minuto; los checkpoints incrementales (solo lo cambiado desde el anterior) y un TTL para el estado que no se toca son las herramientas.
- Exactly-once en streaming
En 02-05 concluimos que "exactly-once" en la entrega no existe, y que lo alcanzable es at-least-once más consumidores idempotentes. En streaming el término se usa con un significado preciso y alcanzable: el estado del motor refleja cada evento exactamente una vez, aunque haya fallos y reprocesamientos. Se consigue con el checkpoint del apartado anterior: tras un fallo, el estado vuelve al del checkpoint y los eventos posteriores se reaplican; como el estado que reflejaba esos eventos se ha descartado, no hay doble conteo.
Lo que el checkpoint no cubre es lo que ya salió del motor: la alerta escrita en el tópico inventario.alertas, la fila insertada en PostgreSQL, antes del fallo y después del último checkpoint. Al reprocesar, esos efectos se producen otra vez. Para que el resultado externo sea también exactly-once (end-to-end), el sink tiene que ser una de dos cosas:
- Idempotente. Escribir de nuevo el mismo resultado no cambia nada:
UPSERTpor clave(ventana, producto, mercado)en PostgreSQL, unPUTen Redis con la misma clave, un producer de Kafka con una clave y un consumidor idempotente aguas abajo (02-05). Es la opción más simple y la que se recomienda siempre que la salida tenga una clave natural. - Transaccional. El sink escribe en una transacción que solo se confirma cuando el checkpoint se completa (two-phase commit sink): Flink con el producer transaccional de Kafka (
DeliveryGuarantee.EXACTLY_ONCE), o con una tabla y una transacción por checkpoint. Entre el pre-commit y el commit los datos existen pero no son visibles para consumidores conisolation.level=read_committed. Añade latencia (el intervalo de checkpoint) y complejidad; se reserva para sinks sin clave natural (un log de eventos en Kafka).
| Garantía | Qué ocurre tras un fallo | Cómo se consigue |
|---|---|---|
| At-most-once | Se pierden eventos entre el fallo y la reanudación | Confirmar offsets antes de procesar; sin checkpoint de estado |
| At-least-once | Se reprocesan eventos; el estado o el sink pueden contarlos dos veces | Checkpoint de offsets sin coordinación con el estado; o sink no idempotente |
| Exactly-once (estado) | El estado es el mismo que sin fallo | Checkpoint consistente de estado + offsets |
| Exactly-once end-to-end | Además, el sink no muestra duplicados | Lo anterior + sink idempotente o transaccional |
- Backpressure
Un DAG de flujos es una cadena de productores y consumidores, y en cualquier momento uno puede ir más despacio que el anterior: el operador de ventanas escribiendo un checkpoint grande, el sink de PostgreSQL saturado, un pico de la Semana del Queso Artesano que triplica los stock.actualizado. Si el operador rápido siguiera enviando, las colas intermedias crecerían hasta agotar la memoria. Backpressure (contrapresión) es el mecanismo por el que la lentitud se propaga hacia atrás: el operador lento deja de aceptar, el anterior llena su búfer de salida y deja de leer de su entrada, y así hasta la fuente, que deja de consumir de Kafka. Los eventos se acumulan en Kafka (que está diseñado para ello, con retención de días), no en la memoria del motor, y el lag del grupo de consumidores (02-04) se convierte en la métrica que indica que el flujo no da abasto.
Flink lo implementa con créditos entre tareas (el receptor anuncia cuánto puede recibir); Kafka Streams lo tiene gratis porque cada instancia hace poll solo cuando ha terminado con el lote anterior; Spark Structured Streaming lo aproxima limitando cuántos offsets lee por microlote (maxOffsetsPerTrigger). Lo que ninguno hace es resolver la causa: si el lag crece de forma sostenida, hay que aumentar el paralelismo (más particiones y más instancias), aligerar el operador o aceptar resultados aproximados. Y ojo con la marca de agua: bajo backpressure, el tiempo de procesamiento se aleja del tiempo de evento, pero las ventanas siguen siendo correctas porque se definen en tiempo de evento; eso es precisamente lo que el apartado 2 compraba.
- Arquitecturas Lambda y Kappa
Cuando el streaming era nuevo y poco fiable, la respuesta a "quiero resultados en tiempo real pero también exactos" fue la arquitectura Lambda (Marz, 2011): mantener dos caminos. La capa batch recalcula cada noche, desde el lago, las vistas completas y exactas; la capa de velocidad calcula con streaming las últimas horas, de forma aproximada; una capa de servicio combina ambas al consultar. Funciona, pero obliga a escribir y mantener la misma lógica dos veces, en dos motores, con dos semánticas, y a reconciliar sus diferencias. La arquitectura Kappa (Kreps, 2014) propone un solo camino: todo es un flujo, con Kafka reteniendo el historial (o un lago de eventos releíble), y el "lote" es simplemente reprocesar el flujo desde un offset antiguo con la misma aplicación de streaming, escribiendo en una tabla nueva y cambiando el puntero cuando alcanza el presente.
| Lambda | Kappa | |
|---|---|---|
| Caminos | Dos: batch (exacto, lento) + velocidad (aproximado, rápido) | Uno: streaming, con reprocesamiento desde el historial |
| Código | Duplicado en dos motores | Una sola aplicación |
| Reprocesar | Relanzar el lote | Relanzar la aplicación desde un offset o desde el lago |
| Exactitud | Batch corrige a velocidad | El streaming debe ser exacto (tiempo de evento, exactly-once) |
| Requisitos | Un lago y un motor de lotes; un motor de flujos | Retención larga en Kafka o lago releíble; motor de flujos con estado |
| Cuándo | Lógica batch compleja que no cabe en streaming (entrenar ALS); resultados históricos que necesitan todo el conjunto | Agregaciones, alertas, materializaciones; cuando la latencia importa y la lógica es la misma |
| Kilómetro Cero | Ventas diarias (batch) + panel en tiempo real (velocidad): Lambda de facto | Alerta de stock y panel de reparto: Kappa puro |
La plataforma de datos de Kilómetro Cero acaba siendo pragmáticamente mixta: ventas_diarias.py y ALS son lotes porque necesitan el conjunto completo y no tienen prisa; el stock y el reparto son flujos porque la latencia es el requisito. Lo que Kappa aporta es el criterio: si la misma lógica se escribe dos veces, algo está mal; que los motores modernos ejecuten el mismo código en lote y en flujo (Spark, Flink) hace que la elección sea de despliegue, no de reescritura.
- Herramientas: Kafka Streams, Flink y Spark Structured Streaming
| Kafka Streams | Apache Flink | Spark Structured Streaming | |
|---|---|---|---|
| Qué es | Una librería Java/Scala: la aplicación es un proceso normal que consume y produce en Kafka | Un motor con clúster propio (JobManager + TaskManagers), o sobre YARN/Kubernetes | El modo streaming del motor Spark (05-03) |
| Modelo | Evento a evento | Evento a evento | Microlotes (100 ms–segundos); modo continuo experimental |
| Fuentes/sinks | Solo Kafka (por diseño) | Kafka, ficheros, JDBC, Kinesis, Pulsar, CDC... | Kafka, ficheros, sockets; sinks Kafka, ficheros, foreachBatch para todo lo demás |
| Tiempo de evento y marcas de agua | Sí, con grace period | Sí, el más completo (allowed lateness, side outputs, temporizadores) | withWatermark; sin allowed lateness ni side outputs |
| Ventanas | Tumbling, hopping, sliding, session | Todas, más ventanas definidas por el usuario | Tumbling, sliding, session |
| Estado | RocksDB local + changelog en Kafka | RocksDB o heap; checkpoints incrementales; savepoints | Estado por microlote en HDFS; RocksDB desde 3.2 |
| Exactly-once | Transacciones de Kafka (processing.guarantee=exactly_once_v2) |
Checkpoint + 2PC sinks | Checkpoint + sinks idempotentes; Kafka sink at-least-once |
| Lenguajes | Java, Scala (Kotlin) | Java, Scala, Python (PyFlink), SQL | Scala, Java, Python, R, SQL |
| Latencia típica | ms | ms | segundos |
| Encaja | Microservicios que transforman tópicos; equipos Java; sin clúster nuevo | Streaming exigente: estado grande, latencia baja, semántica precisa | Equipos que ya usan Spark; lote y flujo con el mismo código |
| Kilómetro Cero | Sería natural dentro de inventario (Java), pero los servicios son Python |
Panel de reparto y alertas (PyFlink) |
Alternativa con ventas_diarias.py reutilizado |
Elegimos Flink como motor de flujos de analitica por la semántica de tiempo de evento y por PyFlink, y mantenemos Structured Streaming como alternativa porque reutiliza la API de 05-03. Kafka Streams queda nombrado: es la opción correcta si el streaming vive dentro de un servicio Java y no en una plataforma de datos.
- Los casos de Kilómetro Cero: stock bajo y panel de reparto
Caso (a): alerta de stock bajo. inventario publica stock.actualizado en pedidos.eventos con cada cambio (04-05 los usaba para invalidar Redis). Datos del evento: producto, mercado, stock_actual, delta, replica (inv-bcn o inv-vlc). Requisito: cada 5 minutos, por producto y mercado, si el stock mínimo observado ha bajado de 10 unidades, emitir una alerta en el tópico inventario.alertas con la ventana, el mínimo y cuántas actualizaciones hubo. Ventana tumbling de 5 minutos por (producto, mercado), en tiempo de evento (el fecha_ms de inventario, porque una réplica con retraso no debe desplazar la alerta), con marca de agua de 2 minutos (medido: el percentil 99 del desfase es 40 s). Sink idempotente: la clave del mensaje de alerta es ventana|producto|mercado, así que un reprocesamiento produce el mismo mensaje y el consumidor de 08-02 lo tratará como el mismo.
Caso (b): panel de reparto. Las posiciones de furgoneta-3 llegan por MQTT (02-01) y un puente las publica en el tópico reparto.posiciones con clave = id de repartidor: {"repartidor":"furgoneta-3","lat":41.9794,"lon":2.8214,"fecha_ms":...}. Dos cálculos: la distancia recorrida en los últimos 5 minutos, refrescada cada minuto (sliding 5/1 por repartidor: detecta repartidores parados o desviados), y repartidores sin señal (session window con gap de 3 minutos: cuando la sesión se cierra, el repartidor lleva 3 minutos sin enviar; cuando se abre una nueva, ha vuelto). Ambos escriben en reparto.panel, que 08-02 empujará a los navegadores de los operadores.
flowchart LR
K1[(Kafka<br/>pedidos.eventos<br/>6 particiones)] --> F1[filtrar<br/>stock.actualizado]
F1 --> WM1[asignar ts + marca de agua<br/>ts − 2 min]
WM1 --> KB1[[keyBy producto, mercado]]
KB1 --> V1[tumbling 5 min<br/>MIN stock, COUNT]
V1 --> H1[HAVING min < 10]
H1 --> K2[(Kafka<br/>inventario.alertas)]
K3[(Kafka<br/>reparto.posiciones<br/>clave = repartidor)] --> WM2[asignar ts + marca de agua<br/>ts − 30 s]
WM2 --> KB2[[keyBy repartidor]]
KB2 --> V2[sliding 5 min / 1 min<br/>distancia]
KB2 --> V3[session gap 3 min<br/>inicio, fin, n]
V2 --> K4[(Kafka<br/>reparto.panel)]
V3 --> K4
- Práctica: simulación en Python, PyFlink y Structured Streaming
10.1 simulaciones/ventana_tumbling.py: un motor de ventanas en 80 líneas
Antes de usar un motor, conviene ver el mecanismo desnudo. Este script procesa una lista de eventos stock.actualizado con dos tiempos cada uno, el de evento (fecha_ms) y el de llegada (llegada_ms), en orden de llegada, manteniendo ventanas tumbling de 5 minutos por producto y una marca de agua con 2 minutos de retraso tolerado. Los tiempos están en minutos desde las 10:00 para leerlos fácilmente.
# km0/simulaciones/ventana_tumbling.py
"""Ventanas tumbling de 5 min en tiempo de evento, con marca de agua y eventos tardíos, en Python puro."""
from collections import defaultdict
VENTANA = 5 # minutos
RETRASO_TOLERADO = 2 # marca de agua = máximo tiempo de evento visto − 2 min
UMBRAL = 10
# (tiempo de evento, tiempo de llegada, producto, stock_actual), en minutos desde las 10:00.
# Están en ORDEN DE LLEGADA, que es el orden en que los ve el operador.
EVENTOS = [
(0.5, 0.6, "queso-curado", 42),
(1.2, 1.3, "queso-curado", 31),
(2.0, 2.1, "tomate-rosa", 120),
(3.8, 4.0, "queso-curado", 12),
(6.1, 6.2, "queso-curado", 6),
(6.5, 6.6, "tomate-rosa", 118),
(4.9, 7.0, "queso-curado", 8), # ocurrió a las 10:04:54, llega a las 10:07 (réplica inv-vlc con retraso)
(7.3, 7.4, "queso-curado", 25), # reposición
(4.2, 9.5, "queso-curado", 9), # ocurrió a las 10:04:12, llega a las 10:09:30: TARDÍO
(10.2, 10.3, "tomate-rosa", 117),
(12.7, 12.8, "queso-curado", 22),
]
def inicio_ventana(t: float) -> int:
return int(t // VENTANA) * VENTANA
def procesar(eventos):
ventanas = defaultdict(lambda: {"min": float("inf"), "n": 0}) # (inicio, producto) -> estado
marca_de_agua = float("-inf")
tardios = 0
for t_evento, t_llegada, producto, stock in eventos:
ini = inicio_ventana(t_evento)
if ini + VENTANA <= marca_de_agua: # la ventana ya se emitió: evento tardío
tardios += 1
print(f" [{t_llegada:5.1f}] TARDÍO: {producto} t={t_evento} (ventana {ini}-{ini + VENTANA} cerrada, "
f"marca de agua {marca_de_agua})")
continue
estado = ventanas[(ini, producto)] # estado por clave y ventana
estado["min"] = min(estado["min"], stock); estado["n"] += 1
marca_de_agua = max(marca_de_agua, t_evento - RETRASO_TOLERADO)
print(f" [{t_llegada:5.1f}] {producto:12s} t={t_evento:4.1f} stock={stock:3d} marca de agua={marca_de_agua:4.1f}")
# emitir toda ventana cuyo fin haya quedado por debajo de la marca de agua
for (v_ini, v_prod) in sorted(k for k in ventanas if k[0] + VENTANA <= marca_de_agua):
e = ventanas.pop((v_ini, v_prod))
alerta = " <-- ALERTA stock bajo" if e["min"] < UMBRAL else ""
print(f" EMITIR ventana {v_ini:2d}-{v_ini + VENTANA:2d} {v_prod:12s} min={e['min']:3d} n={e['n']}{alerta}")
print(f"Fin de la entrada: {len(ventanas)} ventanas abiertas sin emitir, {tardios} eventos tardíos")
if __name__ == "__main__":
procesar(EVENTOS)$ python ventana_tumbling.py
[ 0.6] queso-curado t= 0.5 stock= 42 marca de agua=-1.5
[ 1.3] queso-curado t= 1.2 stock= 31 marca de agua=-0.8
[ 2.1] tomate-rosa t= 2.0 stock=120 marca de agua= 0.0
[ 4.0] queso-curado t= 3.8 stock= 12 marca de agua= 1.8
[ 6.2] queso-curado t= 6.1 stock= 6 marca de agua= 4.1
[ 6.6] tomate-rosa t= 6.5 stock=118 marca de agua= 4.5
[ 7.0] queso-curado t= 4.9 stock= 8 marca de agua= 4.5
[ 7.4] queso-curado t= 7.3 stock= 25 marca de agua= 5.3
EMITIR ventana 0- 5 queso-curado min= 8 n=4 <-- ALERTA stock bajo
EMITIR ventana 0- 5 tomate-rosa min=120 n=1
[ 9.5] TARDÍO: queso-curado t=4.2 (ventana 0-5 cerrada, marca de agua 5.3)
[ 10.3] tomate-rosa t=10.2 stock=117 marca de agua= 8.2
[ 12.8] queso-curado t=12.7 stock= 22 marca de agua=10.7
EMITIR ventana 5-10 queso-curado min= 6 n=2 <-- ALERTA stock bajo
EMITIR ventana 5-10 tomate-rosa min=118 n=1
Fin de la entrada: 2 ventanas abiertas sin emitir, 1 eventos tardíosLo que muestra la traza:
- El evento de
t=4.9que llega a las 10:07 sí entra en la ventana 0–5, porque la marca de agua en ese momento era 4,5 (< 5): la ventana seguía abierta gracias al retraso tolerado. Con tiempo de procesamiento habría caído en 5–10, y el mínimo de la ventana 0–5 habría sido 12, sin alerta. - La ventana 0–5 se emite cuando llega el evento de
t=7.3: la marca de agua pasa a 5,3 ≥ 5. Ni antes (no se sabía si faltaban eventos) ni después (no hace falta esperar más). - El evento de
t=4.2que llega a las 10:09:30 es tardío: su ventana ya se emitió. Aquí se descarta y se cuenta; con allowed lateness se reemitiría la ventana 0–5 conmin=8, n=5(el mínimo no cambia, pero el conteo sí). - Al acabar la entrada quedan dos ventanas abiertas (10–15): en un flujo real no hay "fin", y se emitirán cuando la marca de agua llegue a 15. En un lote, el fin de la entrada dispara la emisión de todo.
Observa también que la marca de agua no avanza con el evento tardío ni retrocede nunca: es monótona. Y que cada línea de EMITIR es una salida que, si el proceso muriera y reprocesara desde el evento 1, se produciría otra vez con los mismos valores: la clave (ventana, producto) la hace idempotente.
10.2 El caso (a) en PyFlink (Table API / SQL)
El clúster Flink se añade al docker-compose.yml con dos servicios (jobmanager y taskmanager de la imagen flink:1.19-python, o una imagen propia con pip install apache-flink), y el trabajo se envía con flink run -py. Con la Table API, el DAG del caso (a) es tres sentencias SQL:
# km0/servicios/analitica/flujos/alertas_stock.py
"""Alerta de stock bajo: tumbling 5 min por (producto, mercado) en tiempo de evento, desde pedidos.eventos."""
from pyflink.table import EnvironmentSettings, TableEnvironment
t_env = TableEnvironment.create(EnvironmentSettings.in_streaming_mode())
t_env.get_config().set("pipeline.name", "km0-alertas-stock")
t_env.get_config().set("execution.checkpointing.interval", "30 s") # checkpoints cada 30 s
t_env.get_config().set("table.exec.source.idle-timeout", "1 min") # particiones sin tráfico no frenan la marca de agua
# Fuente: el tópico de eventos con la envoltura de 02-05. 'datos' es una fila anidada.
t_env.execute_sql("""
CREATE TABLE eventos (
id_evento STRING,
tipo STRING,
version INT,
fecha_ms BIGINT,
origen STRING,
datos ROW<producto STRING, mercado STRING, stock_actual INT, delta INT, replica STRING>,
ts AS TO_TIMESTAMP_LTZ(fecha_ms, 3), -- tiempo de EVENTO, derivado de fecha_ms
WATERMARK FOR ts AS ts - INTERVAL '2' MINUTE -- marca de agua: 2 min de retraso tolerado
) WITH (
'connector' = 'kafka',
'topic' = 'pedidos.eventos',
'properties.bootstrap.servers' = 'kafka:9092',
'properties.group.id' = 'analitica-alertas-stock',
'scan.startup.mode' = 'group-offsets', -- reanuda donde lo dejó (checkpoint)
'format' = 'json',
'json.ignore-parse-errors' = 'true' -- un evento corrupto no tumba el job
)""")
# Sink: alertas con clave (ventana, producto, mercado). El upsert-kafka escribe por clave: idempotente.
t_env.execute_sql("""
CREATE TABLE alertas_stock (
ventana_inicio TIMESTAMP_LTZ(3),
ventana_fin TIMESTAMP_LTZ(3),
producto STRING,
mercado STRING,
stock_min INT,
actualizaciones BIGINT,
PRIMARY KEY (ventana_inicio, producto, mercado) NOT ENFORCED
) WITH (
'connector' = 'upsert-kafka',
'topic' = 'inventario.alertas',
'properties.bootstrap.servers' = 'kafka:9092',
'key.format' = 'json',
'value.format' = 'json'
)""")
# El cálculo: ventana tumbling de 5 minutos sobre ts, por producto y mercado, solo stock.actualizado.
t_env.execute_sql("""
INSERT INTO alertas_stock
SELECT window_start, window_end,
datos.producto, datos.mercado,
MIN(datos.stock_actual) AS stock_min,
COUNT(*) AS actualizaciones
FROM TABLE(TUMBLE(TABLE eventos, DESCRIPTOR(ts), INTERVAL '5' MINUTE))
WHERE tipo = 'stock.actualizado'
GROUP BY window_start, window_end, datos.producto, datos.mercado
HAVING MIN(datos.stock_actual) < 10
""").wait()Cada pieza corresponde a un apartado: ts AS TO_TIMESTAMP_LTZ(fecha_ms, 3) declara el tiempo de evento; WATERMARK FOR ts AS ts - INTERVAL '2' MINUTE la marca de agua (Flink la genera por partición de Kafka y toma el mínimo, con el idle-timeout para particiones paradas); TUMBLE(..., INTERVAL '5' MINUTE) la ventana con window_start/window_end como columnas; el GROUP BY es el keyBy que reparte el estado por clave entre TaskManagers; execution.checkpointing.interval el checkpoint que guarda estado y offsets en el state.checkpoints.dir configurado (HDFS /km0/checkpoints/); y upsert-kafka con PRIMARY KEY el sink idempotente: un reprocesamiento tras un fallo reescribe la misma clave con el mismo valor. Se lanza y se observa así:
docker compose exec jobmanager flink run -py /app/analitica/flujos/alertas_stock.py -d
docker compose exec kafka kafka-console-consumer --bootstrap-server kafka:9092 \
--topic inventario.alertas --property print.key=true --from-beginning
{"ventana_inicio":"2026-09-14 10:00:00Z","producto":"queso-curado","mercado":"girona"} {"ventana_inicio":"2026-09-14 10:00:00Z","ventana_fin":"2026-09-14 10:05:00Z","producto":"queso-curado","mercado":"girona","stock_min":8,"actualizaciones":4}La interfaz web del JobManager (localhost:8081) muestra el DAG desplegado, el paralelismo de cada operador, la marca de agua actual de cada uno, los checkpoints (duración, tamaño) y el backpressure por operador, coloreado.
Para el caso (b), la misma estructura con SESSION para los repartidores sin señal:
INSERT INTO reparto_panel
SELECT repartidor,
SESSION_START(ts, INTERVAL '3' MINUTE) AS desde,
SESSION_END(ts, INTERVAL '3' MINUTE) AS hasta,
COUNT(*) AS posiciones
FROM posiciones -- tabla sobre reparto.posiciones, marca de agua ts - 30 s
GROUP BY repartidor, SESSION(ts, INTERVAL '3' MINUTE)Cada fila emitida significa "el repartidor furgoneta-3 envió posiciones de forma continua entre desde y hasta, y después estuvo al menos 3 minutos sin señal": es el aviso del panel. La ventana deslizante de distancia se escribe con HOP(TABLE posiciones, DESCRIPTOR(ts), INTERVAL '1' MINUTE, INTERVAL '5' MINUTE) y una función de agregación propia (la distancia entre posiciones consecutivas necesita el orden, que en SQL se resuelve con LAG sobre una ventana OVER antes de agregar).
10.3 El caso (a) en Spark Structured Streaming
El mismo cálculo con la API de DataFrames de 05-03, en modo streaming:
# km0/servicios/analitica/flujos/alertas_stock_spark.py
"""Alerta de stock bajo con Spark Structured Streaming: withWatermark + window tumbling de 5 min."""
from pyspark.sql import SparkSession, functions as F, types as T
ESQUEMA = T.StructType([
T.StructField("id_evento", T.StringType()), T.StructField("tipo", T.StringType()),
T.StructField("version", T.IntegerType()), T.StructField("fecha_ms", T.LongType()),
T.StructField("origen", T.StringType()),
T.StructField("datos", T.StructType([
T.StructField("producto", T.StringType()), T.StructField("mercado", T.StringType()),
T.StructField("stock_actual", T.IntegerType()), T.StructField("delta", T.IntegerType()),
T.StructField("replica", T.StringType())])),
])
spark = SparkSession.builder.appName("km0-alertas-stock").getOrCreate()
crudo = (spark.readStream.format("kafka") # fuente NO acotada
.option("kafka.bootstrap.servers", "kafka:9092")
.option("subscribe", "pedidos.eventos")
.option("startingOffsets", "latest")
.option("maxOffsetsPerTrigger", 50000) # backpressure: tope por microlote
.load())
stock = (crudo.select(F.from_json(F.col("value").cast("string"), ESQUEMA).alias("e")).select("e.*")
.filter(F.col("tipo") == "stock.actualizado")
.withColumn("ts", (F.col("fecha_ms") / 1000).cast("timestamp"))) # tiempo de evento
alertas = (stock
.withWatermark("ts", "2 minutes") # marca de agua: 2 min
.groupBy(F.window("ts", "5 minutes").alias("ventana"), # tumbling de 5 min
F.col("datos.producto").alias("producto"), F.col("datos.mercado").alias("mercado"))
.agg(F.min("datos.stock_actual").alias("stock_min"), F.count("*").alias("actualizaciones"))
.filter(F.col("stock_min") < 10)
.select(F.concat_ws("|", F.col("ventana.start"), "producto", "mercado").alias("key"), # clave: idempotente
F.to_json(F.struct(F.col("ventana.start").alias("ventana_inicio"), F.col("ventana.end").alias("ventana_fin"),
"producto", "mercado", "stock_min", "actualizaciones")).alias("value")))
consulta = (alertas.writeStream
.outputMode("append") # emite cada ventana UNA vez, cuando la marca de agua la cierra
.format("kafka")
.option("kafka.bootstrap.servers", "kafka:9092")
.option("topic", "inventario.alertas")
.option("checkpointLocation", "hdfs://namenode:8020/km0/checkpoints/alertas-stock") # offsets + estado
.trigger(processingTime="10 seconds") # un microlote cada 10 s
.start())
consulta.awaitTermination()La correspondencia con Flink es directa: withWatermark ↔ WATERMARK FOR, F.window("ts", "5 minutes") ↔ TUMBLE, checkpointLocation ↔ execution.checkpointing, maxOffsetsPerTrigger ↔ backpressure. Dos diferencias importantes. El outputMode("append") con marca de agua emite cada ventana una sola vez, al cerrarla, y descarta los tardíos sin posibilidad de reemitir (no hay allowed lateness); update emitiría resultados provisionales en cada microlote, que exigen un sink que sepa sobrescribir. Y el sink de Kafka de Spark es at-least-once: el mensaje puede duplicarse tras un fallo, y es la clave key (idéntica en el duplicado) la que hace que el consumidor de 08-02 lo trate como una repetición inocua, exactamente el consumidor idempotente de 02-05. Se lanza con spark-submit --packages org.apache.spark:spark-sql-kafka-0-10_2.12:3.5.1 alertas_stock_spark.py, y en localhost:4040 aparece una pestaña Structured Streaming con la tasa de entrada, la duración de cada microlote y la marca de agua.
Errores Comunes y Consejos
- Usar tiempo de procesamiento porque es más fácil. Los resultados dependen de la carga, los reinicios y el backpressure, y no se pueden reproducir. Si el evento tiene marca de tiempo (y con la envoltura de 02-05 siempre la tiene), usa tiempo de evento.
- Marca de agua sin medir. Un retraso tolerado inventado descarta eventos válidos o retrasa los resultados. Mide el desfase
llegada − fecha_msen producción (percentil 99) y cuenta los tardíos como métrica permanente. - Olvidar las particiones inactivas. Una partición de Kafka sin tráfico retiene la marca de agua global y ninguna ventana se emite. Configura el idle timeout (Flink) o revisa que todas las particiones reciban eventos.
- Estado sin límite. Una agregación por clave sin ventana ni TTL crece para siempre (una clave por
pedido_id, por ejemplo). Ventanas, TTL de estado, o claves acotadas. - Sink no idempotente con reprocesamiento. Tras un fallo, las salidas posteriores al checkpoint se repiten. Clave natural + upsert, o sink transaccional. Nunca un
INSERTsin clave ni unappenda un fichero. - Checkpoint en disco local. Si el nodo muere, el checkpoint muere con él. HDFS, MinIO o un almacenamiento replicado, siempre.
- Cambiar el DAG y reanudar del checkpoint. Un checkpoint guarda estado por operador; si el DAG cambia (nuevo operador, otra clave) puede no ser compatible. Savepoints con identificadores de operador (
uid) en Flink; en Spark, cambios de esquema de estado suelen exigir empezar de cero. - Confundir lag con latencia. El lag del grupo de consumidores (eventos pendientes en Kafka) crece bajo backpressure y es la señal de que falta paralelismo; la latencia de un resultado en tiempo de evento es siempre al menos el retraso tolerado más el tamaño de la ventana.
- Reescribir la lógica dos veces (Lambda por inercia). Si el mismo agregado se calcula en lote y en flujo, saldrá distinto y nadie sabrá cuál es el bueno. Un código, dos despliegues.
Ejercicios
Ejercicio 1: Elegir marca de agua y ventana
Un análisis del tópico reparto.posiciones durante una semana muestra que el desfase entre fecha_ms y la llegada a Kafka tiene esta distribución: mediana 1,2 s; percentil 95, 8 s; percentil 99, 45 s; percentil 99,9, 4 min (túneles, zonas sin cobertura); máximo 22 min (un móvil apagado que reenvió al encender). El panel debe mostrar la posición y la distancia de los últimos 5 minutos con no más de 1 minuto de retraso. Elige el retraso tolerado de la marca de agua, el tipo y tamaño de ventana, y qué hacer con los eventos tardíos. Justifica los porcentajes de eventos que se descartarán y qué pasa con el móvil que reenvía tras 22 minutos.
Ejercicio 2: Trazar la simulación con otra marca de agua
Ejecuta mentalmente (o modificando el script) ventana_tumbling.py con RETRASO_TOLERADO = 0 y con RETRASO_TOLERADO = 5. Para cada caso, indica cuándo se emite la ventana 0–5 de queso-curado, con qué mínimo y conteo, y cuántos eventos tardíos hay. ¿Cuál de las tres configuraciones (0, 2, 5) daría la alerta correcta más pronto?
Ejercicio 3: Fallo y reprocesamiento
El job de PyFlink de alertas lleva checkpoints cada 30 s. A las 10:07:50 el TaskManager que ejecuta la ventana de queso-curado/girona muere; el último checkpoint completado es de las 10:07:30, y a las 10:07:40 el job había emitido la alerta de la ventana 10:00–10:05. Describe qué hace Flink al recuperarse: desde qué offsets lee, qué pasa con el estado de la ventana 10:05–10:10, si la alerta 10:00–10:05 se vuelve a emitir y qué ve el consumidor de inventario.alertas. Después, explica qué cambiaría si el sink fuera un INSERT en PostgreSQL sin clave primaria, y cómo lo arreglarías.
Soluciones
Ejercicio 1.
El presupuesto de latencia es 1 minuto, y la latencia mínima de un resultado es el retraso tolerado (más el intervalo de emisión). Un retraso tolerado de 45 s (el percentil 99) cumple el presupuesto y descarta como tardíos el 1 % de las posiciones; con 8 s (p95) se descartaría el 5 %, demasiado para una traza de posiciones; con 4 min (p99,9) se violaría el requisito de 1 minuto. Ventana deslizante de 5 minutos con avance de 1 minuto por repartidor (o de 30 s si el panel debe refrescar más a menudo, a costa de más ventanas abiertas por clave: 10 en vez de 5). Los tardíos (1 %) van a una salida lateral que se cuenta y se escribe en el lago: para el panel en tiempo real da igual perder una posición de cada cien, pero la distancia recorrida del día que calcula el lote nocturno debe incluirlas, y para eso el lote lee todas las posiciones del lago, tardías incluidas (una Lambda de facto justificada). El móvil que reenvía tras 22 minutos entrega posiciones cuya ventana cerró hace 20 minutos: todas tardías, todas a la salida lateral; el panel las ignora y el lote las incorpora. Alternativa si el negocio lo pidiera: allowed lateness de 5 minutos para reemitir ventanas recientes, no de 22 (mantendría 27 minutos de ventanas en estado por repartidor).
Ejercicio 2.
Con RETRASO_TOLERADO = 0 la marca de agua es el máximo tiempo de evento visto. Al llegar t=6.1 (llegada 6,2) la marca de agua es 6,1 ≥ 5 y la ventana 0–5 se emite con los eventos vistos hasta entonces: 0,5, 1,2, 3,8 y... el de t=4.9 llegó a las 7,0, después, así que la ventana se emite con min=12, n=3: sin alerta, incorrecta. Los eventos de t=4.9 (llegada 7,0) y t=4.2 (llegada 9,5) son ambos tardíos: 2 tardíos. Emisión más temprana, resultado equivocado.
Con RETRASO_TOLERADO = 5, la ventana 0–5 se emite cuando la marca de agua alcanza 5, es decir, al ver un evento con t ≥ 10: el de t=10.2 (llegada 10,3). Para entonces han entrado 0,5, 1,2, 3,8, 4,9 y 4,2 (que llegó a las 9,5 con la ventana aún abierta): min=8, n=5, correcto y completo, 0 tardíos. Pero la alerta sale a las 10:10:18, cinco minutos más tarde que con retraso 2 (10:07:24, min=8, n=4).
La configuración con 2 minutos da la alerta correcta (el mínimo 8 estaba en ambas) más pronto; la de 5 minutos da además el conteo exacto; la de 0 falla. Es el compromiso completitud/latencia del apartado 2, y la razón de que el retraso tolerado se mida y no se adivine.
Ejercicio 3.
Flink detecta la muerte del TaskManager (heartbeat), reinicia el job entero (o la región afectada, con fine-grained recovery) desde el checkpoint de las 10:07:30: restaura el estado de todos los operadores tal como estaba entonces (la ventana 10:05–10:10 con los eventos anteriores a las 10:07:30 ya aplicados; la 10:00–10:05, que aún no se había emitido en ese instante, también restaurada con su estado) y reposiciona el consumidor de Kafka en los offsets guardados en ese checkpoint. Los eventos entre las 10:07:30 y las 10:07:50 se vuelven a leer y a aplicar sobre ese estado: ninguno se cuenta dos veces, porque el estado que los contenía se descartó. La marca de agua vuelve a avanzar y la ventana 10:00–10:05 se emite de nuevo, con el mismo contenido. Como el sink es upsert-kafka con clave (ventana_inicio, producto, mercado), el tópico recibe un segundo mensaje con la misma clave y el mismo valor; el consumidor de 08-02 (o la compactación de Kafka) lo trata como una actualización sin cambios. Exactly-once en el estado, y efectivamente una sola alerta visible.
Con un INSERT en PostgreSQL sin clave, la fila de la alerta 10:00–10:05 existiría dos veces. Arreglos, de menor a mayor esfuerzo: clave primaria (ventana_inicio, producto, mercado) con INSERT ... ON CONFLICT DO UPDATE (sink idempotente, el JDBC sink de Flink lo hace con upsert); o el sink JDBC en modo exactly-once con XA (dos fases, confirma con el checkpoint), que añade la latencia del checkpoint a cada alerta. Para alertas, la primera opción es la correcta.
Conclusión
Procesar un flujo es ejecutar de forma permanente el mismo DAG de operadores que un lote, sobre una entrada que no termina, y eso obliga a responder preguntas que el lote no tenía: cuándo está completa una ventana (la marca de agua, derivada del tiempo de evento y de un retraso tolerado que se mide), qué hacer con lo que llega después (descartar, reemitir o desviar los tardíos), cómo agrupar (ventanas tumbling, sliding y session por clave), cómo hacer durable un estado que vive para siempre (checkpoints consistentes con los offsets, savepoints para desplegar) y cómo no contar dos veces tras un fallo (exactly-once en el estado por checkpoint, y en el sink por idempotencia o transacción). Backpressure protege al motor dejando que Kafka absorba los picos, y el lag es la señal. Lambda y Kappa son dos formas de convivir con los lotes: la primera duplica la lógica, la segunda la unifica, y los motores modernos hacen que la elección sea de despliegue. En Kilómetro Cero, simulaciones/ventana_tumbling.py hizo visible el mecanismo, PyFlink resolvió la alerta de stock bajo con tres sentencias SQL y un sink upsert-kafka, la session window detecta repartidores sin señal, y Spark Structured Streaming demostró que el código de 05-03 sirve casi sin cambios con withWatermark y window.
Con esta lección, analitica tiene sus dos mitades: los lotes de ventas_diarias.py y ALS, y los flujos de alertas y reparto. Pero los lotes no se lanzan solos. Alguien tiene que esperar a que el fichero del día esté completo en HDFS, validarlo, lanzar spark-submit, cargar el resultado en la base de datos de analitica, avisar si algo falla y reintentar, y hacerlo cada día, y para los siete días de la Semana de la Vendimia cuando se corrigió el precio del vino. Hasta ahora eso era un cron y una cadena de scripts. La última lección del módulo trata la planificación de trabajos y los pipelines de datos: cómo expresar esas dependencias como un DAG de tareas en Airflow, con sensores, reintentos, backfill y alertas, para que la plataforma de datos de Kilómetro Cero funcione sin que nadie la lance a mano.
Curso de Arquitecturas Distribuidas
Módulo 1: Introducción a los Sistemas Distribuidos
- Conceptos Básicos de Sistemas Distribuidos
- Modelos de Sistemas Distribuidos
- Ventajas y Desafíos de los Sistemas Distribuidos
- Las Falacias de la Computación Distribuida
- Tiempo, Relojes y Ordenación de Eventos
- Del Monolito a la Plataforma Distribuida: el Caso Kilómetro Cero
Módulo 2: Comunicación en Sistemas Distribuidos
- Protocolos de Comunicación
- RPC y RMI
- gRPC y Serialización de Datos
- Mensajería y Colas de Mensajes
- Patrones de Comunicación Asíncrona
Módulo 3: Consistencia y Replicación
- Modelos de Consistencia
- El Teorema CAP y PACELC
- Algoritmos de Consenso
- Replicación de Datos
- Transacciones Distribuidas y Sagas
Módulo 4: Almacenamiento Distribuido
- Particionado de Datos y Hashing Consistente
- Sistemas de Archivos Distribuidos
- Almacenamiento de Objetos
- Bases de Datos Distribuidas
- Cachés Distribuidos
Módulo 5: Computación Distribuida
- Modelos de Computación Distribuida
- MapReduce y Hadoop
- Spark y Computación en Memoria
- Procesamiento de Flujos de Datos
- Planificación de Trabajos y Pipelines de Datos
Módulo 6: Seguridad en Sistemas Distribuidos
- Autenticación y Autorización
- Cifrado y Protección de Datos
- Gestión de Identidades
- Seguridad entre Servicios: mTLS y Gestión de Secretos
- Puertas de Enlace, Limitación de Tasa y Auditoría
Módulo 7: Monitoreo y Mantenimiento
- Monitoreo de Sistemas Distribuidos
- Logs Centralizados y Trazabilidad Distribuida
- Gestión de Fallos y Recuperación
- Patrones de Resiliencia: Timeouts, Reintentos y Circuit Breaker
- Automatización y Orquestación
- Pruebas en Sistemas Distribuidos e Ingeniería del Caos
