La lección anterior dejó el mapa de contextos de Kilómetro Cero con una flecha que terminaba en Kafka: reparto publica en reparto.posiciones, Flink calcula sobre ese tópico y escribe en reparto.panel (05-04), y ahí se acababa el recorrido. Pero Ana no lee Kafka. Ana tiene abierta la pantalla de seguimiento de su pedido en el móvil, en el autobús, con cobertura irregular, y quiere ver moverse a furgoneta-3 por el mapa sin pulsar "actualizar". Jordi Sala tiene el panel de operadores en un navegador de escritorio con 140 repartidores a la vez. La app de la furgoneta envía su posición cada 5 segundos desde una red móvil que se cae en cada túnel. Y la Quesería Montblanc quiere un aviso en su panel en cuanto queso-curado baje del umbral en Girona. Todos son problemas del último tramo: entre la plataforma (servicios, Kafka, bases de datos, todo en un centro de datos con red fiable) y las personas y dispositivos que están fuera, en redes malas, con miles o decenas de miles de conexiones simultáneas, y sin poder ejecutar un consumidor de Kafka. Esta lección desarrolla lo que 02-01 dejó presentado: MQTT para los dispositivos y WebSockets para los navegadores, junto con las alternativas (polling, long polling, Server-Sent Events), y las une en una arquitectura pub/sub de extremo a extremo con sus problemas propios: el fan-out entre instancias, la autenticación de conexiones largas, el orden y los duplicados en el cliente, la presión sobre clientes lentos y el escalado de conexiones. La nube (08-03) y el procesamiento en el borde (08-04) quedan fuera.

Contenido

  1. Qué significa "tiempo real" aquí y los cuatro casos de Kilómetro Cero
  2. El último tramo frente a la mensajería entre servicios
  3. Técnicas del último tramo: polling, long polling, SSE y WebSockets
  4. MQTT para dispositivos
  5. Arquitectura pub/sub de extremo a extremo
  6. Fan-out con muchas instancias, presencia y suscripciones
  7. Autenticación, orden, deduplicación, backpressure y escalado
  8. Código: del broker MQTT al navegador
  9. SSE como alternativa para las alertas de stock
  10. Cuándo usar cada técnica
  11. Errores Comunes y Consejos
  12. Ejercicios
  13. Conclusión

  1. Qué significa "tiempo real" aquí y los cuatro casos de Kilómetro Cero

"Tiempo real" tiene un significado estricto en ingeniería de control: un sistema de tiempo real duro es el que garantiza que una respuesta llega antes de un plazo, y si no llega, el resultado es un fallo (el airbag, el control de un motor). Nada de eso aplica aquí. En sistemas web, "tiempo real" significa latencia percibida por personas: que la información llegue lo bastante rápido para que quien la mira sienta que ve lo que está pasando, sin que tenga que pedirla. Ese umbral está entre décimas de segundo y unos pocos segundos, y varía por caso. Conviene fijarlo por caso, porque el diseño cambia con él.

Caso Quién recibe Origen del dato Latencia aceptable Volumen Sentido
Ana ve moverse a furgoneta-3 en el mapa Navegador o app móvil de un cliente, en red móvil Posición de la furgoneta (MQTT → Kafka) 2-5 s (una posición cada 5 s) Miles de clientes con seguimiento abierto en campaña; cada uno interesado en una furgoneta Servidor → cliente
Panel de operadores de reparto Navegador de escritorio de Jordi y Marta reparto.panel (Flink, 05-04) y posiciones 1-2 s Decenas de operadores; cada uno interesado en todas las furgonetas de su mercado Servidor → cliente
Avisos de stock a productores Panel web del productor inventario.alertas (Flink alertas_stock.py) 10-30 s Cientos de productores; pocos mensajes al día por productor Servidor → cliente
Chat cliente-productor Navegador de Ana y panel de la quesería Mensajes escritos por personas < 1 s Conversaciones cortas; bidireccional Ambos sentidos

Y un quinto, del lado de origen: la app de la furgoneta publica su posición cada 5 segundos, desde una red móvil que pierde cobertura, con una batería que cuidar, y debe seguir funcionando cuando vuelve la cobertura sin perder lo importante. Tiene requisitos opuestos a los de Ana: pocos bytes, una sola conexión persistente, tolerancia a desconexiones largas.

  1. El último tramo frente a la mensajería entre servicios

La mensajería del Módulo 2 (RabbitMQ, Kafka) resuelve la comunicación entre servicios: procesos de larga vida, en la misma red, con clientes de biblioteca que mantienen conexiones, offsets y grupos de consumidores. El último tramo es distinto en todo:

Aspecto Entre servicios (02-04) Último tramo
Cliente Un proceso Python con confluent_kafka, en Kubernetes Un navegador (solo HTTP y WebSocket), una app móvil, un dispositivo con poca batería
Red Centro de datos: fiable, baja latencia Internet, 4G, Wi-Fi de bar: pérdidas, NAT, proxies, túneles
Número de conexiones Decenas Miles a millones
Vida de la conexión Horas o días Minutos; se corta al cambiar de red o al bloquear el móvil
Estado del consumidor Offset persistido en el broker; reanuda donde lo dejó Ninguno o mínimo; al reconectar, hay que decidir qué se perdió
Seguridad mTLS entre servicios (06-04) Un usuario con un JWT (06-01); el cliente no es de confianza
Fan-out Un grupo de consumidores por servicio Cada persona quiere un subconjunto distinto de los mensajes

Por eso nunca se expone Kafka directamente a un navegador: ni el protocolo lo permite, ni el modelo de seguridad, ni el número de conexiones. Entre Kafka y el navegador hace falta un componente que hable el protocolo del cliente, que mantenga sus conexiones y que reparta a cada uno solo lo que le corresponde. Ese componente es el servidor de tiempo real que construiremos como parte de reparto.

  1. Técnicas del último tramo: polling, long polling, SSE y WebSockets

Todas parten de una limitación: HTTP es petición-respuesta y el servidor no puede hablar primero. Las cuatro técnicas son maneras de darle la vuelta a esa limitación, cada una con un coste.

Polling

El cliente pregunta cada n segundos: GET /api/v1/reparto/pedidos/P-2026-000124/posicion. Es trivial, cacheable en el gateway, compatible con todo, y es exactamente lo que el monolito hacía en el síntoma 4 de 01-06 con el resultado conocido: con 5 000 clientes preguntando cada 3 segundos son 1 700 peticiones por segundo, la mayoría para recibir "nada nuevo", cada una con su handshake TLS si no se reutiliza la conexión (02-01) y sus 300 bytes de cabeceras. La latencia media es la mitad del intervalo.

Long polling

El cliente pregunta, pero el servidor no responde hasta que haya algo nuevo (o hasta un timeout de, por ejemplo, 30 s), y el cliente vuelve a preguntar inmediatamente. Reduce las peticiones vacías y la latencia baja a casi cero, pero mantiene una petición HTTP abierta por cliente (un worker bloqueado en servidores síncronos), sufre con proxies que cortan conexiones inactivas, y cada mensaje sigue costando una petición completa. Fue la técnica dominante antes de WebSockets, y sigue siendo válida como respaldo.

Server-Sent Events (SSE)

