La lliçó anterior va deixar el mapa de contextos de Quilòmetre Zero amb una fletxa que acabava a Kafka: repartiment publica a repartiment.posicions, Flink calcula sobre aquest tòpic i escriu a repartiment.panell (05-04), i aquí s'acabava el recorregut. Però l'Anna no llegeix Kafka. L'Anna té oberta la pantalla de seguiment de la seva comanda al mòbil, a l'autobús, amb cobertura irregular, i vol veure com furgoneta-3 es mou pel mapa sense prémer "actualitza". Jordi Sala té el panell d'operadors en un navegador d'escriptori amb 140 repartidors alhora. L'app de la furgoneta envia la seva posició cada 5 segons des d'una xarxa mòbil que cau a cada túnel. I la Formatgeria Montblanc vol un avís al seu panell tan bon punt formatge-curat baixi del llindar a Girona. Tots són problemes de l'últim tram: entre la plataforma (serveis, Kafka, bases de dades, tot en un centre de dades amb xarxa fiable) i les persones i dispositius que són a fora, en xarxes dolentes, amb milers o desenes de milers de connexions simultànies, i sense poder executar un consumidor de Kafka. Aquesta lliçó desenvolupa el que 02-01 va deixar presentat: MQTT per als dispositius i WebSockets per als navegadors, juntament amb les alternatives (polling, long polling, Server-Sent Events), i els uneix en una arquitectura pub/sub d'extrem a extrem amb els seus problemes propis: el fan-out entre instàncies, l'autenticació de connexions llargues, l'ordre i els duplicats al client, la pressió sobre clients lents i l'escalat de connexions. El núvol (08-03) i el processament a la vora (08-04) en queden fora.
Contingut
- Què significa "temps real" aquí i els quatre casos de Quilòmetre Zero
- L'últim tram davant de la missatgeria entre serveis
- Tècniques de l'últim tram: polling, long polling, SSE i WebSockets
- MQTT per a dispositius
- Arquitectura pub/sub d'extrem a extrem
- Fan-out amb moltes instàncies, presència i subscripcions
- Autenticació, ordre, deduplicació, backpressure i escalat
- Codi: del broker MQTT al navegador
- SSE com a alternativa per a les alertes d'estoc
- Quan usar cada tècnica
- Errors Comuns i Consells
- Exercicis
- Conclusió
- Què significa "temps real" aquí i els quatre casos de Quilòmetre Zero
"Temps real" té un significat estricte en enginyeria de control: un sistema de temps real dur és el que garanteix que una resposta arriba abans d'un termini, i si no arriba, el resultat és una fallada (l'airbag, el control d'un motor). Res d'això no s'aplica aquí. En sistemes web, "temps real" significa latència percebuda per persones: que la informació arribi prou ràpid perquè qui la mira senti que veu el que està passant, sense haver-la de demanar. Aquest llindar és entre dècimes de segon i uns pocs segons, i varia segons el cas. Convé fixar-lo cas per cas, perquè el disseny canvia amb ell.
| Cas | Qui rep | Origen de la dada | Latència acceptable | Volum | Sentit |
|---|---|---|---|---|---|
L'Anna veu moure's furgoneta-3 al mapa |
Navegador o app mòbil d'un client, en xarxa mòbil | Posició de la furgoneta (MQTT → Kafka) | 2-5 s (una posició cada 5 s) | Milers de clients amb el seguiment obert en campanya; cadascun interessat en una furgoneta | Servidor → client |
Panell d'operadors de repartiment |
Navegador d'escriptori del Jordi i la Marta | repartiment.panell (Flink, 05-04) i posicions |
1-2 s | Desenes d'operadors; cadascun interessat en totes les furgonetes del seu mercat | Servidor → client |
| Avisos d'estoc a productors | Panell web del productor | inventari.alertes (Flink alertes_estoc.py) |
10-30 s | Centenars de productors; pocs missatges al dia per productor | Servidor → client |
| Xat client-productor | Navegador de l'Anna i panell de la formatgeria | Missatges escrits per persones | < 1 s | Converses curtes; bidireccional | Tots dos sentits |
I un cinquè, del costat de l'origen: l'app de la furgoneta publica la seva posició cada 5 segons, des d'una xarxa mòbil que perd cobertura, amb una bateria que cal cuidar, i ha de continuar funcionant quan torna la cobertura sense perdre el que és important. Té requisits oposats als de l'Anna: pocs bytes, una sola connexió persistent, tolerància a desconnexions llargues.
- L'últim tram davant de la missatgeria entre serveis
La missatgeria del Mòdul 2 (RabbitMQ, Kafka) resol la comunicació entre serveis: processos de llarga vida, a la mateixa xarxa, amb clients de biblioteca que mantenen connexions, offsets i grups de consumidors. L'últim tram és diferent en tot:
| Aspecte | Entre serveis (02-04) | Últim tram |
|---|---|---|
| Client | Un procés Python amb confluent_kafka, a Kubernetes |
Un navegador (només HTTP i WebSocket), una app mòbil, un dispositiu amb poca bateria |
| Xarxa | Centre de dades: fiable, baixa latència | Internet, 4G, Wi-Fi de bar: pèrdues, NAT, proxies, túnels |
| Nombre de connexions | Desenes | Milers a milions |
| Vida de la connexió | Hores o dies | Minuts; es talla en canviar de xarxa o en bloquejar el mòbil |
| Estat del consumidor | Offset persistit al broker; reprèn on ho va deixar | Cap o mínim; en reconnectar, cal decidir què s'ha perdut |
| Seguretat | mTLS entre serveis (06-04) | Un usuari amb un JWT (06-01); el client no és de confiança |
| Fan-out | Un grup de consumidors per servei | Cada persona vol un subconjunt diferent dels missatges |
Per això mai no s'exposa Kafka directament a un navegador: ni el protocol ho permet, ni el model de seguretat, ni el nombre de connexions. Entre Kafka i el navegador cal un component que parli el protocol del client, que mantingui les seves connexions i que reparteixi a cadascun només el que li correspon. Aquest component és el servidor de temps real que construirem com a part de repartiment.
- Tècniques de l'últim tram: polling, long polling, SSE i WebSockets
Totes parteixen d'una limitació: HTTP és petició-resposta i el servidor no pot parlar primer. Les quatre tècniques són maneres de capgirar aquesta limitació, cadascuna amb un cost.
Polling
El client pregunta cada n segons: GET /api/v1/repartiment/comandes/P-2026-000124/posicio. És trivial, es pot desar en memòria cau al gateway, compatible amb tot, i és exactament el que el monòlit feia al símptoma 4 de 01-06 amb el resultat conegut: amb 5 000 clients preguntant cada 3 segons són 1 700 peticions per segon, la majoria per rebre "res de nou", cadascuna amb el seu handshake TLS si no es reutilitza la connexió (02-01) i els seus 300 bytes de capçaleres. La latència mitjana és la meitat de l'interval.
Long polling
El client pregunta, però el servidor no respon fins que hi hagi alguna cosa nova (o fins a un timeout de, per exemple, 30 s), i el client torna a preguntar immediatament. Redueix les peticions buides i la latència baixa gairebé a zero, però manté una petició HTTP oberta per client (un worker bloquejat en servidors síncrons), pateix amb proxies que tallen connexions inactives, i cada missatge continua costant una petició completa. Va ser la tècnica dominant abans dels WebSockets, i continua sent vàlida com a respatller.
Server-Sent Events (SSE)
El client fa un GET amb Accept: text/event-stream i el servidor deixa la resposta oberta indefinidament, escrivint esdeveniments en un format de text senzill (event:, data:, id:, separats per línia en blanc). És HTTP normal (travessa proxies i gateways, va sobre HTTP/2 multiplexat), el navegador l'implementa amb EventSource, que reconnecta sola i envia la capçalera Last-Event-ID amb l'últim id rebut perquè el servidor reprengui. El seu límit: és unidireccional (servidor → client; el client, si vol enviar, usa peticions HTTP normals) i només text.
WebSockets
El que es va presentar a 02-01: una petició HTTP amb Upgrade: websocket que, si el servidor accepta (101 Switching Protocols), converteix la connexió TCP en un canal bidireccional, persistent i amb trames (binàries o de text), sense capçaleres HTTP per missatge (de 2 a 14 bytes de sobrecàrrega per trama). Convé entendre les quatre peces del protocol que afecten el disseny:
- Handshake: el client envia
Sec-WebSocket-Key; el servidor respon ambSec-WebSocket-Acceptderivat d'ella. És l'únic moment en què hi ha capçaleres HTTP, i per tant l'únic en què Kong pot aplicar els seus plugins (06-05): el JWT es verifica al handshake, no per missatge. - Trames: cada missatge viatja en una o més trames amb un codi d'operació (text, binari, tancament, ping, pong). Les trames del client van emmascarades (XOR amb una clau aleatòria) per una raó de seguretat davant de proxies antics, no de xifratge; el xifratge és TLS (
wss://). - Ping/pong: trames de control que qualsevol dels dos costats envia per comprovar que l'altre continua allà. Sense elles, un mòbil que perd cobertura deixa una connexió "oberta" al servidor durant els minuts que TCP triga a adonar-se'n (02-01). El servidor de Quilòmetre Zero fa ping cada 20 s i tanca si no hi ha pong en 10 s.
- Reconnexió: el protocol no la defineix. Quan la connexió cau, és el client qui reconnecta, amb backoff exponencial i jitter (els mateixos de 07-04, per la mateixa raó: 5 000 clients reconnectant alhora després d'un desplegament és una estampida), i qui ha de dir al servidor "quina és l'última cosa que he vist" per recuperar el que s'ha perdut. Tot això és codi d'aplicació.
| Tècnica | Sentit | Latència | Cost per missatge | Connexions obertes | Travessa proxies | Reconnexió | Ús a Quilòmetre Zero |
|---|---|---|---|---|---|---|---|
| Polling | El client demana | Mitjana = interval / 2 | Una petició HTTP completa (amb dada o sense) | No (o una de reutilitzada) | Sí | Trivial (sense estat) | Respatller si WebSocket falla; dades que canvien cada pocs minuts |
| Long polling | El client demana, el servidor reté | Gairebé immediata | Una petició per missatge | Una per client, en espera | Amb timeouts curts | Trivial | Respatller |
| SSE | Servidor → client | Immediata | Unes desenes de bytes | Una per client | Sí (HTTP/2 millor) | Automàtica (Last-Event-ID) |
Alertes d'estoc a productors |
| WebSockets | Bidireccional | Immediata | 2-14 bytes de trama | Una per client | Sí, amb Upgrade permès al gateway |
Manual (codi de client) | Mapa de l'Anna, panell d'operadors, xat |
| MQTT (sobre TCP o WebSockets) | Pub/sub bidireccional | Immediata | 2 bytes de capçalera fixa | Una per dispositiu | Sobre WebSockets, sí | Del client, amb sessió persistent al broker | App de la furgoneta |
- MQTT per a dispositius
02-01 va presentar MQTT com a protocol pub/sub per a dispositius amb recursos limitats i xarxes dolentes. Aquí el dissenyem per a la flota de Quilòmetre Zero, recolzant-nos en les set característiques que el fan adequat:
- Broker central. Els dispositius no es coneixen entre ells ni coneixen els consumidors: publiquen al broker (Mosquitto per començar; EMQX o HiveMQ per a centenars de milers de connexions) i el broker reparteix. La furgoneta no sap que Kafka existeix.
- Tòpics jeràrquics.
km0/repartiment/furgoneta-3/posicio,km0/repartiment/furgoneta-3/estat,km0/repartiment/furgoneta-3/ordres. Els comodins permeten subscriure's a nivells:+un nivell (km0/repartiment/+/posicio: totes les posicions),#tot el que en penja (km0/repartiment/furgoneta-3/#). La jerarquia és a més la base de les ACL:furgoneta-3només pot publicar sota el seu prefix. - QoS per missatge. QoS 0 (at most once): s'envia i s'oblida; adequat per a la posició cada 5 s, perquè la següent la reemplaça. QoS 1 (at least once): el broker confirma amb
PUBACKi el client reenvia si no arriba; pot duplicar. QoS 2 (exactly once): intercanvi de quatre missatges; car i rarament necessari. Quilòmetre Zero publica posicions amb QoS 1 i no amb 0 per una raó que es veu a l'apartat 7: en reconnectar després d'un túnel, vol que l'última posició coneguda arribi segur, encara que es dupliqui (el consumidor deduplica per seqüència). - Retained messages. En publicar amb
retain=True, el broker desa l'últim missatge del tòpic i el lliura immediatament a qui s'hi subscrigui després. És el que permet que un panell acabat d'obrir vegi l'última posició de cada furgoneta sense esperar 5 s. - Last will (testament). En connectar, el client registra un missatge que el broker publicarà per ell si la connexió es perd sense un
DISCONNECTnet:km0/repartiment/furgoneta-3/estat=desconnectada. És la detecció de fallada de 07-03 feta pel broker. - Sessions persistents. Amb
clean_session=False(MQTT 3.1.1) osession expiry(MQTT 5), el broker recorda les subscripcions del client i encua els missatges QoS 1/2 que li arribin mentre està desconnectat, i els hi lliura en tornar. Per a les ordres que l'operador envia a la furgoneta ("nova parada afegida") és imprescindible: el túnel no les pot perdre. - MQTT sobre WebSockets. Un navegador no obre sockets TCP arbitraris, però sí WebSockets; els brokers exposen un port WebSocket (9001 a Mosquitto) pel qual parlen MQTT dins de trames WebSocket. És la via per la qual el panell d'operadors podria subscriure's directament al broker; en el disseny de Quilòmetre Zero no es fa (apartat 5), però és útil per a eines de diagnòstic.
| Decisió | Posició de la furgoneta | Estat (connectada/desconnectada) | Ordres de l'operador a la furgoneta |
|---|---|---|---|
| Tòpic | km0/repartiment/<id>/posicio |
km0/repartiment/<id>/estat |
km0/repartiment/<id>/ordres |
| QoS | 1 | 1 | 1 |
| Retained | Sí (última posició) | Sí | No |
| Sessió persistent | No cal (només publica) | — | Sí (el dispositiu les ha de rebre en tornar) |
| Last will | — | desconnectada |
— |
- Arquitectura pub/sub d'extrem a extrem
Amb les peces anteriors, el recorregut d'una posició des de la furgoneta fins al mapa de l'Anna és aquest:
flowchart LR
subgraph Carrer[Furgoneta, xarxa mòbil]
APP[App repartidor<br/>furgoneta_mqtt.py<br/>QoS 1, retained]
end
APP -- "MQTT/TLS 8883<br/>km0/repartiment/furgoneta-3/posicio" --> BR[Broker MQTT<br/>Mosquitto / EMQX<br/>ACL per repartidor]
BR -- "subscripció km0/repartiment/+/posicio" --> PU[pont_mqtt_kafka.py]
PU -- "clau = repartidor" --> K[(Kafka<br/>repartiment.posicions)]
K --> FL[Flink alertes i panell<br/>05-04]
FL --> KP[(Kafka<br/>repartiment.panell)]
K --> CS[(Cassandra<br/>posicions, 04-04)]
K --> WS1[repartiment ws_servidor<br/>instància 1]
KP --> WS1
K --> WS2[repartiment ws_servidor<br/>instància 2]
KP --> WS2
WS1 <--> RP[(Redis pub/sub<br/>bus entre instàncies)]
WS2 <--> RP
WS1 -- "wss:// via Kong" --> ANA["Navegador de l'Anna<br/>subscrita a P-2026-000124"]
WS2 -- "wss:// via Kong" --> JOR[Panell del Jordi<br/>subscrit al mercat girona]
Cada salt té una raó:
- La furgoneta parla MQTT, no HTTP ni Kafka. Una connexió persistent barata, amb reconnexió i QoS gestionats per la biblioteca, i un broker que l'autentica i la limita al seu prefix de tòpics.
- El pont MQTT → Kafka existeix perquè la resta de la plataforma consumeix Kafka: Flink (05-04), Cassandra per a l'històric,
analitica. Tradueix el tòpic jeràrquic a un tòpic Kafka amb clau = id de repartidor, cosa que conserva l'ordre per furgoneta (02-04) i permet a Flink la finestra per clau. Alguns brokers (EMQX) porten aquest pont integrat; amb Mosquitto se n'escriu un de petit. repartimentconsumeix Kafka i fa fan-out per WebSockets. És l'únic component que coneix els clients: sap que l'Anna està subscrita a la comandaP-2026-000124, que aquesta comanda la portafurgoneta-3, i que per tant cada posició defurgoneta-3ha d'arribar a la connexió de l'Anna. Aquest coneixement (la taula de subscripcions) és el que cap altra peça no té.- Redis pub/sub entre instàncies resol el problema de l'apartat següent.
Fixeu-vos en el que no es fa: el navegador no es connecta al broker MQTT (encara que podria, sobre WebSockets) perquè això exigiria donar a cada client credencials MQTT, gestionar ACL per comanda al broker i renunciar a l'enriquiment que fa repartiment (l'Anna no ha de rebre la posició crua de furgoneta-3 amb les seves altres dotze comandes, sinó "la teva comanda és a 4 parades i 12 minuts").
- Fan-out amb moltes instàncies, presència i subscripcions
El problema
repartiment corre amb diverses rèpliques a Kubernetes (07-05). L'Anna està connectada a la instància 1; en Marc, que espera un altre lliurament de la mateixa furgoneta-3, a la instància 2. Quan arriba la posició per Kafka, qui la rep? Si les dues instàncies formen un grup de consumidors de Kafka, cada partició la llegeix una d'elles (02-04): la posició arriba a la instància 1, que l'envia a l'Anna, i en Marc no veu res. Si cada instància consumeix tot el tòpic amb un grup diferent, cadascuna rep totes les posicions i les reparteix als seus clients: funciona, però cada instància processa les 28 posicions per segon completes i consumeix el tòpic sencer, cosa que a desenes d'instàncies i amb repartiment.panell inclòs comença a pesar, i no resol missatges que neixen en una instància (el xat: l'Anna escriu a la instància 1 i la formatgeria és a la 2).
Les dues solucions
Sessions enganxoses (sticky sessions): el balancejador (Kong, o l'Ingress) envia sempre el mateix client a la mateixa instància (per cookie o per hash d'IP), i s'encamina cada missatge a la instància correcta. Exigeix que algú sàpiga en quina instància és cada client (un mapa client → instància a Redis) i que les instàncies es parlin. És fràgil: en escalar o en caure una instància el mapa s'invalida, i el hash per IP falla amb el NAT dels operadors mòbils (milers de clients darrere de la mateixa IP).
Bus entre instàncies (l'opció de Quilòmetre Zero): cada instància se subscriu a un canal de Redis pub/sub (o a un tòpic Kafka amb un grup per instància) per cada tema en què té almenys un client interessat. Qui tingui alguna cosa per difondre (el consumidor de Kafka de qualsevol instància, o el gestor del xat) la publica al bus, i cada instància amb subscriptors la rep i la reparteix a les seves connexions. No cal saber on és ningú; les instàncies són intercanviables i s'escalen sense coordinació.
flowchart TB
K[(Kafka repartiment.posicions)] --> C1[Consumidor Kafka<br/>grup repartiment-ws<br/>a la instància que toqui]
C1 -- "PUBLISH furgoneta:furgoneta-3" --> R[(Redis pub/sub)]
R -- "SUBSCRIBE furgoneta:furgoneta-3" --> I1[Instància 1<br/>subscriptors locals:<br/>Anna]
R -- "SUBSCRIBE furgoneta:furgoneta-3" --> I2[Instància 2<br/>subscriptors locals:<br/>Marc, Jordi]
R -. "ningú subscrit: no se subscriu" .- I3[Instància 3<br/>sense interessats]
I1 --> A[Anna]
I2 --> M[Marc]
I2 --> J[Jordi]
Redis pub/sub és fire-and-forget: si una instància està desconnectada del bus en aquell instant, el missatge es perd, i no hi ha històric. Per a les posicions és igual (la següent arriba en 5 s i el client pot demanar l'última coneguda en connectar). Per al xat no: els missatges es persisteixen primer (a Cassandra o PostgreSQL) i el bus només notifica; en reconnectar, el client demana "tot des de l'id X" a l'API. Quan el volum o la garantia exigeixen més, el bus passa a ser Kafka amb un grup de consumidors per instància, o Redis Streams.
Presència i subscripcions
La taula de subscripcions viu en memòria a cada instància ({tema: {connexions}}): és local, ràpida i es perd amb la instància, que és el correcte perquè les connexions també es perden. La presència (qui està connectat ara, útil per a "la Marta està en línia" al xat o perquè l'operador sàpiga quins clients estan mirant) es desa a Redis amb TTL (presencia:u-anna → instancia-1, renovat amb cada ping): si la instància mor, la clau caduca sola. I la resolució d'a quina furgoneta correspon la comanda de l'Anna es fa en el moment de subscriure's, consultant el model de repartiment, i es torna a fer si arriba un esdeveniment de reassignació.
- Autenticació, ordre, deduplicació, backpressure i escalat
Autenticació del WebSocket amb el JWT
La connexió WebSocket dura minuts; el JWT de 06-01 caduca als 15. Tres decisions:
- El JWT es verifica al handshake (Kong ho fa amb el plugin
jwtcom amb qualsevol ruta;repartimentel torna a verificar, defensa en profunditat). Es passa com a paràmetre de consulta (wss://.../ws?token=...) o, millor, com a primer missatge després de connectar, perquè les URL acaben als logs (07-02) i el token no hi ha d'acabar. - Autorització per subscripció: cada missatge
{"accio": "subscriure", "comanda": "P-2026-000124"}s'autoritza contra els claims (sub=u-annaha de ser el client de la comanda; unroles: ["operador"]pot subscriure's a un mercat). No hi ha autorització per cada missatge enviat pel servidor: es va fer en subscriure. - Caducitat: quan el token caduca, el servidor no talla la connexió de cop (el mapa es congelaria a mig trajecte); envia
{"tipus": "renovar"}i el client envia un{"accio": "token", "jwt": "..."}nou. Si no ho fa en 60 s, es tanca amb codi4401.
Ordre i deduplicació al client
QoS 1 pot duplicar; Kafka amb clau ordena per partició però un reintent del pont pot repetir; la reconnexió del client pot fer que rebi l'"última posició" retinguda i alhora la posició nova. La solució és la de 01-05: cada posició porta un número de seqüència per furgoneta (seq, un comptador monòton que l'app incrementa a cada enviament i que sobreviu a reinicis perquè es persisteix localment) a més de data_ms. El client desa ultima_seq per furgoneta i descarta tot el que no sigui estrictament més gran. No cal ordenar: una posició vella no aporta res, es llença. El servidor fa el mateix abans de difondre, per no gastar amplada de banda en duplicats.
Backpressure cap a clients lents
Un navegador en 3G no consumeix 28 missatges per segon (el panell del Jordi els rebria tots). Si el servidor escriu sense control, el búfer d'enviament d'aquesta connexió creix sense límit i acaba amb la memòria de la instància, que s'endú per davant tots els altres clients: el mateix problema de 05-04, a l'últim tram. Tres defenses, de menor a major agressivitat: coalescència (per a cada furgoneta només importa l'última posició: si el client en té tres de pendents, se n'envia una), cua acotada per connexió (per exemple 100 missatges; si s'omple, es descarten els més antics i es compta en una mètrica km0_ws_descarts_total), i tancament (si la cua està plena durant més de 30 s, el client no pot seguir el ritme: es tanca amb codi 1013 Try again later i el client reconnecta amb backoff, potser demanant menys freqüència).
Escalat: connexions per instància i límits del sistema operatiu
Una connexió WebSocket inactiva costa poc: un descriptor de fitxer, uns 20-50 KB entre kernel i aplicació amb asyncio. Una instància de repartiment amb 2 GB pot sostenir de l'ordre de 20 000 a 50 000 connexions, sempre que el sistema operatiu ho permeti: ulimit -n (descriptors per procés, per defecte 1 024; al contenidor es puja a 65 536 o més), net.core.somaxconn i net.ipv4.ip_local_port_range al node, i el balancejador del davant (Kong, l'Ingress), que també manté una connexió per client i té els seus propis límits. El que sí que costa és el trànsit: 50 000 clients rebent 1 missatge/s són 50 000 escriptures per segon per instància; aquí el límit és la CPU de serialitzar i l'amplada de banda, i s'escala afegint instàncies (l'HPA de 07-05 amb la mètrica km0_ws_connexions en lloc de CPU). Els desplegaments són el moment delicat: un rolling update talla totes les connexions de la instància que es retira, així que l'apagada ordenada (07-05) envia un tancament 1001 Going away esglaonat durant 30 s i confia en la reconnexió amb jitter del client.
- Codi: del broker MQTT al navegador
Broker Mosquitto amb ACL per repartidor
# km0/docker-compose.yml (fragment afegit en aquesta lliçó)
services:
mosquitto:
image: eclipse-mosquitto:2
ports:
- "8883:8883" # MQTT sobre TLS per a les furgonetes
- "9001:9001" # MQTT sobre WebSockets (diagnòstic)
volumes:
- ./vora/mosquitto/mosquitto.conf:/mosquitto/config/mosquitto.conf:ro
- ./vora/mosquitto/acl:/mosquitto/config/acl:ro
- ./vora/mosquitto/passwd:/mosquitto/config/passwd:ro
- ./certs:/mosquitto/certs:ro # emesos per km0-ca (06-04)
- mosquitto-data:/mosquitto/data # sessions persistents i retained
volumes:
mosquitto-data:# km0/vora/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/vora/mosquitto/acl # Cada repartidor només publica sota el seu propi prefix i només llegeix les seves ordres. # %u se substitueix pel nom d'usuari amb què s'ha autenticat. pattern write km0/repartiment/%u/posicio pattern write km0/repartiment/%u/estat pattern read km0/repartiment/%u/ordres # El pont llegeix totes les posicions i estats, i escriu ordres. user pont topic read km0/repartiment/+/posicio topic read km0/repartiment/+/estat topic write km0/repartiment/+/ordres
El fitxer passwd es genera amb mosquitto_passwd -c passwd furgoneta-3 (una contrasenya per repartidor, emesa en donar d'alta el dispositiu i desada a Vault, 06-04). Amb l'ACL, una app compromesa de furgoneta-3 no pot publicar posicions falses de furgoneta-7 ni llegir les seves ordres: el broker ho rebutja abans que arribi a la plataforma.
L'app de la furgoneta: furgoneta_mqtt.py
# km0/serveis/repartiment/furgoneta_mqtt.py
"""Client MQTT de l'app del repartidor. Publica posicions amb QoS 1 i retained,
registra un last will i persisteix el número de seqüència entre reinicis."""
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/repartiment/{REPARTIDOR}/posicio"
T_EST = f"km0/repartiment/{REPARTIDOR}/estat"
T_CMD = f"km0/repartiment/{REPARTIDOR}/ordres"
FITXER_SEQ = f"/var/lib/km0/{REPARTIDOR}.seq" # el comptador sobreviu a reinicis (01-05)
def llegir_seq() -> int:
try:
return int(open(FITXER_SEQ).read())
except FileNotFoundError:
return 0
def desar_seq(seq: int) -> None:
with open(FITXER_SEQ + ".tmp", "w") as f:
f.write(str(seq))
os.replace(FITXER_SEQ + ".tmp", FITXER_SEQ) # escriptura atòmica
def en_connectar(client, userdata, flags, rc, properties=None):
print("connectat, sessió prèvia:", flags.get("session present", flags))
client.subscribe(T_CMD, qos=1) # amb clean_session=False, el broker la recorda
client.publish(T_EST, "connectada", qos=1, retain=True)
def en_missatge(client, userdata, msg):
ordre = json.loads(msg.payload)
print("ordre rebuda:", ordre["tipus"]) # p. ex. {"tipus": "nova_parada", "comanda": "P-2026-000126"}
def crear_client() -> 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, "desconnectada", qos=1, retain=True) # last will: el publica el broker si desapareixem
c.on_connect = en_connectar
c.on_message = en_missatge
c.reconnect_delay_set(min_delay=1, max_delay=60) # backoff de reconnexió, el gestiona paho
return c
def bucle_posicions(client: mqtt.Client, gps) -> None:
seq = llegir_seq()
while True:
lat, lon = gps.llegir()
seq += 1
desar_seq(seq)
carrega = json.dumps({"repartidor": REPARTIDOR, "seq": seq, "lat": lat, "lon": lon,
"data_ms": int(time.time() * 1000)})
# QoS 1: paho reenvia si no hi ha PUBACK; si som en un túnel, encua i publica en tornar.
# retain=True: qui es subscrigui després rep l'última posició sense esperar 5 s.
client.publish(T_POS, carrega, qos=1, retain=True)
time.sleep(5)
if __name__ == "__main__":
client = crear_client()
client.connect_async(BROKER, 8883, keepalive=30) # keepalive: PINGREQ cada 30 s; el broker detecta la caiguda en 45 s
client.loop_start() # fil de xarxa: gestiona reconnexions i reenviaments
from serveis.repartiment.gps import GPS # abstracció del receptor del dispositiu
bucle_posicions(client, GPS())Detalls que importen: clean_session=False amb un client_id fix fa que el broker desi la subscripció a ordres i encui les ordres QoS 1 durant el túnel; keepalive=30 és el que dona vida al last will (sense PINGREQ durant 1,5 × keepalive, el broker publica desconnectada); i el comptador seq es persisteix abans de publicar, perquè un reinici de l'app no reutilitzi un número (cosa que faria que el client descartés posicions noves com si fossin antigues). Que la biblioteca encui mentre no hi ha xarxa significa que, després d'un túnel de 3 minuts, arribaran de cop 36 posicions antigues: repartiment les processarà per ordre de seq i a l'Anna només li arribarà l'última, gràcies a la coalescència.
El pont MQTT → Kafka: pont_mqtt_kafka.py
# km0/serveis/repartiment/pont_mqtt_kafka.py
"""Se subscriu a km0/repartiment/+/posicio i publica cada posició a Kafka repartiment.posicions
amb clau = repartidor. És un procés sense estat: se'n poden córrer diversos (el broker
reparteix amb subscripcions compartides a MQTT 5: $share/pont/km0/repartiment/+/posicio)."""
import json, ssl, os
import paho.mqtt.client as mqtt
from confluent_kafka import Producer
from serveis.comu import metriques, logs
log = logs.obtenir("repartiment.pont")
productor = Producer({"bootstrap.servers": os.environ["KAFKA_BOOTSTRAP"],
"enable.idempotence": True, "acks": "all"}) # 02-05: sense duplicats per reintent
publicades = metriques.comptador("km0_pont_posicions_total")
def lliurat(err, msg):
if err is not None:
log.error("fallada en publicar a Kafka", error=str(err), clau=msg.key())
def en_missatge(client, userdata, msg):
parts = msg.topic.split("/") # ["km0", "repartiment", "furgoneta-3", "posicio"]
repartidor = parts[2]
pos = json.loads(msg.payload)
if pos.get("repartidor") != repartidor: # l'ACL ja ho impedeix, però no ens refiem del payload
log.warning("payload amb repartidor diferent del tòpic", topic=msg.topic)
return
embolcall = {"id_esdeveniment": f"{repartidor}-{pos['seq']}", # determinista: mateix esdeveniment, mateix id (02-05)
"tipus": "posicio.actualitzada", "versio": 1,
"data_ms": pos["data_ms"], "origen": "pont-mqtt", "dades": pos}
productor.produce("repartiment.posicions", key=repartidor.encode(),
value=json.dumps(embolcall).encode(), callback=lliurat)
productor.poll(0)
publicades.inc()
client = mqtt.Client(client_id=f"pont-{os.getpid()}", protocol=mqtt.MQTTv5)
client.username_pw_set("pont", os.environ["KM0_MQTT_PASS_PONT"])
client.tls_set(ca_certs="/etc/km0/km0-ca.crt", tls_version=ssl.PROTOCOL_TLS_CLIENT)
client.on_message = en_missatge
client.connect(os.environ.get("KM0_MQTT_HOST", "mosquitto"), 8883)
client.subscribe("$share/pont/km0/repartiment/+/posicio", qos=1) # subscripció compartida: diverses rèpliques
client.loop_forever()L'id_esdeveniment es construeix de manera determinista amb repartidor-seq en lloc d'un UUID aleatori: així, si el broker reenvia un PUBLISH (QoS 1) o el pont es reinicia i el reprocessa, l'esdeveniment a Kafka porta el mateix id i els consumidors idempotents de 02-05 el descarten. És un detall petit que elimina tota una classe de duplicats.
El servidor WebSocket: ws_servidor.py
# km0/serveis/repartiment/ws_servidor.py
"""Servidor WebSocket de repartiment amb FastAPI: autenticació per JWT, subscripció per
comanda o per mercat, fan-out entre instàncies amb Redis pub/sub, ping/pong,
deduplicació per seq i cua acotada per connexió."""
import asyncio, json, os
from collections import defaultdict
from fastapi import FastAPI, WebSocket, WebSocketDisconnect
import redis.asyncio as redis
from serveis.comu import metriques, logs
from serveis.comu.jwt import verificar_jwt, JwtInvalid # 06-01
from serveis.repartiment.assignacions import furgoneta_de_comanda # model de repartiment
app = FastAPI()
log = logs.obtenir("repartiment.ws")
r = redis.from_url(os.environ["REDIS_URL"])
connexions_g = metriques.gauge("km0_ws_connexions")
descarts = metriques.comptador("km0_ws_descarts_total")
CUA_MAX = 100
class Connexio:
def __init__(self, ws: WebSocket, claims: dict):
self.ws, self.claims = ws, claims
self.cua: asyncio.Queue = asyncio.Queue(maxsize=CUA_MAX)
self.ultima_seq: dict[str, int] = {} # per furgoneta: deduplicació (01-05)
self.temes: set[str] = set()
def encuar(self, tema: str, missatge: dict) -> None:
furgo, seq = missatge.get("repartidor"), missatge.get("seq")
if furgo and seq is not None:
if seq <= self.ultima_seq.get(furgo, -1):
return # duplicat o antic: es descarta
self.ultima_seq[furgo] = seq
if self.cua.full(): # backpressure: client lent
self.cua.get_nowait() # llençar el més antic
descarts.inc()
self.cua.put_nowait(missatge)
# Taula de subscripcions LOCAL a aquesta instància: tema -> connexions
subscriptors: dict[str, set[Connexio]] = defaultdict(set)
pubsub = r.pubsub()
async def bucle_bus() -> None:
"""Rep del bus Redis el que ha publicat qualsevol instància i ho reparteix localment."""
async for m in pubsub.listen():
if m["type"] not in ("message", "pmessage"):
continue
tema = m["channel"].decode()
missatge = json.loads(m["data"])
for c in list(subscriptors.get(tema, ())):
c.encuar(tema, missatge)
@app.on_event("startup")
async def arrencar():
asyncio.create_task(bucle_bus())
async def subscriure(c: Connexio, tema: str) -> None:
if not subscriptors[tema]:
await pubsub.subscribe(tema) # primera connexió local interessada: subscriure al bus
subscriptors[tema].add(c)
c.temes.add(tema)
async def dessubscriure_tot(c: Connexio) -> None:
for tema in c.temes:
subscriptors[tema].discard(c)
if not subscriptors[tema]:
await pubsub.unsubscribe(tema) # ningú més aquí: deixar de rebre del bus
del subscriptors[tema]
def autoritzar(claims: dict, accio: dict) -> str | None:
"""Retorna el tema del bus al qual es tradueix la subscripció, o None si no està permès."""
if "comanda" in accio:
comanda = accio["comanda"]
assignacio = furgoneta_de_comanda(comanda) # {"client": "u-anna", "furgoneta": "furgoneta-3"}
if assignacio and assignacio["client"] == claims["sub"]:
return f"furgoneta:{assignacio['furgoneta']}"
if "mercat" in accio and "operador" in claims.get("roles", []):
return f"mercat:{accio['mercat']}"
return None
async def emissor(c: Connexio) -> None:
"""Tasca que buida la cua cap al socket. Separada de la recepció per no bloquejar-la."""
while True:
missatge = await c.cua.get()
await c.ws.send_text(json.dumps(missatge))
@app.websocket("/ws")
async def ws_endpoint(ws: WebSocket):
await ws.accept()
# 1. Primer missatge: el token (no a l'URL, perquè no acabi als logs de Kong).
try:
primer = json.loads(await asyncio.wait_for(ws.receive_text(), timeout=5))
claims = verificar_jwt(primer["jwt"], audiencia="repartiment")
except (asyncio.TimeoutError, KeyError, JwtInvalid):
await ws.close(code=4401); return
c = Connexio(ws, claims)
connexions_g.inc()
tasca_emissor = asyncio.create_task(emissor(c))
log.info("ws connectat", sub=claims["sub"])
try:
while True:
accio = json.loads(await ws.receive_text())
if accio.get("accio") == "subscriure":
tema = autoritzar(claims, accio)
if tema is None:
await ws.send_text(json.dumps({"tipus": "error", "codi": "no_autoritzat"})); continue
await subscriure(c, tema)
ultima = await r.get(f"ultima:{tema}") # última posició coneguda (com el retained d'MQTT)
if ultima:
c.encuar(tema, json.loads(ultima))
elif accio.get("accio") == "token":
claims = verificar_jwt(accio["jwt"], audiencia="repartiment") # renovació sense tallar
c.claims = claims
except WebSocketDisconnect:
pass
finally:
tasca_emissor.cancel()
await dessubscriure_tot(c)
connexions_g.dec()
log.info("ws tancat", sub=claims["sub"])I el consumidor de Kafka que alimenta el bus (una tasca a cada instància, totes al mateix grup de consumidors, de manera que cada partició la llegeix una):
# km0/serveis/repartiment/ws_consumidor_kafka.py
"""Llegeix repartiment.posicions i repartiment.panell (grup repartiment-ws) i publica a 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": "repartiment-ws",
"auto.offset.reset": "latest"}) # a ningú no li interessen posicions velles en arrencar
c.subscribe(["repartiment.posicions", "repartiment.panell"])
while True:
msg = c.poll(1.0)
if msg is None or msg.error():
continue
ev = json.loads(msg.value())
d = ev["dades"]
if msg.topic() == "repartiment.posicions":
carrega = json.dumps({"tipus": "posicio", **d})
r.publish(f"furgoneta:{d['repartidor']}", carrega)
r.publish(f"mercat:{d.get('mercat', 'desconegut')}", carrega)
r.set(f"ultima:furgoneta:{d['repartidor']}", carrega, ex=600) # equivalent al retained
else:
r.publish(f"mercat:{d['mercat']}", json.dumps({"tipus": "panell", **d}))
c.commit(msg)El client del navegador, amb reconnexió i backoff
// km0/vora/web/seguiment.js — client WebSocket mínim per a la pantalla de seguiment de l'Anna
class Seguiment {
constructor(comanda, obtenirJwt, enPosicio) {
this.comanda = comanda; this.obtenirJwt = obtenirJwt; this.enPosicio = enPosicio;
this.intent = 0; this.ultimaSeq = -1; this.tancat = false;
this.connectar();
}
connectar() {
this.ws = new WebSocket("wss://api.km0.example/ws"); // passa per Kong (Upgrade permès)
this.ws.onopen = () => {
this.intent = 0; // connexió aconseguida: reiniciar backoff
this.ws.send(JSON.stringify({ jwt: this.obtenirJwt() })); // 1r: autenticar
this.ws.send(JSON.stringify({ accio: "subscriure", comanda: this.comanda }));
};
this.ws.onmessage = (e) => {
const m = JSON.parse(e.data);
if (m.tipus === "renovar") { this.ws.send(JSON.stringify({ accio: "token", jwt: this.obtenirJwt() })); return; }
if (m.tipus === "posicio") {
if (m.seq <= this.ultimaSeq) return; // duplicat o antic (01-05)
this.ultimaSeq = m.seq;
this.enPosicio(m.lat, m.lon, m.data_ms);
}
};
this.ws.onclose = (e) => {
if (this.tancat || e.code === 4401) return; // tancament voluntari o no autoritzat: no reintentar
const base = Math.min(30000, 1000 * 2 ** this.intent++); // exponencial amb sostre de 30 s
const espera = base / 2 + Math.random() * base / 2; // jitter: evitar l'estampida (07-04)
setTimeout(() => this.connectar(), espera);
};
this.ws.onerror = () => this.ws.close();
}
tancar() { this.tancat = true; this.ws.close(1000); }
}
// Ús: new Seguiment("P-2026-000124", () => sessio.jwt, (lat, lon) => mapa.moure(lat, lon));El navegador no implementa ping/pong visible des de JavaScript (ho fa per sota quan el servidor envia ping), així que la detecció d'una connexió "zombi" des del client es recolza en el servidor: si no arriba cap missatge en 60 s, el client pot tancar i reconnectar. ultimaSeq al client és l'última línia de defensa contra duplicats: encara que el servidor ja filtri, en reconnectar a una altra instància es rep l'"última coneguda", que pot ser una de ja vista.
- SSE com a alternativa per a les alertes d'estoc
Les alertes d'inventari.alertes a productors són el cas oposat al mapa: pocs missatges, només servidor → client, i un panell que pot estar obert hores. WebSockets seria sobredimensionar; SSE encaixa exactament, amb reconnexió gratis i Last-Event-ID per no perdre alertes durant una desconnexió.
# km0/serveis/inventari/sse_alertes.py
from fastapi import FastAPI, Request, Depends
from sse_starlette.sse import EventSourceResponse
import asyncio, json, redis.asyncio as redis
from serveis.comu.jwt import claims_de_peticio # 06-01: Bearer a la capçalera, com qualsevol GET
app = FastAPI()
r = redis.from_url("redis://redis:6379")
@app.get("/api/v1/inventari/alertes/stream")
async def stream(request: Request, claims=Depends(claims_de_peticio)):
productor = claims["productor_id"] # p. ex. "formatgeria-montblanc"
ultim = request.headers.get("Last-Event-ID", "0-0") # el navegador l'envia en reconnectar
async def generador():
cursor = ultim
while not await request.is_disconnected():
# Redis Streams (no pub/sub): té històric, així que la reconnexió no perd alertes.
res = await r.xread({f"alertes:{productor}": cursor}, block=15000, count=10)
if not res:
yield {"comment": "keepalive"} # evita que els proxies tanquin per inactivitat
continue
for _, entrades in res:
for id_, camps in entrades:
cursor = id_
yield {"id": id_, "event": "estoc_baix", "data": camps[b"json"].decode()}
return EventSourceResponse(generador())// Panell del productor: el navegador reconnecta sol i reenvia Last-Event-ID
const es = new EventSource("/api/v1/inventari/alertes/stream"); // el JWT va en cookie o via Kong
es.addEventListener("estoc_baix", (e) => { const a = JSON.parse(e.data); mostrarAvis(a.producte, a.mercat, a.disponible); });Un consumidor d'inventari.alertes (el que a 05-04 escrivia les alertes de Flink) afegeix cada alerta al stream alertes:<productor> amb XADD i un MAXLEN d'uns centenars d'entrades. La diferència amb el mapa és que aquí sí que importa no perdre'n cap: per això Redis Streams (amb històric i ids monòtons) i no pub/sub.
- Quan usar cada tècnica
| Necessitat | Tècnica | Per què |
|---|---|---|
| Dada que canvia cada diversos minuts, pocs clients | Polling amb Cache-Control |
Simplicitat; el gateway ho desa en memòria cau |
| Servidor → client, text, no perdre missatges, sense volum | SSE + Redis Streams | Reconnexió i represa de sèrie; HTTP normal |
| Bidireccional o alt volum cap al navegador | WebSockets + bus entre instàncies | Trames barates, un canal per a tot (mapa, xat, ordres) |
| Dispositius en xarxes dolentes, amb bateria | MQTT (QoS 1, sessions, last will) | Dissenyat per a això; el broker aïlla i autentica |
| Navegador que ha de parlar amb un broker MQTT | MQTT sobre WebSockets | Únic transport disponible al navegador |
| Missatges que han de sobreviure a la desconnexió del client | Persistir primer (Streams, base de dades); el canal només notifica | Ni WebSocket ni Redis pub/sub no desen res |
| Respatller quan WebSocket està bloquejat (proxies corporatius) | Long polling | Funciona on no funciona res més |
Errors Comuns i Consells
- Exposar Kafka o el broker MQTT al navegador. Ni protocol, ni seguretat, ni nombre de connexions. Sempre un servidor d'últim tram que autentica, autoritza per subscripció i enriqueix.
- Passar el JWT a l'URL del WebSocket. Acaba als logs de Kong i del servidor (07-02). Primer missatge després de connectar, i renovació sense tallar.
- Un grup de consumidors de Kafka per instància de WebSocket sense bus. O cada instància rep només una part i els seus clients es queden sense dades, o cadascuna ho consumeix tot i no escala. Bus entre instàncies (Redis pub/sub o Kafka).
- Escriure al socket sense cua acotada. Un client en 3G esgota la memòria de la instància i tomba tothom. Cua per connexió, coalescència, mètrica de descarts i tancament
1013. - Reconnexió sense backoff ni jitter. Un desplegament de
repartimentprovoca que milers de clients reconnectin en el mateix segon. Exponencial amb sostre i jitter, i apagada ordenada esglaonada al servidor. - QoS 2 "per seguretat". Quatre missatges per posició, sessió pesada al broker, i de tota manera el consumidor ha de deduplicar per altres raons. QoS 1 i
seq. - Oblidar el last will i el keepalive. Sense ells, un repartidor sense cobertura apareix "connectat" durant els minuts que TCP triga a assabentar-se'n.
keepalive=30i testament aestat. - Consell: mesura
km0_ws_connexions,km0_ws_descarts_total, el lag del gruprepartiment-wsi la latència d'extrem a extrem (data_msde la furgoneta davant de l'hora de recepció al navegador, amb el rellotge corregit): és l'SLO real del "temps real". - Consell: en el disseny, comença per la taula de l'apartat 1. Latència acceptable, volum i sentit decideixen la tècnica abans que qualsevol preferència pels WebSockets.
Exercicis
Exercici 1: el xat client-productor
Dissenya el xat entre l'Anna i la Formatgeria Montblanc sobre la infraestructura d'aquesta lliçó: (a) quin transport usa cada extrem i per què; (b) què es persisteix, on i en quin ordre respecte de la notificació pel bus; (c) què fa el client en reconnectar després de 2 minuts sense xarxa per no perdre ni duplicar missatges; (d) quin tema del bus i quina regla d'autorització aplica ws_servidor.py.
Exercici 2: el túnel de tres minuts
furgoneta-3 travessa un túnel de 3 minuts. Descriu, pas a pas i anomenant els mecanismes, què passa a: l'app (furgoneta_mqtt.py), el broker (last will, retained, sessió), el pont, repartiment.posicions, Flink (la finestra de sessió de 05-04), ws_servidor.py i la pantalla de l'Anna. Quants missatges rep l'Anna en sortir del túnel i per què?
Exercici 3: dimensionar el panell d'operadors
En campanya hi ha 140 furgonetes enviant cada 5 s i 40 operadors amb el panell obert, cadascun subscrit al mercat complet (35 furgonetes de mitjana). (a) Quants missatges per segon rep cada operador i quants n'escriu en total el conjunt d'instàncies de ws_servidor? (b) Si un operador és en una connexió que només admet 2 missatges/s, què passa amb la cua de 100 i amb quin mecanisme es resol sense tancar la connexió? (c) Proposa un canvi al consumidor de Kafka o al servidor que redueixi el trànsit al panell a 1 missatge/s per operador mantenint la informació útil.
Solucions
Exercici 1.
(a) Tots dos extrems són navegadors (l'Anna al mòbil, la formatgeria al seu panell): WebSockets en tots dos, perquè el xat és bidireccional i de baixa latència, i perquè l'Anna ja té oberta la connexió del seguiment: es reutilitza el mateix canal amb un altre tipus de missatge ({"accio": "xat", "conversa": "P-2026-000124", "text": "..."}).
(b) Cada missatge es persisteix primer a comandes o en un context missatges (taula missatges_per_conversa a Cassandra amb clau de partició = conversa i clustering per un id monòton, per exemple un TimeUUID o un seq per conversa assignat pel servidor), i només després es publica al bus conversa:P-2026-000124. Si es publiqués abans de persistir i el procés morís entremig, l'altre extrem veuria un missatge que no existeix.
(c) En reconnectar, el client envia {"accio": "subscriure", "conversa": "...", "des_de": <últim id vist>}; el servidor respon amb els missatges persistits posteriors (una consulta per clau de partició, barata) i després subscriu al bus. El client deduplica per id de missatge (en pot rebre un tant per la recuperació com pel bus si ha arribat just en aquell instant). És el mateix Last-Event-ID d'SSE, fet a mà.
(d) Tema conversa:<comanda>; autorització: claims["sub"] és el client de la comanda, o claims["productor_id"] és el productor d'alguna línia de la comanda (consulta al model de comandes en el moment de subscriure, no per missatge).
Exercici 2.
- App: perd el
PINGREQ; paho detecta la caiguda i entra en reconnexió amb backoff (1 s → 60 s). El bucle de posicions continua: incrementaseq, persisteix i cridapublishamb QoS 1; paho encua en memòria (~36 missatges en 3 minuts). - Broker: als 45 s sense keepalive (1,5 × 30) dona la connexió per morta i publica el last will
km0/repartiment/furgoneta-3/estat = desconnectada(retained). La sessió persistent conserva la subscripció aordresi encua qualsevol ordre QoS 1 que l'operador enviï. - Pont: rep l'estat
desconnectadai el publica a Kafka (com a esdevenimentestat.actualitzat); no rep posicions. repartiment.posicions: cap esdeveniment defurgoneta-3durant 3 minuts.- Flink: la finestra de sessió amb gap de 3 minuts es tanca (segons la marca d'aigua) i emet "repartidor sense senyal" a
repartiment.panell; el panell del Jordi ho mostra. ws_servidor: no envia res a l'Anna; la seva connexió continua viva (ping/pong amb el servidor, que sí que té xarxa). El mapa mostra l'última posició i, si el client ho implementa, "últim senyal fa 2 min".- Sortida del túnel: paho reconnecta (la sessió hi era present), rep les ordres encuades i buida la seva cua: 36
PUBLISHQoS 1 en ràfega, en ordre deseq. El broker actualitza el retained amb l'última. El pont publica els 36 a Kafka ambid_esdevenimentdeterminista. Flink obre una sessió nova ("ha tornat"). El consumidor dews_servidorpublica 36 missatges al bus; a laConnexiode l'Anna,encuaraccepta cadascun perquèseqés creixent, però com que arriben en mil·lisegons i la cua té 100 de capacitat, s'encuen tots: l'Anna rebria 36 missatges en ràfega, i el mapa "saltaria" pel túnel. Perquè només rebi l'última caldria la coalescència per furgoneta (mantenir a la cua només l'última posició de cada furgoneta), que el codi de l'apartat 8 no implementa i que l'exercici 3(c) introdueix. Sense ella, no hi ha error, només trànsit inútil.
Exercici 3.
(a) 140 furgonetes / 5 s = 28 posicions/s en total. Cada operador, amb 35 furgonetes: 7 missatges/s. Escriptures totals: 40 × 7 = 280 missatges/s entre totes les instàncies (més les dels clients amb seguiment). És poc: una sola instància ho sosté; el problema no és el volum sinó la robustesa.
(b) N'entren 7/s i en surten 2/s: la cua creix 5/s i s'omple en 20 s. A partir d'aquí, encuar descarta el més antic per cada nou (km0_ws_descarts_total puja 5/s) i l'operador veu posicions amb fins a 100/7 ≈ 14 s de retard, però la connexió continua. Es resol sense tancar amb coalescència: si d'una mateixa furgoneta hi ha una posició pendent a la cua, la nova la reemplaça en lloc d'afegir-s'hi. Amb 35 furgonetes, la cua mai no supera 35 entrades i l'operador rep sempre la posició més recent de cadascuna, amb retard acotat.
(c) Dues opcions. Al servidor: un emissor per connexió que, en lloc de buidar la cua missatge a missatge, cada segon agrupa el que hi ha pendent en un únic missatge {"tipus": "posicions", "items": [...]} amb l'última posició de cada furgoneta (coalescència + agrupació: 1 missatge/s, 35 posicions a dins). Al consumidor de Kafka: no publicar a mercat:<m> cada posició sinó mantenir a Redis un hash mercat:<m>:posicions (HSET furgoneta-3 <json>) i publicar un "tick" per segon; el servidor, en rebre el tick, llegeix el hash i envia l'estat complet. La primera és més simple i manté el bus sense canvis; la segona desacobla la freqüència del panell de la de les furgonetes i és més adequada si el panell creix a centenars d'operadors.
Conclusió
L'últim tram té regles pròpies perquè els seus clients són navegadors, mòbils i dispositius en xarxes que fallen, en nombres que cap servei intern no assoleix, i sense la confiança ni les biblioteques d'un consumidor de Kafka. "Temps real" aquí és latència percebuda, fixada cas per cas: segons per al mapa de l'Anna, un o dos per al panell del Jordi, desenes per a les alertes de la formatgeria, menys d'un per al xat. Amb aquesta taula al davant, la tècnica es tria sola: polling i long polling com a respatller, SSE per a l'unidireccional que no es pot perdre (amb Redis Streams i Last-Event-ID), WebSockets per al bidireccional i voluminós, i MQTT per als dispositius, amb tòpics jeràrquics que són també ACL, QoS 1 amb deduplicació per seqüència, missatges retinguts, testament i sessions persistents que sobreviuen al túnel. L'arquitectura d'extrem a extrem encadena la furgoneta, Mosquitto, el pont cap a repartiment.posicions, Flink i el servidor WebSocket de repartiment, que és l'únic que sap qui vol què; el fan-out entre les seves instàncies es resol amb un bus (Redis pub/sub) en lloc de sessions enganxoses, el JWT es verifica al handshake i es renova sense tallar, l'ordre i els duplicats es resolen amb el seq de 01-05 al servidor i al client, la pressió dels clients lents amb cues acotades i coalescència, i l'escalat amb instàncies intercanviables, límits del sistema operatiu ajustats i apagades esglaonades.
Tot el que s'ha construït fins aquí, des del broker MQTT fins al clúster de Kubernetes, corre en màquines que algú ha comprat, instal·lat i manté. La lliçó següent canvia aquesta premissa: què passa quan la infraestructura es converteix en una API d'un proveïdor de núvol, quins serveis gestionats substitueixen cada peça que hem muntat a mà, com es descriu tot amb Terraform, i què costa al mes. És la lliçó d'Aplicacions al Núvol.
Curs d'Arquitectures Distribuïdes
Mòdul 1: Introducció als Sistemes Distribuïts
- Conceptes Bàsics de Sistemes Distribuïts
- Models de Sistemes Distribuïts
- Avantatges i Desafiaments dels Sistemes Distribuïts
- Les Fal·làcies de la Computació Distribuïda
- Temps, Rellotges i Ordenació d'Esdeveniments
- Del Monòlit a la Plataforma Distribuïda: el Cas Quilòmetre Zero
Mòdul 2: Comunicació en Sistemes Distribuïts
- Protocols de Comunicació
- RPC i RMI
- gRPC i Serialització de Dades
- Missatgeria i Cues de Missatges
- Patrons de Comunicació Asíncrona
Mòdul 3: Consistència i Replicació
- Models de Consistència
- El Teorema CAP i PACELC
- Algorismes de Consens
- Replicació de Dades
- Transaccions Distribuïdes i Sagues
Mòdul 4: Emmagatzematge Distribuït
- Particionament de Dades i Hashing Consistent
- Sistemes de Fitxers Distribuïts
- Emmagatzematge d'Objectes
- Bases de Dades Distribuïdes
- Memòries Cau Distribuïdes
Mòdul 5: Computació Distribuïda
- Models de Computació Distribuïda
- MapReduce i Hadoop
- Spark i Computació en Memòria
- Processament de Fluxos de Dades
- Planificació de Treballs i Pipelines de Dades
Mòdul 6: Seguretat en Sistemes Distribuïts
- Autenticació i Autorització
- Xifratge i Protecció de Dades
- Gestió d'Identitats
- Seguretat entre Serveis: mTLS i Gestió de Secrets
- Passarel·les d'API, Limitació de Taxa i Auditoria
Mòdul 7: Monitoratge i Manteniment
- Monitoratge de Sistemes Distribuïts
- Logs Centralitzats i Traçabilitat Distribuïda
- Gestió de Fallades i Recuperació
- Patrons de Resiliència: Timeouts, Reintents i Circuit Breaker
- Automatització i Orquestració
- Proves en Sistemes Distribuïts i Enginyeria del Caos