El cliente hace un GET con Accept: text/event-stream y el servidor deja la respuesta abierta indefinidamente, escribiendo eventos en un formato de texto sencillo (event:, data:, id:, separados por línea en blanco). Es HTTP normal (atraviesa proxies y gateways, va sobre HTTP/2 multiplexado), el navegador lo implementa con EventSource, que reconecta sola y envía la cabecera Last-Event-ID con el último id recibido para que el servidor reanude. Su límite: es unidireccional (servidor → cliente; el cliente, si quiere enviar, usa peticiones HTTP normales) y solo texto.

WebSockets

Lo presentado en 02-01: una petición HTTP con Upgrade: websocket que, si el servidor acepta (101 Switching Protocols), convierte la conexión TCP en un canal bidireccional, persistente y con tramas (binarias o de texto), sin cabeceras HTTP por mensaje (2 a 14 bytes de sobrecarga por trama). Conviene entender las cuatro piezas del protocolo que afectan al diseño:

  • Handshake: el cliente envía Sec-WebSocket-Key; el servidor responde con Sec-WebSocket-Accept derivado de ella. Es el único momento en que hay cabeceras HTTP, y por tanto el único en que Kong puede aplicar sus plugins (06-05): el JWT se verifica en el handshake, no por mensaje.
  • Tramas: cada mensaje viaja en una o más tramas con un código de operación (texto, binario, cierre, ping, pong). Las tramas del cliente van enmascaradas (XOR con una clave aleatoria) por una razón de seguridad frente a proxies antiguos, no de cifrado; el cifrado es TLS (wss://).
  • Ping/pong: tramas de control que cualquiera de los dos lados envía para comprobar que el otro sigue ahí. Sin ellas, un móvil que pierde cobertura deja una conexión "abierta" en el servidor durante los minutos que TCP tarda en darse cuenta (02-01). El servidor de Kilómetro Cero hace ping cada 20 s y cierra si no hay pong en 10 s.
  • Reconexión: el protocolo no la define. Cuando la conexión cae, es el cliente quien reconecta, con backoff exponencial y jitter (los mismos de 07-04, por la misma razón: 5 000 clientes reconectando a la vez tras un despliegue es una estampida), y quien tiene que decir al servidor "qué es lo último que vi" para recuperar lo perdido. Todo eso es código de aplicación.
Técnica Sentido Latencia Coste por mensaje Conexiones abiertas Atraviesa proxies Reconexión Uso en Kilómetro Cero
Polling Cliente pide Media = intervalo / 2 Una petición HTTP completa (con o sin dato) No (o una reutilizada) Sí Trivial (sin estado) Respaldo si WebSocket falla; datos que cambian cada minutos
Long polling Cliente pide, servidor retiene Casi inmediata Una petición por mensaje Una por cliente, en espera Con timeouts cortos Trivial Respaldo
SSE Servidor → cliente Inmediata Unas decenas de bytes Una por cliente Sí (HTTP/2 mejor) Automática (Last-Event-ID) Alertas de stock a productores
WebSockets Bidireccional Inmediata 2-14 bytes de trama Una por cliente Sí, con Upgrade permitido en gateway Manual (código de cliente) Mapa de Ana, panel de operadores, chat
MQTT (sobre TCP o WebSockets) Pub/sub bidireccional Inmediata 2 bytes de cabecera fija Una por dispositivo Sobre WebSockets, sí Del cliente, con sesión persistente en el broker App de la furgoneta

  1. MQTT para dispositivos

02-01 presentó MQTT como protocolo pub/sub para dispositivos con recursos limitados y redes malas. Aquí lo diseñamos para la flota de Kilómetro Cero, apoyándonos en las siete características que lo hacen adecuado:

  1. Broker central. Los dispositivos no se conocen entre sí ni conocen a los consumidores: publican en el broker (Mosquitto para empezar; EMQX o HiveMQ para cientos de miles de conexiones) y el broker reparte. La furgoneta no sabe que existe Kafka.
  2. Tópicos jerárquicos. km0/reparto/furgoneta-3/posicion, km0/reparto/furgoneta-3/estado, km0/reparto/furgoneta-3/comandos. Los comodines permiten suscribirse a niveles: + un nivel (km0/reparto/+/posicion: todas las posiciones), # todo lo que cuelga (km0/reparto/furgoneta-3/#). La jerarquía es además la base de las ACL: furgoneta-3 solo puede publicar bajo su prefijo.
  3. QoS por mensaje. QoS 0 (at most once): se envía y se olvida; adecuado para la posición cada 5 s, porque la siguiente la reemplaza. QoS 1 (at least once): el broker confirma con PUBACK y el cliente reenvía si no llega; puede duplicar. QoS 2 (exactly once): intercambio de cuatro mensajes; caro y raramente necesario. Kilómetro Cero publica posiciones con QoS 1 y no con 0 por una razón que se ve en el apartado 7: al reconectar tras un túnel, quiere que la última posición conocida llegue seguro, aunque se duplique (el consumidor deduplica por secuencia).
  4. Retained messages. Al publicar con retain=True, el broker guarda el último mensaje del tópico y lo entrega inmediatamente a quien se suscriba después. Es lo que permite que un panel recién abierto vea la última posición de cada furgoneta sin esperar 5 s.
  5. Last will (testamento). Al conectar, el cliente registra un mensaje que el broker publicará por él si la conexión se pierde sin un DISCONNECT limpio: km0/reparto/furgoneta-3/estado = desconectado. Es la detección de fallo de 07-03 hecha por el broker.
  6. Sesiones persistentes. Con clean_session=False (MQTT 3.1.1) o session expiry (MQTT 5), el broker recuerda las suscripciones del cliente y encola los mensajes QoS 1/2 que le lleguen mientras está desconectado, y se los entrega al volver. Para los comandos que el operador envía a la furgoneta ("nueva parada añadida") es imprescindible: el túnel no puede perderlos.
  7. MQTT sobre WebSockets. Un navegador no abre sockets TCP arbitrarios, pero sí WebSockets; los brokers exponen un puerto WebSocket (9001 en Mosquitto) por el que hablan MQTT dentro de tramas WebSocket. Es la vía por la que el panel de operadores podría suscribirse directamente al broker; en el diseño de Kilómetro Cero no se hace (apartado 5), pero es útil para herramientas de diagnóstico.
Decisión Posición de la furgoneta Estado (conectada/desconectada) Comandos del operador a la furgoneta
Tópico km0/reparto/<id>/posicion km0/reparto/<id>/estado km0/reparto/<id>/comandos
QoS 1 1 1
Retained Sí (última posición) Sí No
Sesión persistente No hace falta (solo publica) — Sí (el dispositivo debe recibirlos al volver)
Last will — desconectado —

  1. Arquitectura pub/sub de extremo a extremo

Con las piezas anteriores, el recorrido de una posición desde la furgoneta hasta el mapa de Ana es este:

flowchart LR
    subgraph Calle[Furgoneta, red móvil]
        APP[App repartidor<br/>furgoneta_mqtt.py<br/>QoS 1, retained]
    end
    APP -- "MQTT/TLS 8883<br/>km0/reparto/furgoneta-3/posicion" --> BR[Broker MQTT<br/>Mosquitto / EMQX<br/>ACL por repartidor]
    BR -- "suscripción km0/reparto/+/posicion" --> PU[puente_mqtt_kafka.py]
    PU -- "clave = repartidor" --> K[(Kafka<br/>reparto.posiciones)]
    K --> FL[Flink alertas y panel<br/>05-04]
    FL --> KP[(Kafka<br/>reparto.panel)]
    K --> CS[(Cassandra<br/>posiciones, 04-04)]
    K --> WS1[reparto ws_servidor<br/>instancia 1]
    KP --> WS1
    K --> WS2[reparto ws_servidor<br/>instancia 2]
    KP --> WS2
    WS1 <--> RP[(Redis pub/sub<br/>bus entre instancias)]
    WS2 <--> RP
    WS1 -- "wss:// vía Kong" --> ANA[Navegador de Ana<br/>suscrita a P-2026-000124]
    WS2 -- "wss:// vía Kong" --> JOR[Panel de Jordi<br/>suscrito a mercado girona]

Cada salto tiene una razón:

  • La furgoneta habla MQTT, no HTTP ni Kafka. Una conexión persistente barata, con reconexión y QoS gestionados por la biblioteca, y un broker que la autentica y limita a su prefijo de tópicos.
  • El puente MQTT → Kafka existe porque el resto de la plataforma consume Kafka: Flink (05-04), Cassandra para el histórico, analitica. Traduce el tópico jerárquico a un tópico Kafka con clave = id de repartidor, lo que conserva el orden por furgoneta (02-04) y permite a Flink la ventana por clave. Algunos brokers (EMQX) traen este puente integrado; con Mosquitto se escribe uno pequeño.
  • reparto consume Kafka y hace fan-out por WebSockets. Es el único componente que conoce a los clientes: sabe que Ana está suscrita al pedido P-2026-000124, que ese pedido lo lleva furgoneta-3, y que por tanto cada posición de furgoneta-3 debe llegar a la conexión de Ana. Ese conocimiento (la tabla de suscripciones) es lo que ninguna otra pieza tiene.
  • Redis pub/sub entre instancias resuelve el problema del apartado siguiente.

Nótese lo que no se hace: el navegador no se conecta al broker MQTT (aunque podría, sobre WebSockets) porque eso exigiría dar a cada cliente credenciales MQTT, gestionar ACL por pedido en el broker y renunciar al enriquecimiento que hace reparto (Ana no debe recibir la posición cruda de furgoneta-3 con sus otros doce pedidos, sino "tu pedido está a 4 paradas y 12 minutos").

  1. Fan-out con muchas instancias, presencia y suscripciones

El problema

reparto corre con varias réplicas en Kubernetes (07-05). Ana está conectada a la instancia 1; Marc, que espera otra entrega de la misma furgoneta-3, a la instancia 2. Cuando llega la posición por Kafka, ¿quién la recibe? Si las dos instancias forman un grupo de consumidores de Kafka, cada partición la lee una de ellas (02-04): la posición llega a la instancia 1, que la envía a Ana, y Marc no ve nada. Si cada instancia consume todo el tópico con un grupo distinto, cada una recibe todas las posiciones y las reparte a sus clientes: funciona, pero cada instancia procesa las 28 posiciones por segundo completas y consume el tópico entero, lo que a decenas de instancias y con reparto.panel incluido empieza a pesar, y no resuelve mensajes que nacen en una instancia (el chat: Ana escribe en la instancia 1 y la quesería está en la 2).

Las dos soluciones

Sesiones pegajosas (sticky sessions): el balanceador (Kong, o el Ingress) envía siempre al mismo cliente a la misma instancia (por cookie o por hash de IP), y se enruta cada mensaje a la instancia correcta. Exige que alguien sepa en qué instancia está cada cliente (un mapa cliente → instancia en Redis) y que las instancias se hablen. Es frágil: al escalar o al caer una instancia el mapa se invalida, y el hash por IP falla con NAT de operadores móviles (miles de clientes tras la misma IP).

Bus entre instancias (la opción de Kilómetro Cero): cada instancia se suscribe a un canal de Redis pub/sub (o a un tópico Kafka con un grupo por instancia) por cada tema en el que tiene al menos un cliente interesado. Quien tenga algo que difundir (el consumidor de Kafka de cualquier instancia, o el manejador del chat) lo publica en el bus, y cada instancia con suscriptores lo recibe y lo reparte a sus conexiones. No hace falta saber dónde está nadie; las instancias son intercambiables y se escalan sin coordinación.

flowchart TB
    K[(Kafka reparto.posiciones)] --> C1[Consumidor Kafka<br/>grupo reparto-ws<br/>en la instancia que toque]
    C1 -- "PUBLISH furgoneta:furgoneta-3" --> R[(Redis pub/sub)]
    R -- "SUBSCRIBE furgoneta:furgoneta-3" --> I1[Instancia 1<br/>suscriptores locales:<br/>Ana]
    R -- "SUBSCRIBE furgoneta:furgoneta-3" --> I2[Instancia 2<br/>suscriptores locales:<br/>Marc, Jordi]
    R -. "nadie suscrito: no se suscribe" .- I3[Instancia 3<br/>sin interesados]
    I1 --> A[Ana]
    I2 --> M[Marc]
    I2 --> J[Jordi]

Redis pub/sub es fire-and-forget: si una instancia está desconectada del bus en ese instante, se pierde el mensaje, y no hay histórico. Para las posiciones da igual (la siguiente llega en 5 s y el cliente puede pedir la última conocida al conectar). Para el chat no: los mensajes se persisten primero (en Cassandra o PostgreSQL) y el bus solo notifica; al reconectar, el cliente pide "todo desde el id X" a la API. Cuando el volumen o la garantía exigen más, el bus pasa a ser Kafka con un grupo de consumidores por instancia, o Redis Streams.

Presencia y suscripciones

La tabla de suscripciones vive en memoria en cada instancia ({tema: {conexiones}}): es local, rápida y se pierde con la instancia, que es lo correcto porque las conexiones también se pierden. La presencia (quién está conectado ahora, útil para "Marta está en línea" en el chat o para que el operador sepa qué clientes están mirando) se guarda en Redis con TTL (presencia:u-ana → instancia-1, renovado con cada ping): si la instancia muere, la clave caduca sola. Y la resolución de a qué furgoneta corresponde el pedido de Ana se hace en el momento de suscribirse, consultando el modelo de reparto, y se vuelve a hacer si llega un evento de reasignación.

  1. Autenticación, orden, deduplicación, backpressure y escalado

Autenticación del WebSocket con el JWT

La conexión WebSocket dura minutos; el JWT de 06-01 caduca a los 15. Tres decisiones:

  • El JWT se verifica en el handshake (Kong lo hace con el plugin jwt como con cualquier ruta; reparto lo vuelve a verificar, defensa en profundidad). Se pasa como parámetro de consulta (wss://.../ws?token=...) o, mejor, como primer mensaje tras conectar, porque las URL acaban en logs (07-02) y el token no debe.
  • Autorización por suscripción: cada mensaje {"accion": "suscribir", "pedido": "P-2026-000124"} se autoriza contra los claims (sub = u-ana debe ser el cliente del pedido; un roles: ["operador"] puede suscribirse a un mercado). No hay autorización por cada mensaje enviado por el servidor: se hizo al suscribir.
  • Caducidad: cuando el token caduca, el servidor no corta la conexión de golpe (el mapa se congelaría en mitad del trayecto); envía {"tipo": "renovar"} y el cliente manda un {"accion": "token", "jwt": "..."} nuevo. Si no lo hace en 60 s, se cierra con código 4401.

Orden y deduplicación en el cliente

QoS 1 puede duplicar; Kafka con clave ordena por partición pero un reintento del puente puede repetir; la reconexión del cliente puede hacer que reciba la "última posición" retenida y a la vez la posición nueva. La solución es la de 01-05: cada posición lleva un número de secuencia por furgoneta (seq, un contador monótono que la app incrementa en cada envío y que sobrevive a reinicios porque se persiste localmente) además de fecha_ms. El cliente guarda ultima_seq por furgoneta y descarta todo lo que no sea estrictamente mayor. No hace falta ordenar: una posición vieja no aporta nada, se tira. El servidor hace lo mismo antes de difundir, para no gastar ancho de banda en duplicados.

Backpressure hacia clientes lentos

Un navegador en 3G no consume 28 mensajes por segundo (el panel de Jordi los recibiría todos). Si el servidor escribe sin control, el búfer de envío de esa conexión crece sin límite y acaba con la memoria de la instancia, que se lleva por delante a todos los demás clientes: el mismo problema de 05-04, en el último tramo. Tres defensas, de menor a mayor agresividad: coalescencia (para cada furgoneta solo importa la última posición: si el cliente tiene tres pendientes, se envía una), cola acotada por conexión (por ejemplo 100 mensajes; si se llena, se descartan los más antiguos y se cuenta en una métrica km0_ws_descartes_total), y cierre (si la cola está llena durante más de 30 s, el cliente no puede seguir el ritmo: se cierra con código 1013 Try again later y el cliente reconecta con backoff, quizá pidiendo menor frecuencia).

Escalado: conexiones por instancia y límites del sistema operativo

Una conexión WebSocket inactiva cuesta poco: un descriptor de fichero, unos 20-50 KB entre kernel y aplicación con asyncio. Una instancia de reparto con 2 GB puede sostener del orden de 20 000 a 50 000 conexiones, siempre que el sistema operativo lo permita: ulimit -n (descriptores por proceso, por defecto 1 024; en el contenedor se sube a 65 536 o más), net.core.somaxconn y net.ipv4.ip_local_port_range en el nodo, y el balanceador delante (Kong, el Ingress) que también mantiene una conexión por cliente y tiene sus propios límites. Lo que sí cuesta es el tráfico: 50 000 clientes recibiendo 1 mensaje/s son 50 000 escrituras por segundo por instancia; ahí el límite es la CPU de serializar y el ancho de banda, y se escala añadiendo instancias (el HPA de 07-05 con la métrica km0_ws_conexiones en lugar de CPU). Los despliegues son el momento delicado: un rolling update corta todas las conexiones de la instancia que se retira, así que el apagado ordenado (07-05) envía un cierre 1001 Going away escalonado durante 30 s y confía en la reconexión con jitter del cliente.

  1. Código: del broker MQTT al navegador

Broker Mosquitto con ACL por repartidor

# km0/docker-compose.yml (fragmento añadido en esta lección)
services:
  mosquitto:
    image: eclipse-mosquitto:2
    ports:
      - "8883:8883"     # MQTT sobre TLS para las furgonetas
      - "9001:9001"     # MQTT sobre WebSockets (diagnóstico)
    volumes:
      - ./borde/mosquitto/mosquitto.conf:/mosquitto/config/mosquitto.conf:ro
      - ./borde/mosquitto/acl:/mosquitto/config/acl:ro
      - ./borde/mosquitto/passwd:/mosquitto/config/passwd:ro
      - ./certs:/mosquitto/certs:ro          # emitidos por km0-ca (06-04)
      - mosquitto-data:/mosquitto/data       # sesiones persistentes y retained
volumes:
  mosquitto-data:
# km0/borde/mosquitto/mosquitto.conf
persistence true
persistence_location /mosquitto/data/
allow_anonymous false
password_file /mosquitto/config/passwd
acl_file /mosquitto/config/acl

listener 8883
cafile   /mosquitto/certs/km0-ca.crt
certfile /mosquitto/certs/mosquitto.crt
keyfile  /mosquitto/certs/mosquitto.key

listener 9001
protocol websockets
# km0/borde/mosquitto/acl
# Cada repartidor solo publica bajo su propio prefijo y solo lee sus comandos.
# %u se sustituye por el nombre de usuario con el que se autenticó.
pattern write km0/reparto/%u/posicion
pattern write km0/reparto/%u/estado
pattern read  km0/reparto/%u/comandos

# El puente lee todas las posiciones y estados, y escribe comandos.
user puente
topic read  km0/reparto/+/posicion
topic read  km0/reparto/+/estado
topic write km0/reparto/+/comandos

El fichero passwd se genera con mosquitto_passwd -c passwd furgoneta-3 (una contraseña por repartidor, emitida al dar de alta el dispositivo y guardada en Vault, 06-04). Con la ACL, una app comprometida de furgoneta-3 no puede publicar posiciones falsas de furgoneta-7 ni leer sus comandos: el broker lo rechaza antes de que llegue a la plataforma.

La app de la furgoneta: furgoneta_mqtt.py

# km0/servicios/reparto/furgoneta_mqtt.py
"""Cliente MQTT de la app del repartidor. Publica posiciones con QoS 1 y retained,
registra un last will y persiste el número de secuencia entre reinicios."""
import json, ssl, time, os
import paho.mqtt.client as mqtt

REPARTIDOR = os.environ["KM0_REPARTIDOR"]            # "furgoneta-3"
BROKER = os.environ.get("KM0_MQTT_HOST", "mqtt.km0.example")
T_POS = f"km0/reparto/{REPARTIDOR}/posicion"
T_EST = f"km0/reparto/{REPARTIDOR}/estado"
T_CMD = f"km0/reparto/{REPARTIDOR}/comandos"
FICHERO_SEQ = f"/var/lib/km0/{REPARTIDOR}.seq"        # el contador sobrevive a reinicios (01-05)


def leer_seq() -> int:
    try:
        return int(open(FICHERO_SEQ).read())
    except FileNotFoundError:
        return 0


def guardar_seq(seq: int) -> None:
    with open(FICHERO_SEQ + ".tmp", "w") as f:
        f.write(str(seq))
    os.replace(FICHERO_SEQ + ".tmp", FICHERO_SEQ)   # escritura atómica


def al_conectar(cliente, userdata, flags, rc, properties=None):
    print("conectado, sesión previa:", flags.get("session present", flags))
    cliente.subscribe(T_CMD, qos=1)                 # con clean_session=False, el broker la recuerda
    cliente.publish(T_EST, "conectada", qos=1, retain=True)


def al_mensaje(cliente, userdata, msg):
    comando = json.loads(msg.payload)
    print("comando recibido:", comando["tipo"])     # p. ej. {"tipo": "nueva_parada", "pedido": "P-2026-000126"}


def crear_cliente() -> mqtt.Client:
    c = mqtt.Client(client_id=REPARTIDOR, clean_session=False, protocol=mqtt.MQTTv311)
    c.username_pw_set(REPARTIDOR, os.environ["KM0_MQTT_PASS"])
    c.tls_set(ca_certs="/etc/km0/km0-ca.crt", tls_version=ssl.PROTOCOL_TLS_CLIENT)
    c.will_set(T_EST, "desconectada", qos=1, retain=True)   # last will: lo publica el broker si desaparecemos
    c.on_connect = al_conectar
    c.on_message = al_mensaje
    c.reconnect_delay_set(min_delay=1, max_delay=60)         # backoff de reconexión, lo gestiona paho
    return c


def bucle_posiciones(cliente: mqtt.Client, gps) -> None:
    seq = leer_seq()
    while True:
        lat, lon = gps.leer()
        seq += 1
        guardar_seq(seq)
        carga = json.dumps({"repartidor": REPARTIDOR, "seq": seq, "lat": lat, "lon": lon,
                            "fecha_ms": int(time.time() * 1000)})
        # QoS 1: paho reenvía si no hay PUBACK; si estamos en un túnel, encola y publica al volver.
        # retain=True: quien se suscriba después recibe la última posición sin esperar 5 s.
        cliente.publish(T_POS, carga, qos=1, retain=True)
        time.sleep(5)


if __name__ == "__main__":
    cliente = crear_cliente()
    cliente.connect_async(BROKER, 8883, keepalive=30)   # keepalive: PINGREQ cada 30 s; el broker detecta la caída en 45 s
    cliente.loop_start()                                # hilo de red: gestiona reconexiones y reenvíos
    from servicios.reparto.gps import GPS               # abstracción del receptor del dispositivo
    bucle_posiciones(cliente, GPS())

Detalles que importan: clean_session=False con un client_id fijo hace que el broker guarde la suscripción a comandos y encole los comandos QoS 1 durante el túnel; keepalive=30 es lo que da vida al last will (sin PINGREQ durante 1,5 × keepalive, el broker publica desconectada); y el contador seq se persiste antes de publicar, para que un reinicio de la app no reutilice un número (lo que haría que el cliente descartara posiciones nuevas como antiguas). Que la biblioteca encole mientras no hay red significa que, tras un túnel de 3 minutos, llegarán de golpe 36 posiciones antiguas: reparto las procesará por orden de seq y a Ana solo le llegará la última, gracias a la coalescencia.

El puente MQTT → Kafka: puente_mqtt_kafka.py

# km0/servicios/reparto/puente_mqtt_kafka.py
"""Suscribe a km0/reparto/+/posicion y publica cada posición en Kafka reparto.posiciones
con clave = repartidor. Es un proceso sin estado: se pueden correr varios (el broker
reparte con suscripciones compartidas en MQTT 5: $share/puente/km0/reparto/+/posicion)."""
import json, ssl, os
import paho.mqtt.client as mqtt
from confluent_kafka import Producer
from servicios.comun import metricas, logs

log = logs.obtener("reparto.puente")
productor = Producer({"bootstrap.servers": os.environ["KAFKA_BOOTSTRAP"],
                      "enable.idempotence": True, "acks": "all"})   # 02-05: sin duplicados por reintento
publicadas = metricas.contador("km0_puente_posiciones_total")


def entregado(err, msg):
    if err is not None:
        log.error("fallo al publicar en Kafka", error=str(err), clave=msg.key())


def al_mensaje(cliente, userdata, msg):
    partes = msg.topic.split("/")                  # ["km0", "reparto", "furgoneta-3", "posicion"]
    repartidor = partes[2]
    pos = json.loads(msg.payload)
    if pos.get("repartidor") != repartidor:        # la ACL ya lo impide, pero no nos fiamos del payload
        log.warning("payload con repartidor distinto del tópico", topico=msg.topic)
        return
    envoltura = {"id_evento": f"{repartidor}-{pos['seq']}",   # determinista: mismo evento, mismo id (02-05)
                 "tipo": "posicion.actualizada", "version": 1,
                 "fecha_ms": pos["fecha_ms"], "origen": "puente-mqtt", "datos": pos}
    productor.produce("reparto.posiciones", key=repartidor.encode(),
                      value=json.dumps(envoltura).encode(), callback=entregado)
    productor.poll(0)
    publicadas.inc()


cliente = mqtt.Client(client_id=f"puente-{os.getpid()}", protocol=mqtt.MQTTv5)
cliente.username_pw_set("puente", os.environ["KM0_MQTT_PASS_PUENTE"])
cliente.tls_set(ca_certs="/etc/km0/km0-ca.crt", tls_version=ssl.PROTOCOL_TLS_CLIENT)
cliente.on_message = al_mensaje
cliente.connect(os.environ.get("KM0_MQTT_HOST", "mosquitto"), 8883)
cliente.subscribe("$share/puente/km0/reparto/+/posicion", qos=1)   # suscripción compartida: varias réplicas
cliente.loop_forever()

El id_evento se construye de forma determinista con repartidor-seq en lugar de un UUID aleatorio: así, si el broker reenvía un PUBLISH (QoS 1) o el puente se reinicia y lo reprocesa, el evento en Kafka lleva el mismo id y los consumidores idempotentes de 02-05 lo descartan. Es un detalle pequeño que elimina toda una clase de duplicados.

El servidor WebSocket: ws_servidor.py

# km0/servicios/reparto/ws_servidor.py
"""Servidor WebSocket de reparto con FastAPI: autenticación por JWT, suscripción por
pedido o por mercado, fan-out entre instancias con Redis pub/sub, ping/pong,
deduplicación por seq y cola acotada por conexión."""
import asyncio, json, os
from collections import defaultdict
from fastapi import FastAPI, WebSocket, WebSocketDisconnect
import redis.asyncio as redis
from servicios.comun import metricas, logs
from servicios.comun.jwt import verificar_jwt, JwtInvalido      # 06-01
from servicios.reparto.asignaciones import furgoneta_de_pedido   # modelo de reparto

app = FastAPI()
log = logs.obtener("reparto.ws")
r = redis.from_url(os.environ["REDIS_URL"])
conexiones_g = metricas.gauge("km0_ws_conexiones")
descartes = metricas.contador("km0_ws_descartes_total")
COLA_MAX = 100


class Conexion:
    def __init__(self, ws: WebSocket, claims: dict):
        self.ws, self.claims = ws, claims
        self.cola: asyncio.Queue = asyncio.Queue(maxsize=COLA_MAX)
        self.ultima_seq: dict[str, int] = {}       # por furgoneta: deduplicación (01-05)
        self.temas: set[str] = set()

    def encolar(self, tema: str, mensaje: dict) -> None:
        furgo, seq = mensaje.get("repartidor"), mensaje.get("seq")
        if furgo and seq is not None:
            if seq <= self.ultima_seq.get(furgo, -1):
                return                              # duplicado o antiguo: se descarta
            self.ultima_seq[furgo] = seq
        if self.cola.full():                        # backpressure: cliente lento
            self.cola.get_nowait()                  # tirar el más antiguo
            descartes.inc()
        self.cola.put_nowait(mensaje)


# Tabla de suscripciones LOCAL a esta instancia: tema -> conexiones
suscriptores: dict[str, set[Conexion]] = defaultdict(set)
pubsub = r.pubsub()


async def bucle_bus() -> None:
    """Recibe del bus Redis lo publicado por cualquier instancia y lo reparte localmente."""
    async for m in pubsub.listen():
        if m["type"] not in ("message", "pmessage"):
            continue
        tema = m["channel"].decode()
        mensaje = json.loads(m["data"])
        for c in list(suscriptores.get(tema, ())):
            c.encolar(tema, mensaje)


@app.on_event("startup")
async def arrancar():
    asyncio.create_task(bucle_bus())


async def suscribir(c: Conexion, tema: str) -> None:
    if not suscriptores[tema]:
        await pubsub.subscribe(tema)                # primera conexión local interesada: suscribir al bus
    suscriptores[tema].add(c)
    c.temas.add(tema)


async def desuscribir_todo(c: Conexion) -> None:
    for tema in c.temas:
        suscriptores[tema].discard(c)
        if not suscriptores[tema]:
            await pubsub.unsubscribe(tema)          # nadie más aquí: dejar de recibir del bus
            del suscriptores[tema]


def autorizar(claims: dict, accion: dict) -> str | None:
    """Devuelve el tema del bus al que se traduce la suscripción, o None si no está permitido."""
    if "pedido" in accion:
        pedido = accion["pedido"]
        asignacion = furgoneta_de_pedido(pedido)                 # {"cliente": "u-ana", "furgoneta": "furgoneta-3"}
        if asignacion and asignacion["cliente"] == claims["sub"]:
            return f"furgoneta:{asignacion['furgoneta']}"
    if "mercado" in accion and "operador" in claims.get("roles", []):
        return f"mercado:{accion['mercado']}"
    return None


async def emisor(c: Conexion) -> None:
    """Tarea que vacía la cola hacia el socket. Separada de la recepción para no bloquearla."""
    while True:
        mensaje = await c.cola.get()
        await c.ws.send_text(json.dumps(mensaje))


@app.websocket("/ws")
async def ws_endpoint(ws: WebSocket):
    await ws.accept()
    # 1. Primer mensaje: el token (no en la URL, para que no acabe en los logs de Kong).
    try:
        primero = json.loads(await asyncio.wait_for(ws.receive_text(), timeout=5))
        claims = verificar_jwt(primero["jwt"], audiencia="reparto")
    except (asyncio.TimeoutError, KeyError, JwtInvalido):
        await ws.close(code=4401); return
    c = Conexion(ws, claims)
    conexiones_g.inc()
    tarea_emisor = asyncio.create_task(emisor(c))
    log.info("ws conectado", sub=claims["sub"])
    try:
        while True:
            accion = json.loads(await ws.receive_text())
            if accion.get("accion") == "suscribir":
                tema = autorizar(claims, accion)
                if tema is None:
                    await ws.send_text(json.dumps({"tipo": "error", "codigo": "no_autorizado"})); continue
                await suscribir(c, tema)
                ultima = await r.get(f"ultima:{tema}")          # última posición conocida (como el retained de MQTT)
                if ultima:
                    c.encolar(tema, json.loads(ultima))
            elif accion.get("accion") == "token":
                claims = verificar_jwt(accion["jwt"], audiencia="reparto")   # renovación sin cortar
                c.claims = claims
    except WebSocketDisconnect:
        pass
    finally:
        tarea_emisor.cancel()
        await desuscribir_todo(c)
        conexiones_g.dec()
        log.info("ws cerrado", sub=claims["sub"])

Y el consumidor de Kafka que alimenta el bus (una tarea en cada instancia, todas en el mismo grupo de consumidores, de modo que cada partición la lee una):

# km0/servicios/reparto/ws_consumidor_kafka.py
"""Lee reparto.posiciones y reparto.panel (grupo reparto-ws) y publica en Redis pub/sub."""
import json, os
from confluent_kafka import Consumer
import redis

r = redis.from_url(os.environ["REDIS_URL"])
c = Consumer({"bootstrap.servers": os.environ["KAFKA_BOOTSTRAP"], "group.id": "reparto-ws",
              "auto.offset.reset": "latest"})       # a nadie le interesan posiciones viejas al arrancar
c.subscribe(["reparto.posiciones", "reparto.panel"])
while True:
    msg = c.poll(1.0)
    if msg is None or msg.error():
        continue
    ev = json.loads(msg.value())
    d = ev["datos"]
    if msg.topic() == "reparto.posiciones":
        carga = json.dumps({"tipo": "posicion", **d})
        r.publish(f"furgoneta:{d['repartidor']}", carga)
        r.publish(f"mercado:{d.get('mercado', 'desconocido')}", carga)
        r.set(f"ultima:furgoneta:{d['repartidor']}", carga, ex=600)   # equivalente al retained
    else:
        r.publish(f"mercado:{d['mercado']}", json.dumps({"tipo": "panel", **d}))
    c.commit(msg)

El cliente del navegador, con reconexión y backoff

// km0/borde/web/seguimiento.js — cliente WebSocket mínimo para la pantalla de seguimiento de Ana
class Seguimiento {
  constructor(pedido, obtenerJwt, alPosicion) {
    this.pedido = pedido; this.obtenerJwt = obtenerJwt; this.alPosicion = alPosicion;
    this.intento = 0; this.ultimaSeq = -1; this.cerrado = false;
    this.conectar();
  }
  conectar() {
    this.ws = new WebSocket("wss://api.km0.example/ws");        // pasa por Kong (Upgrade permitido)
    this.ws.onopen = () => {
      this.intento = 0;                                          // conexión lograda: reiniciar backoff
      this.ws.send(JSON.stringify({ jwt: this.obtenerJwt() }));  // 1º: autenticar
      this.ws.send(JSON.stringify({ accion: "suscribir", pedido: this.pedido }));
    };
    this.ws.onmessage = (e) => {
      const m = JSON.parse(e.data);
      if (m.tipo === "renovar") { this.ws.send(JSON.stringify({ accion: "token", jwt: this.obtenerJwt() })); return; }
      if (m.tipo === "posicion") {
        if (m.seq <= this.ultimaSeq) return;                     // duplicado o antiguo (01-05)
        this.ultimaSeq = m.seq;
        this.alPosicion(m.lat, m.lon, m.fecha_ms);
      }
    };
    this.ws.onclose = (e) => {
      if (this.cerrado || e.code === 4401) return;               // cierre voluntario o no autorizado: no reintentar
      const base = Math.min(30000, 1000 * 2 ** this.intento++);  // exponencial con tope de 30 s
      const espera = base / 2 + Math.random() * base / 2;        // jitter: evitar la estampida (07-04)
      setTimeout(() => this.conectar(), espera);
    };
    this.ws.onerror = () => this.ws.close();
  }
  cerrar() { this.cerrado = true; this.ws.close(1000); }
}
// Uso: new Seguimiento("P-2026-000124", () => sesion.jwt, (lat, lon) => mapa.mover(lat, lon));

El navegador no implementa ping/pong visible desde JavaScript (lo hace por debajo cuando el servidor envía ping), así que la detección de una conexión "zombi" desde el cliente se apoya en el servidor: si no llega ningún mensaje en 60 s, el cliente puede cerrar y reconectar. ultimaSeq en el cliente es la última línea de defensa contra duplicados: aunque el servidor ya filtre, al reconectar a otra instancia se recibe la "última conocida" que puede ser una ya vista.

  1. SSE como alternativa para las alertas de stock

Las alertas de inventario.alertas a productores son el caso opuesto al mapa: pocos mensajes, solo servidor → cliente, y un panel que puede estar abierto horas. WebSockets sería sobredimensionar; SSE encaja exactamente, con reconexión gratis y Last-Event-ID para no perder alertas durante una desconexión.

# km0/servicios/inventario/sse_alertas.py
from fastapi import FastAPI, Request, Depends
from sse_starlette.sse import EventSourceResponse
import asyncio, json, redis.asyncio as redis
from servicios.comun.jwt import claims_de_peticion   # 06-01: Bearer en la cabecera, como cualquier GET

app = FastAPI()
r = redis.from_url("redis://redis:6379")


@app.get("/api/v1/inventario/alertas/stream")
async def stream(request: Request, claims=Depends(claims_de_peticion)):
    productor = claims["productor_id"]                            # p. ej. "queseria-montblanc"
    ultimo = request.headers.get("Last-Event-ID", "0-0")      # el navegador lo envía al reconectar

    async def generador():
        cursor = ultimo
        while not await request.is_disconnected():
            # Redis Streams (no pub/sub): tiene histórico, así que la reconexión no pierde alertas.
            res = await r.xread({f"alertas:{productor}": cursor}, block=15000, count=10)
            if not res:
                yield {"comment": "keepalive"}                 # evita que proxies cierren por inactividad
                continue
            for _, entradas in res:
                for id_, campos in entradas:
                    cursor = id_
                    yield {"id": id_, "event": "stock_bajo", "data": campos[b"json"].decode()}
    return EventSourceResponse(generador())
// Panel del productor: el navegador reconecta solo y reenvía Last-Event-ID
const es = new EventSource("/api/v1/inventario/alertas/stream");   // el JWT va en cookie o vía Kong
es.addEventListener("stock_bajo", (e) => { const a = JSON.parse(e.data); mostrarAviso(a.producto, a.mercado, a.disponible); });

Un consumidor de inventario.alertas (el que en 05-04 escribía las alertas de Flink) añade cada alerta al stream alertas:<productor> con XADD y un MAXLEN de unos cientos de entradas. La diferencia con el mapa es que aquí sí importa no perder ninguna: por eso Redis Streams (con histórico e ids monótonos) y no pub/sub.

  1. Cuándo usar cada técnica

Necesidad Técnica Por qué
Dato que cambia cada varios minutos, pocos clientes Polling con Cache-Control Simplicidad; el gateway cachea
Servidor → cliente, texto, no perder mensajes, sin volumen SSE + Redis Streams Reconexión y reanudación de serie; HTTP normal
Bidireccional o alto volumen hacia el navegador WebSockets + bus entre instancias Tramas baratas, un canal para todo (mapa, chat, comandos)
Dispositivos en redes malas, con batería MQTT (QoS 1, sesiones, last will) Diseñado para eso; el broker aísla y autentica
Navegador que debe hablar con un broker MQTT MQTT sobre WebSockets Único transporte disponible en el navegador
Mensajes que deben sobrevivir a la desconexión del cliente Persistir primero (Streams, base de datos); el canal solo notifica Ni WebSocket ni Redis pub/sub guardan nada
Respaldo cuando WebSocket está bloqueado (proxies corporativos) Long polling Funciona donde nada más funciona

Errores Comunes y Consejos

  • Exponer Kafka o el broker MQTT al navegador. Ni protocolo, ni seguridad, ni número de conexiones. Siempre un servidor de último tramo que autentica, autoriza por suscripción y enriquece.
  • Pasar el JWT en la URL del WebSocket. Acaba en los logs de Kong y del servidor (07-02). Primer mensaje tras conectar, y renovación sin cortar.
  • Un grupo de consumidores de Kafka por instancia de WebSocket sin bus. O cada instancia recibe solo una parte y sus clientes se quedan sin datos, o cada una consume todo y no escala. Bus entre instancias (Redis pub/sub o Kafka).
  • Escribir al socket sin cola acotada. Un cliente en 3G agota la memoria de la instancia y tumba a todos. Cola por conexión, coalescencia, métrica de descartes y cierre 1013.
  • Reconexión sin backoff ni jitter. Un despliegue de reparto provoca que miles de clientes reconecten en el mismo segundo. Exponencial con tope y jitter, y apagado ordenado escalonado en el servidor.
  • QoS 2 "por seguridad". Cuatro mensajes por posición, sesión pesada en el broker, y de todos modos el consumidor tiene que deduplicar por otras razones. QoS 1 y seq.
  • Olvidar el last will y el keepalive. Sin ellos, un repartidor sin cobertura aparece "conectado" durante los minutos que TCP tarda en enterarse. keepalive=30 y testamento en estado.
  • Consejo: mide km0_ws_conexiones, km0_ws_descartes_total, el lag del grupo reparto-ws y la latencia de extremo a extremo (fecha_ms de la furgoneta frente a la hora de recepción en el navegador, con el reloj corregido): es el SLO real del "tiempo real".
  • Consejo: en el diseño, empieza por la tabla del apartado 1. Latencia aceptable, volumen y sentido deciden la técnica antes que cualquier preferencia por WebSockets.

Ejercicios

Ejercicio 1: el chat cliente-productor

Diseña el chat entre Ana y la Quesería Montblanc sobre la infraestructura de esta lección: (a) qué transporte usa cada extremo y por qué; (b) qué se persiste, dónde y en qué orden respecto a la notificación por el bus; (c) qué hace el cliente al reconectar tras 2 minutos sin red para no perder ni duplicar mensajes; (d) qué tema del bus y qué regla de autorización aplica ws_servidor.py.

Ejercicio 2: el túnel de tres minutos

furgoneta-3 atraviesa un túnel de 3 minutos. Describe, paso a paso y nombrando los mecanismos, qué ocurre en: la app (furgoneta_mqtt.py), el broker (last will, retained, sesión), el puente, reparto.posiciones, Flink (la ventana de sesión de 05-04), ws_servidor.py y la pantalla de Ana. ¿Cuántos mensajes recibe Ana al salir del túnel y por qué?

Ejercicio 3: dimensionar el panel de operadores

En campaña hay 140 furgonetas enviando cada 5 s y 40 operadores con el panel abierto, cada uno suscrito al mercado completo (35 furgonetas de media). (a) ¿Cuántos mensajes por segundo recibe cada operador y cuántos escribe en total el conjunto de instancias de ws_servidor? (b) Si un operador está en una conexión que solo admite 2 mensajes/s, ¿qué ocurre con la cola de 100 y con qué mecanismo se resuelve sin cerrar la conexión? (c) Propón un cambio en el consumidor de Kafka o en el servidor que reduzca el tráfico al panel a 1 mensaje/s por operador manteniendo la información útil.

Soluciones

Ejercicio 1.

(a) Ambos extremos son navegadores (Ana en el móvil, la quesería en su panel): WebSockets en los dos, porque el chat es bidireccional y de baja latencia, y porque Ana ya tiene abierta la conexión del seguimiento: se reutiliza el mismo canal con otro tipo de mensaje ({"accion": "chat", "conversacion": "P-2026-000124", "texto": "..."}).

(b) Cada mensaje se persiste primero en pedidos o en un contexto mensajes (tabla mensajes_por_conversacion en Cassandra con clave de partición = conversación y clustering por un id monótono, por ejemplo un TimeUUID o un seq por conversación asignado por el servidor), y solo después se publica en el bus conversacion:P-2026-000124. Si se publicara antes de persistir y el proceso muriese entre medias, el otro extremo vería un mensaje que no existe.

(c) Al reconectar, el cliente envía {"accion": "suscribir", "conversacion": "...", "desde": <último id visto>}; el servidor responde con los mensajes persistidos posteriores (una consulta por clave de partición, barata) y después suscribe al bus. El cliente deduplica por id de mensaje (puede recibir uno tanto por la recuperación como por el bus si llegó justo en ese instante). Es el mismo Last-Event-ID de SSE, hecho a mano.

(d) Tema conversacion:<pedido>; autorización: claims["sub"] es el cliente del pedido, o claims["productor_id"] es el productor de alguna línea del pedido (consulta al modelo de pedidos en el momento de suscribir, no por mensaje).

Ejercicio 2.

  1. App: pierde el PINGREQ; paho detecta la caída y entra en reconexión con backoff (1 s → 60 s). El bucle de posiciones sigue: incrementa seq, persiste y llama a publish con QoS 1; paho encola en memoria (~36 mensajes en 3 minutos).
  2. Broker: a los 45 s sin keepalive (1,5 × 30) da por muerta la conexión y publica el last will km0/reparto/furgoneta-3/estado = desconectada (retained). La sesión persistente conserva la suscripción a comandos y encola cualquier comando QoS 1 que el operador envíe.
  3. Puente: recibe el estado desconectada y lo publica en Kafka (como evento estado.actualizado); no recibe posiciones.
  4. reparto.posiciones: ningún evento de furgoneta-3 durante 3 minutos.
  5. Flink: la ventana de sesión con gap de 3 minutos se cierra (según la marca de agua) y emite "repartidor sin señal" a reparto.panel; el panel de Jordi lo muestra.
  6. ws_servidor: no envía nada a Ana; su conexión sigue viva (ping/pong con el servidor, que sí tiene red). El mapa muestra la última posición y, si el cliente lo implementa, "última señal hace 2 min".
  7. Salida del túnel: paho reconecta (la sesión estaba presente), recibe los comandos encolados y vacía su cola: 36 PUBLISH QoS 1 en ráfaga, en orden de seq. El broker actualiza el retained con la última. El puente publica los 36 en Kafka con id_evento determinista. Flink abre una sesión nueva ("ha vuelto"). El consumidor de ws_servidor publica 36 mensajes en el bus; en la Conexion de Ana, encolar acepta cada uno porque seq es creciente, pero como llegan en milisegundos y la cola tiene 100 de capacidad, se encolan todos: Ana recibiría 36 mensajes en ráfaga, y el mapa "saltaría" por el túnel. Para que reciba solo la última haría falta la coalescencia por furgoneta (mantener en la cola solo la última posición de cada furgoneta), que el código del apartado 8 no implementa y que el ejercicio 3(c) introduce. Sin ella, no hay error, solo tráfico inútil.

Ejercicio 3.

(a) 140 furgonetas / 5 s = 28 posiciones/s en total. Cada operador, con 35 furgonetas: 7 mensajes/s. Escrituras totales: 40 × 7 = 280 mensajes/s entre todas las instancias (más las de los clientes con seguimiento). Es poco: una sola instancia lo sostiene; el problema no es el volumen sino la robustez.

(b) Entran 7/s y salen 2/s: la cola crece 5/s y se llena en 20 s. A partir de ahí, encolar descarta el más antiguo por cada nuevo (km0_ws_descartes_total sube 5/s) y el operador ve posiciones con hasta 100/7 ≈ 14 s de retraso, pero la conexión sigue. Se resuelve sin cerrar con coalescencia: si de una misma furgoneta hay una posición pendiente en la cola, la nueva la reemplaza en lugar de añadirse. Con 35 furgonetas, la cola nunca supera 35 entradas y el operador recibe siempre la posición más reciente de cada una, con retraso acotado.

(c) Dos opciones. En el servidor: un emisor por conexión que, en lugar de vaciar la cola mensaje a mensaje, cada segundo agrupa lo pendiente en un único mensaje {"tipo": "posiciones", "items": [...]} con la última posición de cada furgoneta (coalescencia + agrupación: 1 mensaje/s, 35 posiciones dentro). En el consumidor de Kafka: no publicar en mercado:<m> cada posición sino mantener en Redis un hash mercado:<m>:posiciones (HSET furgoneta-3 <json>) y publicar un "tick" por segundo; el servidor, al recibir el tick, lee el hash y envía el estado completo. La primera es más simple y mantiene el bus sin cambios; la segunda desacopla la frecuencia del panel de la de las furgonetas y es más adecuada si el panel crece a cientos de operadores.

Conclusión

El último tramo tiene reglas propias porque sus clientes son navegadores, móviles y dispositivos en redes que fallan, en números que ningún servicio interno alcanza, y sin la confianza ni las bibliotecas de un consumidor de Kafka. "Tiempo real" aquí es latencia percibida, fijada por caso: segundos para el mapa de Ana, uno o dos para el panel de Jordi, decenas para las alertas de la quesería, menos de uno para el chat. Con esa tabla delante, la técnica se elige sola: polling y long polling como respaldo, SSE para lo unidireccional que no puede perderse (con Redis Streams y Last-Event-ID), WebSockets para lo bidireccional y voluminoso, y MQTT para los dispositivos, con tópicos jerárquicos que son también ACL, QoS 1 con deduplicación por secuencia, mensajes retenidos, testamento y sesiones persistentes que sobreviven al túnel. La arquitectura de extremo a extremo encadena la furgoneta, Mosquitto, el puente hacia reparto.posiciones, Flink y el servidor WebSocket de reparto, que es el único que sabe quién quiere qué; el fan-out entre sus instancias se resuelve con un bus (Redis pub/sub) en lugar de sesiones pegajosas, el JWT se verifica en el handshake y se renueva sin cortar, el orden y los duplicados se resuelven con el seq de 01-05 en servidor y cliente, la presión de los clientes lentos con colas acotadas y coalescencia, y el escalado con instancias intercambiables, límites del sistema operativo ajustados y apagados escalonados.

Todo lo construido hasta aquí, desde el broker MQTT hasta el clúster de Kubernetes, corre en máquinas que alguien ha comprado, instalado y mantiene. La siguiente lección cambia esa premisa: qué ocurre cuando la infraestructura se convierte en una API de un proveedor de nube, qué servicios gestionados sustituyen a cada pieza que hemos montado a mano, cómo se describe todo con Terraform, y qué cuesta al mes. Es la lección de Aplicaciones en la Nube.

Curso de Arquitecturas Distribuidas

Módulo 1: Introducción a los Sistemas Distribuidos

Módulo 2: Comunicación en Sistemas Distribuidos

Módulo 3: Consistencia y Replicación

Módulo 4: Almacenamiento Distribuido

Módulo 5: Computación Distribuida

Módulo 6: Seguridad en Sistemas Distribuidos

Módulo 7: Monitoreo y Mantenimiento

Módulo 8: Casos de Estudio y Aplicaciones

© Copyright 2026. Todos los derechos reservados