Tot el que hem calculat fins ara en el mòdul partia d'un conjunt tancat: el fitxer del 14 de setembre, els clics de set dies. Però les dades de Quilòmetre Zero no neixen tancades. comandes.esdeveniments rep centenars d'esdeveniments estoc.actualitzat per segon durant una campanya, els 140 repartidors (com furgoneta-3) envien 2,4 milions de posicions al dia, i hi ha preguntes que no poden esperar el lot nocturn: quins productes s'estan quedant sense estoc a Girona ara, on és cada repartidor ara, quin repartidor porta tres minuts sense donar senyal. El processament de fluxos respon a aquestes preguntes executant el càlcul de manera contínua, esdeveniment a esdeveniment, sobre un conjunt que no acaba mai. Això trenca suposicions que en els lots eren gratuïtes: no es pot "esperar a tenir-ho tot" abans d'agrupar, els esdeveniments arriben desordenats i amb retard respecte del moment en què van passar, l'estat del càlcul ha de sobreviure a les fallades sense poder reexecutar "des del principi", i l'entrada pot arribar més ràpid del que es processa. Aquesta lliçó posa nom a cadascun d'aquests problemes (temps d'esdeveniment, marques d'aigua, finestres, estat i checkpoints, exactly-once, backpressure), compara les arquitectures Lambda i Kappa i les tres eines habituals (Kafka Streams, Flink, Spark Structured Streaming), i resol els dos casos d'analitica en temps real: l'alerta d'estoc baix i el panell de repartiment. El lliurament d'aquests resultats al navegador de l'operador queda per a 08-02; aquí acabem en un tòpic de Kafka.
Contingut
- Dels lots als fluxos: què canvia
- Temps d'esdeveniment, temps de processament i marques d'aigua
- Finestres: tumbling, sliding i session
- Estat i checkpoints
- Exactly-once en streaming
- Backpressure
- Arquitectures Lambda i Kappa
- Eines: Kafka Streams, Flink i Spark Structured Streaming
- Els casos de Quilòmetre Zero: estoc baix i panell de repartiment
- Pràctica: simulació en Python, PyFlink i Structured Streaming
- Errors Comuns i Consells
- Exercicis
- Conclusió
- Dels lots als fluxos: què canvia
Recuperem la taula de 05-01 i l'afinem amb el que hem après des de llavors:
| Lot | Flux | |
|---|---|---|
| Entrada | Acotada: un fitxer, un directori, "el dia 14" | No acotada: un tòpic de Kafka que no acaba |
| Unitat | El conjunt sencer | L'esdeveniment: un fet immutable amb marca de temps (comanda.creada, estoc.actualitzat, una posició) |
| Quan es produeix la sortida | En acabar | Contínuament: cada esdeveniment, cada finestra, cada segon |
| "Totes les dades" | Existeix: s'espera a llegir-les | No existeix: cal decidir quan se n'ha vist prou (marques d'aigua) |
| Ordre | Es pot ordenar tot abans de calcular | Els esdeveniments arriben desordenats i tard |
| Estat | Viu al job; si falla, es reexecuta des de l'entrada | Viu per sempre a l'operador; cal persistir-lo (checkpoints) |
| Tolerància a fallades | Reexecutar el job (05-02, 05-03) | Restaurar l'estat i reprendre des de l'offset del checkpoint |
| Latència | Minuts a hores | Mil·lisegons a segons |
| Quilòmetre Zero | vendes_diaries.py |
Alerta d'estoc baix, panell de repartiment |
Un motor de fluxos executa el mateix tipus de DAG d'operadors que Spark (05-03), però desplegat de manera permanent: cada operador és un procés (o diversos, en paral·lel per clau) que espera esdeveniments, els processa i els emet cap al següent, sense que el DAG acabi mai. La font típica és Kafka (02-04), que aporta el que un flux necessita: particions (paral·lelisme), offsets (posició reprenible) i retenció (rellegir des d'enrere). El sink és un altre tòpic, una base de dades o un panell.
Hi ha dues maneres de moure esdeveniments pel DAG: esdeveniment a esdeveniment (Flink, Kafka Streams: latència de mil·lisegons) i per microlots (Spark Structured Streaming: cada pocs centenars de mil·lisegons o segons es pren el que ha arribat i s'executa com un lot petit). La segona reutilitza tota la maquinària de lots i simplifica l'exactly-once; la primera dona latències menors. Per al panell de repartiment i l'alerta d'estoc, segons són suficients, així que totes dues serveixen.
- Temps d'esdeveniment, temps de processament i marques d'aigua
A 01-05 vam distingir el moment en què una cosa va passar del moment en què un altre node se n'assabenta. En fluxos aquesta distinció és la que té més conseqüències:
- Temps d'esdeveniment (event time): quan va passar el fet, segons el rellotge de qui el va generar. És el
data_msde l'embolcall de 02-05. La posició defurgoneta-3presa a les 10:04:58 té aquest temps encara que el mòbil l'enviï a les 10:07 per falta de cobertura. - Temps de processament (processing time): quan l'operador veu l'esdeveniment, segons el seu rellotge. Depèn de la cua a Kafka, de la xarxa, de la càrrega.
La pregunta "quants estoc.actualitzat de formatge-curat hi va haver entre les 10:00 i les 10:05?" només té una resposta correcta en temps d'esdeveniment. Amb temps de processament, un reinici del consumidor a les 10:03 posaria a la finestra 10:05–10:10 esdeveniments que van passar a les 10:02; el resultat dependria del comportament del sistema, no dels fets. Però calcular en temps d'esdeveniment crea un problema nou: quan arriba un esdeveniment amb data_ms = 10:04:58 a les 10:07, la finestra 10:00–10:05 ja s'havia emès? Cal reobrir-la? Quant esperar abans de donar-la per tancada?
La resposta és la marca d'aigua (watermark): una afirmació que el sistema fa fluir pel DAG que diu "no espero més esdeveniments amb temps d'esdeveniment anterior a T". Es genera a la font amb una heurística, normalment màxim temps d'esdeveniment vist − retard tolerat: amb un retard tolerat de 2 minuts, després de veure un esdeveniment de les 10:07:00 la marca d'aigua és 10:05:00, i en aquell moment la finestra 10:00–10:05 es considera completa i s'emet. El retard tolerat és una decisió de negoci i de mesura: quant triguen de debò els esdeveniments a arribar (el percentil 99 del desfasament arribada − data_ms), contra quanta latència s'accepta en el resultat. Poc retard tolerat: resultats ràpids que descarten més esdeveniments; molt: resultats més complets però tardans.
Els esdeveniments que arriben després que la marca d'aigua hagi passat la seva finestra són esdeveniments tardans (late events). Cada motor ofereix tres destins: descartar-los (el comportament per defecte), incorporar-los reemetent la finestra actualitzada durant un marge addicional (allowed lateness a Flink, que obliga a mantenir la finestra a l'estat més temps), o desviar-los a una sortida lateral (side output) per comptar-los, registrar-los o reprocessar-los per lots. Comptar els tardans és obligatori: és la mètrica que diu si el retard tolerat està ben triat.
flowchart LR
subgraph W1[Finestra 10:00–10:05]
e1[e1 10:01]
e2[e2 10:03]
e4[e4 10:04:58<br/>arriba a les 10:08]
end
subgraph W2[Finestra 10:05–10:10]
e3[e3 10:06]
e5[e5 10:07]
end
e3 -. "marca d'aigua = 10:06 − 2 min = 10:04<br/>W1 continua oberta" .-> WM1[ ]
e5 -. "marca d'aigua = 10:07 − 2 min = 10:05<br/>W1 es tanca i s'emet" .-> WM2[ ]
e4 -. "arriba després del tancament: TARDÀ" .-> L[descartar / reemetre / sortida lateral]
Amb diverses particions de Kafka, cada partició té la seva pròpia marca d'aigua i la de l'operador és el mínim de totes: una partició sense trànsit (un mercat tancat de nit) reté la marca d'aigua global i bloqueja l'emissió de finestres de tothom. Els motors ho tracten amb idle timeouts que exclouen les particions inactives del mínim.
- Finestres: tumbling, sliding i session
Sobre un flux infinit, qualsevol agregació ("quants", "el mínim") necessita un límit: la finestra. Les tres formes bàsiques:
| Finestra | Definició | Un esdeveniment pertany a | Exemple a Quilòmetre Zero | Mida de l'estat |
|---|---|---|---|---|
| Tumbling (fixa, a salts) | Intervals consecutius i disjunts de mida fixa: 10:00–10:05, 10:05–10:10 | Exactament una finestra | Estoc mínim per producte i mercat cada 5 minuts | Una finestra oberta per clau (més les que esperen la marca d'aigua) |
| Sliding (lliscant) | Mida fixa, avanç menor: cada 1 minut, els últims 5 | Diverses finestres (mida / avanç) | Distància recorreguda per furgoneta-3 en els últims 5 minuts, refrescada cada minut |
Mida/avanç finestres per clau: 5 |
| Session (sessió) | Sense mida fixa: s'obre amb un esdeveniment i es tanca després d'un gap sense esdeveniments | Una sessió, que creix | "Repartidor sense senyal": sessió de posicions que es tanca després de 3 minuts sense rebre'n cap | Una sessió oberta per clau; les sessions es fusionen si un esdeveniment tardà les uneix |
| Global + trigger | Tota la història, amb disparadors explícits | Una | Comptador total de comandes del dia, emès cada 10 s | Un acumulador per clau |
Les finestres es combinen amb una clau: "tumbling de 5 minuts per (producte, mercat)" manté una finestra per cada parella, i el paral·lelisme de l'operador és per clau (totes les formatge-curat/girona van a la mateixa instància, com al shuffle de 05-02). I amb el temps d'esdeveniment: és el data_ms el que decideix a quina finestra cau un esdeveniment, i la marca d'aigua la que decideix quan s'emet.
flowchart TB
subgraph T[Tumbling 5 min]
direction LR
t1[10:00–10:05] --- t2[10:05–10:10] --- t3[10:10–10:15]
end
subgraph S[Sliding 5 min cada 1 min]
direction LR
s1[10:00–10:05]
s2[10:01–10:06]
s3[10:02–10:07]
end
subgraph G[Session gap 3 min]
direction LR
g1[10:00:10 … 10:04:50] -- "gap > 3 min: sense senyal" --- g2[10:09:30 … 10:21:00]
end
- Estat i checkpoints
Una finestra oberta, un comptador per clau, l'última posició coneguda de cada repartidor: tot això és estat, i en un flux viu indefinidament a l'operador. Els motors el desen en un state backend local a l'operador (memòria, o RocksDB a disc local per a estats de gigabytes) particionat per clau, exactament com un actor per clau (05-01). El problema és la durabilitat: si el node mor, el seu estat local es perd, i no es pot "reexecutar des del principi" perquè el principi va ser fa mesos.
La solució és el checkpoint: periòdicament (cada 10 s, cada minut) el motor escriu una còpia consistent de l'estat de tots els operadors, juntament amb els offsets de Kafka que aquell estat reflecteix, en emmagatzematge durador (HDFS, MinIO). Davant d'una fallada, restaura l'últim checkpoint en nodes sans i reprèn el consum des d'aquells offsets: els esdeveniments posteriors al checkpoint es tornen a processar, i l'estat acaba sent el mateix que si no hi hagués hagut fallada.
La paraula difícil és consistent: l'estat de tots els operadors ha de correspondre al mateix punt del flux, encara que cadascun vagi per un esdeveniment diferent. Flink ho aconsegueix amb l'algorisme de Chandy-Lamport adaptat (les barreres de checkpoint): la font injecta al flux un marcador amb el número de checkpoint; cada operador, en rebre'l per totes les seves entrades, desa el seu estat i reenvia el marcador; l'estat desat reflecteix exactament els esdeveniments anteriors al marcador. És una instantània distribuïda sense aturar el flux. Spark Structured Streaming ho té més fàcil: cada microlot és una unitat, i el checkpoint registra quins microlots s'han completat amb els seus rangs d'offsets. Kafka Streams desa l'estat en tòpics de Kafka (changelog topics) i el reconstrueix rellegint-los.
Un savepoint és un checkpoint disparat a mà, amb format estable, que es fa servir per aturar el treball, canviar el codi o el paral·lelisme, i reprendre des del mateix estat: l'equivalent d'un desplegament sense perdre "els últims cinc minuts". I la mida de l'estat importa: una session window per repartidor és petita, però "totes les comandes de les últimes 24 hores per client" són gigabytes que cal checkpointejar cada minut; els checkpoints incrementals (només el que ha canviat des de l'anterior) i un TTL per a l'estat que no es toca són les eines.
- Exactly-once en streaming
A 02-05 vam concloure que l'"exactly-once" en el lliurament no existeix, i que l'assolible és at-least-once més consumidors idempotents. En streaming el terme es fa servir amb un significat precís i assolible: l'estat del motor reflecteix cada esdeveniment exactament una vegada, encara que hi hagi fallades i reprocessaments. S'aconsegueix amb el checkpoint de l'apartat anterior: després d'una fallada, l'estat torna al del checkpoint i els esdeveniments posteriors es reapliquen; com que l'estat que reflectia aquests esdeveniments s'ha descartat, no hi ha doble compte.
El que el checkpoint no cobreix és el que ja va sortir del motor: l'alerta escrita al tòpic inventari.alertes, la fila inserida a PostgreSQL, abans de la fallada i després de l'últim checkpoint. En reprocessar, aquests efectes es produeixen una altra vegada. Perquè el resultat extern sigui també exactly-once (end-to-end), el sink ha de ser una de dues coses:
- Idempotent. Escriure de nou el mateix resultat no canvia res:
UPSERTper clau(finestra, producte, mercat)a PostgreSQL, unPUTa Redis amb la mateixa clau, un producer de Kafka amb una clau i un consumidor idempotent aigües avall (02-05). És l'opció més simple i la que es recomana sempre que la sortida tingui una clau natural. - Transaccional. El sink escriu en una transacció que només es confirma quan el checkpoint es completa (two-phase commit sink): Flink amb el producer transaccional de Kafka (
DeliveryGuarantee.EXACTLY_ONCE), o amb una taula i una transacció per checkpoint. Entre el pre-commit i el commit les dades existeixen però no són visibles per a consumidors ambisolation.level=read_committed. Afegeix latència (l'interval de checkpoint) i complexitat; es reserva per a sinks sense clau natural (un log d'esdeveniments a Kafka).
| Garantia | Què passa després d'una fallada | Com s'aconsegueix |
|---|---|---|
| At-most-once | Es perden esdeveniments entre la fallada i la represa | Confirmar offsets abans de processar; sense checkpoint d'estat |
| At-least-once | Es reprocessen esdeveniments; l'estat o el sink els poden comptar dues vegades | Checkpoint d'offsets sense coordinació amb l'estat; o sink no idempotent |
| Exactly-once (estat) | L'estat és el mateix que sense fallada | Checkpoint consistent d'estat + offsets |
| Exactly-once end-to-end | A més, el sink no mostra duplicats | L'anterior + sink idempotent o transaccional |
- Backpressure
Un DAG de fluxos és una cadena de productors i consumidors, i en qualsevol moment un pot anar més a poc a poc que l'anterior: l'operador de finestres escrivint un checkpoint gran, el sink de PostgreSQL saturat, un pic de la Setmana del Formatge Artesà que triplica els estoc.actualitzat. Si l'operador ràpid continués enviant, les cues intermèdies creixerien fins a esgotar la memòria. Backpressure (contrapressió) és el mecanisme pel qual la lentitud es propaga cap enrere: l'operador lent deixa d'acceptar, l'anterior omple el seu búfer de sortida i deixa de llegir de la seva entrada, i així fins a la font, que deixa de consumir de Kafka. Els esdeveniments s'acumulen a Kafka (que està dissenyat per a això, amb retenció de dies), no a la memòria del motor, i el lag del grup de consumidors (02-04) es converteix en la mètrica que indica que el flux no dona l'abast.
Flink ho implementa amb crèdits entre tasques (el receptor anuncia quant pot rebre); Kafka Streams ho té gratis perquè cada instància fa poll només quan ha acabat amb el lot anterior; Spark Structured Streaming ho aproxima limitant quants offsets llegeix per microlot (maxOffsetsPerTrigger). El que cap d'ells no fa és resoldre la causa: si el lag creix de manera sostinguda, cal augmentar el paral·lelisme (més particions i més instàncies), alleugerir l'operador o acceptar resultats aproximats. I compte amb la marca d'aigua: sota backpressure, el temps de processament s'allunya del temps d'esdeveniment, però les finestres continuen sent correctes perquè es defineixen en temps d'esdeveniment; això és precisament el que l'apartat 2 comprava.
- Arquitectures Lambda i Kappa
Quan l'streaming era nou i poc fiable, la resposta a "vull resultats en temps real però també exactes" va ser l'arquitectura Lambda (Marz, 2011): mantenir dos camins. La capa batch recalcula cada nit, des del llac, les vistes completes i exactes; la capa de velocitat calcula amb streaming les últimes hores, de manera aproximada; una capa de servei combina totes dues en consultar. Funciona, però obliga a escriure i mantenir la mateixa lògica dues vegades, en dos motors, amb dues semàntiques, i a reconciliar-ne les diferències. L'arquitectura Kappa (Kreps, 2014) proposa un sol camí: tot és un flux, amb Kafka retenint l'historial (o un llac d'esdeveniments rellegible), i el "lot" és simplement reprocessar el flux des d'un offset antic amb la mateixa aplicació d'streaming, escrivint en una taula nova i canviant el punter quan arriba al present.
| Lambda | Kappa | |
|---|---|---|
| Camins | Dos: batch (exacte, lent) + velocitat (aproximat, ràpid) | Un: streaming, amb reprocessament des de l'historial |
| Codi | Duplicat en dos motors | Una sola aplicació |
| Reprocessar | Rellançar el lot | Rellançar l'aplicació des d'un offset o des del llac |
| Exactitud | Batch corregeix velocitat | L'streaming ha de ser exacte (temps d'esdeveniment, exactly-once) |
| Requisits | Un llac i un motor de lots; un motor de fluxos | Retenció llarga a Kafka o llac rellegible; motor de fluxos amb estat |
| Quan | Lògica batch complexa que no cap en streaming (entrenar ALS); resultats històrics que necessiten tot el conjunt | Agregacions, alertes, materialitzacions; quan la latència importa i la lògica és la mateixa |
| Quilòmetre Zero | Vendes diàries (batch) + panell en temps real (velocitat): Lambda de facto | Alerta d'estoc i panell de repartiment: Kappa pur |
La plataforma de dades de Quilòmetre Zero acaba sent pragmàticament mixta: vendes_diaries.py i ALS són lots perquè necessiten el conjunt complet i no tenen pressa; l'estoc i el repartiment són fluxos perquè la latència és el requisit. El que Kappa aporta és el criteri: si la mateixa lògica s'escriu dues vegades, alguna cosa està malament; que els motors moderns executin el mateix codi en lot i en flux (Spark, Flink) fa que la tria sigui de desplegament, no de reescriptura.
- Eines: Kafka Streams, Flink i Spark Structured Streaming
| Kafka Streams | Apache Flink | Spark Structured Streaming | |
|---|---|---|---|
| Què és | Una llibreria Java/Scala: l'aplicació és un procés normal que consumeix i produeix a Kafka | Un motor amb clúster propi (JobManager + TaskManagers), o sobre YARN/Kubernetes | El mode streaming del motor Spark (05-03) |
| Model | Esdeveniment a esdeveniment | Esdeveniment a esdeveniment | Microlots (100 ms–segons); mode continu experimental |
| Fonts/sinks | Només Kafka (per disseny) | Kafka, fitxers, JDBC, Kinesis, Pulsar, CDC... | Kafka, fitxers, sockets; sinks Kafka, fitxers, foreachBatch per a tota la resta |
| Temps d'esdeveniment i marques d'aigua | Sí, amb grace period | Sí, el més complet (allowed lateness, side outputs, temporitzadors) | withWatermark; sense allowed lateness ni side outputs |
| Finestres | Tumbling, hopping, sliding, session | Totes, més finestres definides per l'usuari | Tumbling, sliding, session |
| Estat | RocksDB local + changelog a Kafka | RocksDB o heap; checkpoints incrementals; savepoints | Estat per microlot a HDFS; RocksDB des de 3.2 |
| Exactly-once | Transaccions de Kafka (processing.guarantee=exactly_once_v2) |
Checkpoint + 2PC sinks | Checkpoint + sinks idempotents; Kafka sink at-least-once |
| Llenguatges | Java, Scala (Kotlin) | Java, Scala, Python (PyFlink), SQL | Scala, Java, Python, R, SQL |
| Latència típica | ms | ms | segons |
| Encaixa | Microserveis que transformen tòpics; equips Java; sense clúster nou | Streaming exigent: estat gran, latència baixa, semàntica precisa | Equips que ja fan servir Spark; lot i flux amb el mateix codi |
| Quilòmetre Zero | Seria natural dins d'inventari (Java), però els serveis són Python |
Panell de repartiment i alertes (PyFlink) |
Alternativa amb vendes_diaries.py reutilitzat |
Triem Flink com a motor de fluxos d'analitica per la semàntica de temps d'esdeveniment i per PyFlink, i mantenim Structured Streaming com a alternativa perquè reutilitza l'API de 05-03. Kafka Streams queda anomenat: és l'opció correcta si l'streaming viu dins d'un servei Java i no en una plataforma de dades.
- Els casos de Quilòmetre Zero: estoc baix i panell de repartiment
Cas (a): alerta d'estoc baix. inventari publica estoc.actualitzat a comandes.esdeveniments amb cada canvi (04-05 els feia servir per invalidar Redis). Dades de l'esdeveniment: producte, mercat, estoc_actual, delta, replica (inv-bcn o inv-vlc). Requisit: cada 5 minuts, per producte i mercat, si l'estoc mínim observat ha baixat de 10 unitats, emetre una alerta al tòpic inventari.alertes amb la finestra, el mínim i quantes actualitzacions hi va haver. Finestra tumbling de 5 minuts per (producte, mercat), en temps d'esdeveniment (el data_ms d'inventari, perquè una rèplica amb retard no ha de desplaçar l'alerta), amb marca d'aigua de 2 minuts (mesurat: el percentil 99 del desfasament és 40 s). Sink idempotent: la clau del missatge d'alerta és finestra|producte|mercat, així que un reprocessament produeix el mateix missatge i el consumidor de 08-02 el tractarà com el mateix.
Cas (b): panell de repartiment. Les posicions de furgoneta-3 arriben per MQTT (02-01) i un pont les publica al tòpic repartiment.posicions amb clau = id de repartidor: {"repartidor":"furgoneta-3","lat":41.9794,"lon":2.8214,"data_ms":...}. Dos càlculs: la distància recorreguda en els últims 5 minuts, refrescada cada minut (sliding 5/1 per repartidor: detecta repartidors aturats o desviats), i repartidors sense senyal (session window amb gap de 3 minuts: quan la sessió es tanca, el repartidor porta 3 minuts sense enviar; quan se n'obre una de nova, ha tornat). Tots dos escriuen a repartiment.panell, que 08-02 empenyerà als navegadors dels operadors.
flowchart LR
K1[(Kafka<br/>comandes.esdeveniments<br/>6 particions)] --> F1[filtrar<br/>estoc.actualitzat]
F1 --> WM1[assignar ts + marca d'aigua<br/>ts − 2 min]
WM1 --> KB1[[keyBy producte, mercat]]
KB1 --> V1[tumbling 5 min<br/>MIN estoc, COUNT]
V1 --> H1[HAVING min < 10]
H1 --> K2[(Kafka<br/>inventari.alertes)]
K3[(Kafka<br/>repartiment.posicions<br/>clau = repartidor)] --> WM2[assignar ts + marca d'aigua<br/>ts − 30 s]
WM2 --> KB2[[keyBy repartidor]]
KB2 --> V2[sliding 5 min / 1 min<br/>distància]
KB2 --> V3[session gap 3 min<br/>inici, fi, n]
V2 --> K4[(Kafka<br/>repartiment.panell)]
V3 --> K4
- Pràctica: simulació en Python, PyFlink i Structured Streaming
10.1 simulacions/finestra_tumbling.py: un motor de finestres en 80 línies
Abans de fer servir un motor, convé veure el mecanisme nu. Aquest script processa una llista d'esdeveniments estoc.actualitzat amb dos temps cadascun, el d'esdeveniment (data_ms) i el d'arribada (arribada_ms), en ordre d'arribada, mantenint finestres tumbling de 5 minuts per producte i una marca d'aigua amb 2 minuts de retard tolerat. Els temps són en minuts des de les 10:00 per llegir-los fàcilment.
# km0/simulacions/finestra_tumbling.py
"""Finestres tumbling de 5 min en temps d'esdeveniment, amb marca d'aigua i esdeveniments tardans, en Python pur."""
from collections import defaultdict
FINESTRA = 5 # minuts
RETARD_TOLERAT = 2 # marca d'aigua = màxim temps d'esdeveniment vist − 2 min
LLINDAR = 10
# (temps d'esdeveniment, temps d'arribada, producte, estoc_actual), en minuts des de les 10:00.
# Són en ORDRE D'ARRIBADA, que és l'ordre en què els veu l'operador.
ESDEVENIMENTS = [
(0.5, 0.6, "formatge-curat", 42),
(1.2, 1.3, "formatge-curat", 31),
(2.0, 2.1, "tomaquet-rosa", 120),
(3.8, 4.0, "formatge-curat", 12),
(6.1, 6.2, "formatge-curat", 6),
(6.5, 6.6, "tomaquet-rosa", 118),
(4.9, 7.0, "formatge-curat", 8), # va passar a les 10:04:54, arriba a les 10:07 (rèplica inv-vlc amb retard)
(7.3, 7.4, "formatge-curat", 25), # reposició
(4.2, 9.5, "formatge-curat", 9), # va passar a les 10:04:12, arriba a les 10:09:30: TARDÀ
(10.2, 10.3, "tomaquet-rosa", 117),
(12.7, 12.8, "formatge-curat", 22),
]
def inici_finestra(t: float) -> int:
return int(t // FINESTRA) * FINESTRA
def processar(esdeveniments):
finestres = defaultdict(lambda: {"min": float("inf"), "n": 0}) # (inici, producte) -> estat
marca_aigua = float("-inf")
tardans = 0
for t_esdeveniment, t_arribada, producte, estoc in esdeveniments:
ini = inici_finestra(t_esdeveniment)
if ini + FINESTRA <= marca_aigua: # la finestra ja s'ha emès: esdeveniment tardà
tardans += 1
print(f" [{t_arribada:5.1f}] TARDÀ: {producte} t={t_esdeveniment} (finestra {ini}-{ini + FINESTRA} tancada, "
f"marca d'aigua {marca_aigua})")
continue
estat = finestres[(ini, producte)] # estat per clau i finestra
estat["min"] = min(estat["min"], estoc); estat["n"] += 1
marca_aigua = max(marca_aigua, t_esdeveniment - RETARD_TOLERAT)
print(f" [{t_arribada:5.1f}] {producte:14s} t={t_esdeveniment:4.1f} estoc={estoc:3d} marca d'aigua={marca_aigua:4.1f}")
# emetre tota finestra el final de la qual hagi quedat per sota de la marca d'aigua
for (f_ini, f_prod) in sorted(k for k in finestres if k[0] + FINESTRA <= marca_aigua):
e = finestres.pop((f_ini, f_prod))
alerta = " <-- ALERTA estoc baix" if e["min"] < LLINDAR else ""
print(f" EMETRE finestra {f_ini:2d}-{f_ini + FINESTRA:2d} {f_prod:14s} min={e['min']:3d} n={e['n']}{alerta}")
print(f"Fi de l'entrada: {len(finestres)} finestres obertes sense emetre, {tardans} esdeveniments tardans")
if __name__ == "__main__":
processar(ESDEVENIMENTS)$ python finestra_tumbling.py
[ 0.6] formatge-curat t= 0.5 estoc= 42 marca d'aigua=-1.5
[ 1.3] formatge-curat t= 1.2 estoc= 31 marca d'aigua=-0.8
[ 2.1] tomaquet-rosa t= 2.0 estoc=120 marca d'aigua= 0.0
[ 4.0] formatge-curat t= 3.8 estoc= 12 marca d'aigua= 1.8
[ 6.2] formatge-curat t= 6.1 estoc= 6 marca d'aigua= 4.1
[ 6.6] tomaquet-rosa t= 6.5 estoc=118 marca d'aigua= 4.5
[ 7.0] formatge-curat t= 4.9 estoc= 8 marca d'aigua= 4.5
[ 7.4] formatge-curat t= 7.3 estoc= 25 marca d'aigua= 5.3
EMETRE finestra 0- 5 formatge-curat min= 8 n=4 <-- ALERTA estoc baix
EMETRE finestra 0- 5 tomaquet-rosa min=120 n=1
[ 9.5] TARDÀ: formatge-curat t=4.2 (finestra 0-5 tancada, marca d'aigua 5.3)
[ 10.3] tomaquet-rosa t=10.2 estoc=117 marca d'aigua= 8.2
[ 12.8] formatge-curat t=12.7 estoc= 22 marca d'aigua=10.7
EMETRE finestra 5-10 formatge-curat min= 6 n=2 <-- ALERTA estoc baix
EMETRE finestra 5-10 tomaquet-rosa min=118 n=1
Fi de l'entrada: 2 finestres obertes sense emetre, 1 esdeveniments tardansEl que mostra la traça:
- L'esdeveniment de
t=4.9que arriba a les 10:07 sí que entra a la finestra 0–5, perquè la marca d'aigua en aquell moment era 4,5 (< 5): la finestra continuava oberta gràcies al retard tolerat. Amb temps de processament hauria caigut a 5–10, i el mínim de la finestra 0–5 hauria estat 12, sense alerta. - La finestra 0–5 s'emet quan arriba l'esdeveniment de
t=7.3: la marca d'aigua passa a 5,3 ≥ 5. Ni abans (no se sabia si faltaven esdeveniments) ni després (no cal esperar més). - L'esdeveniment de
t=4.2que arriba a les 10:09:30 és tardà: la seva finestra ja s'ha emès. Aquí es descarta i es compta; amb allowed lateness es reemetria la finestra 0–5 ambmin=8, n=5(el mínim no canvia, però el compte sí). - En acabar l'entrada queden dues finestres obertes (10–15): en un flux real no hi ha "fi", i s'emetran quan la marca d'aigua arribi a 15. En un lot, el final de l'entrada dispara l'emissió de tot.
Observa també que la marca d'aigua no avança amb l'esdeveniment tardà ni retrocedeix mai: és monòtona. I que cada línia d'EMETRE és una sortida que, si el procés morís i reprocessés des de l'esdeveniment 1, es produiria una altra vegada amb els mateixos valors: la clau (finestra, producte) la fa idempotent.
10.2 El cas (a) a PyFlink (Table API / SQL)
El clúster Flink s'afegeix al docker-compose.yml amb dos serveis (jobmanager i taskmanager de la imatge flink:1.19-python, o una imatge pròpia amb pip install apache-flink), i el treball s'envia amb flink run -py. Amb la Table API, el DAG del cas (a) són tres sentències SQL:
# km0/serveis/analitica/fluxos/alertes_estoc.py
"""Alerta d'estoc baix: tumbling 5 min per (producte, mercat) en temps d'esdeveniment, des de comandes.esdeveniments."""
from pyflink.table import EnvironmentSettings, TableEnvironment
t_env = TableEnvironment.create(EnvironmentSettings.in_streaming_mode())
t_env.get_config().set("pipeline.name", "km0-alertes-estoc")
t_env.get_config().set("execution.checkpointing.interval", "30 s") # checkpoints cada 30 s
t_env.get_config().set("table.exec.source.idle-timeout", "1 min") # particions sense trànsit no frenen la marca d'aigua
# Font: el tòpic d'esdeveniments amb l'embolcall de 02-05. 'dades' és una fila niada.
t_env.execute_sql("""
CREATE TABLE esdeveniments (
id_esdeveniment STRING,
tipus STRING,
versio INT,
data_ms BIGINT,
origen STRING,
dades ROW<producte STRING, mercat STRING, estoc_actual INT, delta INT, replica STRING>,
ts AS TO_TIMESTAMP_LTZ(data_ms, 3), -- temps d'ESDEVENIMENT, derivat de data_ms
WATERMARK FOR ts AS ts - INTERVAL '2' MINUTE -- marca d'aigua: 2 min de retard tolerat
) WITH (
'connector' = 'kafka',
'topic' = 'comandes.esdeveniments',
'properties.bootstrap.servers' = 'kafka:9092',
'properties.group.id' = 'analitica-alertes-estoc',
'scan.startup.mode' = 'group-offsets', -- reprèn on ho va deixar (checkpoint)
'format' = 'json',
'json.ignore-parse-errors' = 'true' -- un esdeveniment corrupte no tomba el job
)""")
# Sink: alertes amb clau (finestra, producte, mercat). L'upsert-kafka escriu per clau: idempotent.
t_env.execute_sql("""
CREATE TABLE alertes_estoc (
finestra_inici TIMESTAMP_LTZ(3),
finestra_fi TIMESTAMP_LTZ(3),
producte STRING,
mercat STRING,
estoc_min INT,
actualitzacions BIGINT,
PRIMARY KEY (finestra_inici, producte, mercat) NOT ENFORCED
) WITH (
'connector' = 'upsert-kafka',
'topic' = 'inventari.alertes',
'properties.bootstrap.servers' = 'kafka:9092',
'key.format' = 'json',
'value.format' = 'json'
)""")
# El càlcul: finestra tumbling de 5 minuts sobre ts, per producte i mercat, només estoc.actualitzat.
t_env.execute_sql("""
INSERT INTO alertes_estoc
SELECT window_start, window_end,
dades.producte, dades.mercat,
MIN(dades.estoc_actual) AS estoc_min,
COUNT(*) AS actualitzacions
FROM TABLE(TUMBLE(TABLE esdeveniments, DESCRIPTOR(ts), INTERVAL '5' MINUTE))
WHERE tipus = 'estoc.actualitzat'
GROUP BY window_start, window_end, dades.producte, dades.mercat
HAVING MIN(dades.estoc_actual) < 10
""").wait()Cada peça correspon a un apartat: ts AS TO_TIMESTAMP_LTZ(data_ms, 3) declara el temps d'esdeveniment; WATERMARK FOR ts AS ts - INTERVAL '2' MINUTE la marca d'aigua (Flink la genera per partició de Kafka i en pren el mínim, amb l'idle-timeout per a particions aturades); TUMBLE(..., INTERVAL '5' MINUTE) la finestra amb window_start/window_end com a columnes; el GROUP BY és el keyBy que reparteix l'estat per clau entre TaskManagers; execution.checkpointing.interval el checkpoint que desa estat i offsets a l'state.checkpoints.dir configurat (HDFS /km0/checkpoints/); i upsert-kafka amb PRIMARY KEY el sink idempotent: un reprocessament després d'una fallada reescriu la mateixa clau amb el mateix valor. Es llança i s'observa així:
docker compose exec jobmanager flink run -py /app/analitica/fluxos/alertes_estoc.py -d
docker compose exec kafka kafka-console-consumer --bootstrap-server kafka:9092 \
--topic inventari.alertes --property print.key=true --from-beginning
{"finestra_inici":"2026-09-14 10:00:00Z","producte":"formatge-curat","mercat":"girona"} {"finestra_inici":"2026-09-14 10:00:00Z","finestra_fi":"2026-09-14 10:05:00Z","producte":"formatge-curat","mercat":"girona","estoc_min":8,"actualitzacions":4}La interfície web del JobManager (localhost:8081) mostra el DAG desplegat, el paral·lelisme de cada operador, la marca d'aigua actual de cadascun, els checkpoints (durada, mida) i el backpressure per operador, acolorit.
Per al cas (b), la mateixa estructura amb SESSION per als repartidors sense senyal:
INSERT INTO repartiment_panell
SELECT repartidor,
SESSION_START(ts, INTERVAL '3' MINUTE) AS inici,
SESSION_END(ts, INTERVAL '3' MINUTE) AS fi,
COUNT(*) AS posicions
FROM posicions -- taula sobre repartiment.posicions, marca d'aigua ts - 30 s
GROUP BY repartidor, SESSION(ts, INTERVAL '3' MINUTE)Cada fila emesa significa "el repartidor furgoneta-3 va enviar posicions de manera contínua entre inici i fi, i després va estar com a mínim 3 minuts sense senyal": és l'avís del panell. La finestra lliscant de distància s'escriu amb HOP(TABLE posicions, DESCRIPTOR(ts), INTERVAL '1' MINUTE, INTERVAL '5' MINUTE) i una funció d'agregació pròpia (la distància entre posicions consecutives necessita l'ordre, que en SQL es resol amb LAG sobre una finestra OVER abans d'agregar).
10.3 El cas (a) a Spark Structured Streaming
El mateix càlcul amb l'API de DataFrames de 05-03, en mode streaming:
# km0/serveis/analitica/fluxos/alertes_estoc_spark.py
"""Alerta d'estoc baix amb Spark Structured Streaming: withWatermark + window tumbling de 5 min."""
from pyspark.sql import SparkSession, functions as F, types as T
ESQUEMA = T.StructType([
T.StructField("id_esdeveniment", T.StringType()), T.StructField("tipus", T.StringType()),
T.StructField("versio", T.IntegerType()), T.StructField("data_ms", T.LongType()),
T.StructField("origen", T.StringType()),
T.StructField("dades", T.StructType([
T.StructField("producte", T.StringType()), T.StructField("mercat", T.StringType()),
T.StructField("estoc_actual", T.IntegerType()), T.StructField("delta", T.IntegerType()),
T.StructField("replica", T.StringType())])),
])
spark = SparkSession.builder.appName("km0-alertes-estoc").getOrCreate()
cru = (spark.readStream.format("kafka") # font NO acotada
.option("kafka.bootstrap.servers", "kafka:9092")
.option("subscribe", "comandes.esdeveniments")
.option("startingOffsets", "latest")
.option("maxOffsetsPerTrigger", 50000) # backpressure: sostre per microlot
.load())
estoc = (cru.select(F.from_json(F.col("value").cast("string"), ESQUEMA).alias("e")).select("e.*")
.filter(F.col("tipus") == "estoc.actualitzat")
.withColumn("ts", (F.col("data_ms") / 1000).cast("timestamp"))) # temps d'esdeveniment
alertes = (estoc
.withWatermark("ts", "2 minutes") # marca d'aigua: 2 min
.groupBy(F.window("ts", "5 minutes").alias("finestra"), # tumbling de 5 min
F.col("dades.producte").alias("producte"), F.col("dades.mercat").alias("mercat"))
.agg(F.min("dades.estoc_actual").alias("estoc_min"), F.count("*").alias("actualitzacions"))
.filter(F.col("estoc_min") < 10)
.select(F.concat_ws("|", F.col("finestra.start"), "producte", "mercat").alias("key"), # clau: idempotent
F.to_json(F.struct(F.col("finestra.start").alias("finestra_inici"), F.col("finestra.end").alias("finestra_fi"),
"producte", "mercat", "estoc_min", "actualitzacions")).alias("value")))
consulta = (alertes.writeStream
.outputMode("append") # emet cada finestra UNA vegada, quan la marca d'aigua la tanca
.format("kafka")
.option("kafka.bootstrap.servers", "kafka:9092")
.option("topic", "inventari.alertes")
.option("checkpointLocation", "hdfs://namenode:8020/km0/checkpoints/alertes-estoc") # offsets + estat
.trigger(processingTime="10 seconds") # un microlot cada 10 s
.start())
consulta.awaitTermination()La correspondència amb Flink és directa: withWatermark ↔ WATERMARK FOR, F.window("ts", "5 minutes") ↔ TUMBLE, checkpointLocation ↔ execution.checkpointing, maxOffsetsPerTrigger ↔ backpressure. Dues diferències importants. L'outputMode("append") amb marca d'aigua emet cada finestra una sola vegada, en tancar-la, i descarta els tardans sense possibilitat de reemetre (no hi ha allowed lateness); update emetria resultats provisionals a cada microlot, que exigeixen un sink que sàpiga sobreescriure. I el sink de Kafka d'Spark és at-least-once: el missatge es pot duplicar després d'una fallada, i és la clau key (idèntica al duplicat) la que fa que el consumidor de 08-02 el tracti com una repetició innòcua, exactament el consumidor idempotent de 02-05. Es llança amb spark-submit --packages org.apache.spark:spark-sql-kafka-0-10_2.12:3.5.1 alertes_estoc_spark.py, i a localhost:4040 apareix una pestanya Structured Streaming amb la taxa d'entrada, la durada de cada microlot i la marca d'aigua.
Errors Comuns i Consells
- Fer servir temps de processament perquè és més fàcil. Els resultats depenen de la càrrega, els reinicis i el backpressure, i no es poden reproduir. Si l'esdeveniment té marca de temps (i amb l'embolcall de 02-05 sempre la té), fes servir temps d'esdeveniment.
- Marca d'aigua sense mesurar. Un retard tolerat inventat descarta esdeveniments vàlids o retarda els resultats. Mesura el desfasament
arribada − data_msen producció (percentil 99) i compta els tardans com a mètrica permanent. - Oblidar les particions inactives. Una partició de Kafka sense trànsit reté la marca d'aigua global i cap finestra no s'emet. Configura l'idle timeout (Flink) o revisa que totes les particions rebin esdeveniments.
- Estat sense límit. Una agregació per clau sense finestra ni TTL creix per sempre (una clau per
comanda_id, per exemple). Finestres, TTL d'estat, o claus acotades. - Sink no idempotent amb reprocessament. Després d'una fallada, les sortides posteriors al checkpoint es repeteixen. Clau natural + upsert, o sink transaccional. Mai un
INSERTsense clau ni unappenda un fitxer. - Checkpoint a disc local. Si el node mor, el checkpoint mor amb ell. HDFS, MinIO o un emmagatzematge replicat, sempre.
- Canviar el DAG i reprendre del checkpoint. Un checkpoint desa estat per operador; si el DAG canvia (operador nou, una altra clau) pot no ser compatible. Savepoints amb identificadors d'operador (
uid) a Flink; a Spark, els canvis d'esquema d'estat solen exigir començar de zero. - Confondre lag amb latència. El lag del grup de consumidors (esdeveniments pendents a Kafka) creix sota backpressure i és el senyal que falta paral·lelisme; la latència d'un resultat en temps d'esdeveniment és sempre com a mínim el retard tolerat més la mida de la finestra.
- Reescriure la lògica dues vegades (Lambda per inèrcia). Si el mateix agregat es calcula en lot i en flux, sortirà diferent i ningú no sabrà quin és el bo. Un codi, dos desplegaments.
Exercicis
Exercici 1: Triar marca d'aigua i finestra
Una anàlisi del tòpic repartiment.posicions durant una setmana mostra que el desfasament entre data_ms i l'arribada a Kafka té aquesta distribució: mediana 1,2 s; percentil 95, 8 s; percentil 99, 45 s; percentil 99,9, 4 min (túnels, zones sense cobertura); màxim 22 min (un mòbil apagat que va reenviar en encendre's). El panell ha de mostrar la posició i la distància dels últims 5 minuts amb no més d'1 minut de retard. Tria el retard tolerat de la marca d'aigua, el tipus i mida de finestra, i què fer amb els esdeveniments tardans. Justifica els percentatges d'esdeveniments que es descartaran i què passa amb el mòbil que reenvia després de 22 minuts.
Exercici 2: Traçar la simulació amb una altra marca d'aigua
Executa mentalment (o modificant l'script) finestra_tumbling.py amb RETARD_TOLERAT = 0 i amb RETARD_TOLERAT = 5. Per a cada cas, indica quan s'emet la finestra 0–5 de formatge-curat, amb quin mínim i compte, i quants esdeveniments tardans hi ha. Quina de les tres configuracions (0, 2, 5) donaria l'alerta correcta més aviat?
Exercici 3: Fallada i reprocessament
El job de PyFlink d'alertes porta checkpoints cada 30 s. A les 10:07:50 el TaskManager que executa la finestra de formatge-curat/girona mor; l'últim checkpoint completat és de les 10:07:30, i a les 10:07:40 el job havia emès l'alerta de la finestra 10:00–10:05. Descriu què fa Flink en recuperar-se: des de quins offsets llegeix, què passa amb l'estat de la finestra 10:05–10:10, si l'alerta 10:00–10:05 es torna a emetre i què veu el consumidor d'inventari.alertes. Després, explica què canviaria si el sink fos un INSERT a PostgreSQL sense clau primària, i com ho arreglaries.
Solucions
Exercici 1.
El pressupost de latència és 1 minut, i la latència mínima d'un resultat és el retard tolerat (més l'interval d'emissió). Un retard tolerat de 45 s (el percentil 99) compleix el pressupost i descarta com a tardanes l'1 % de les posicions; amb 8 s (p95) se'n descartaria el 5 %, massa per a una traça de posicions; amb 4 min (p99,9) es violaria el requisit d'1 minut. Finestra lliscant de 5 minuts amb avanç d'1 minut per repartidor (o de 30 s si el panell s'ha de refrescar més sovint, a costa de més finestres obertes per clau: 10 en comptes de 5). Els tardans (1 %) van a una sortida lateral que es compta i s'escriu al llac: per al panell en temps real tant li fa perdre una posició de cada cent, però la distància recorreguda del dia que calcula el lot nocturn les ha d'incloure, i per a això el lot llegeix totes les posicions del llac, tardanes incloses (una Lambda de facto justificada). El mòbil que reenvia després de 22 minuts lliura posicions la finestra de les quals es va tancar fa 20 minuts: totes tardanes, totes a la sortida lateral; el panell les ignora i el lot les incorpora. Alternativa si el negoci ho demanés: allowed lateness de 5 minuts per reemetre finestres recents, no de 22 (mantindria 27 minuts de finestres a l'estat per repartidor).
Exercici 2.
Amb RETARD_TOLERAT = 0 la marca d'aigua és el màxim temps d'esdeveniment vist. En arribar t=6.1 (arribada 6,2) la marca d'aigua és 6,1 ≥ 5 i la finestra 0–5 s'emet amb els esdeveniments vistos fins llavors: 0,5, 1,2, 3,8 i... el de t=4.9 va arribar a les 7,0, després, així que la finestra s'emet amb min=12, n=3: sense alerta, incorrecta. Els esdeveniments de t=4.9 (arribada 7,0) i t=4.2 (arribada 9,5) són tots dos tardans: 2 tardans. Emissió més primerenca, resultat equivocat.
Amb RETARD_TOLERAT = 5, la finestra 0–5 s'emet quan la marca d'aigua arriba a 5, és a dir, en veure un esdeveniment amb t ≥ 10: el de t=10.2 (arribada 10,3). Per llavors han entrat 0,5, 1,2, 3,8, 4,9 i 4,2 (que va arribar a les 9,5 amb la finestra encara oberta): min=8, n=5, correcte i complet, 0 tardans. Però l'alerta surt a les 10:10:18, cinc minuts més tard que amb retard 2 (10:07:24, min=8, n=4).
La configuració amb 2 minuts dona l'alerta correcta (el mínim 8 era a totes dues) més aviat; la de 5 minuts dona a més el compte exacte; la de 0 falla. És el compromís completesa/latència de l'apartat 2, i la raó que el retard tolerat es mesuri i no s'endevini.
Exercici 3.
Flink detecta la mort del TaskManager (heartbeat), reinicia el job sencer (o la regió afectada, amb fine-grained recovery) des del checkpoint de les 10:07:30: restaura l'estat de tots els operadors tal com era llavors (la finestra 10:05–10:10 amb els esdeveniments anteriors a les 10:07:30 ja aplicats; la 10:00–10:05, que encara no s'havia emès en aquell instant, també restaurada amb el seu estat) i reposiciona el consumidor de Kafka als offsets desats en aquell checkpoint. Els esdeveniments entre les 10:07:30 i les 10:07:50 es tornen a llegir i a aplicar sobre aquell estat: cap no es compta dues vegades, perquè l'estat que els contenia s'ha descartat. La marca d'aigua torna a avançar i la finestra 10:00–10:05 s'emet de nou, amb el mateix contingut. Com que el sink és upsert-kafka amb clau (finestra_inici, producte, mercat), el tòpic rep un segon missatge amb la mateixa clau i el mateix valor; el consumidor de 08-02 (o la compactació de Kafka) el tracta com una actualització sense canvis. Exactly-once a l'estat, i efectivament una sola alerta visible.
Amb un INSERT a PostgreSQL sense clau, la fila de l'alerta 10:00–10:05 existiria dues vegades. Arranjaments, de menor a major esforç: clau primària (finestra_inici, producte, mercat) amb INSERT ... ON CONFLICT DO UPDATE (sink idempotent, el JDBC sink de Flink ho fa amb upsert); o el sink JDBC en mode exactly-once amb XA (dues fases, confirma amb el checkpoint), que afegeix la latència del checkpoint a cada alerta. Per a alertes, la primera opció és la correcta.
Conclusió
Processar un flux és executar de manera permanent el mateix DAG d'operadors que un lot, sobre una entrada que no acaba, i això obliga a respondre preguntes que el lot no tenia: quan és completa una finestra (la marca d'aigua, derivada del temps d'esdeveniment i d'un retard tolerat que es mesura), què fer amb el que arriba després (descartar, reemetre o desviar els tardans), com agrupar (finestres tumbling, sliding i session per clau), com fer durador un estat que viu per sempre (checkpoints consistents amb els offsets, savepoints per desplegar) i com no comptar dues vegades després d'una fallada (exactly-once a l'estat per checkpoint, i al sink per idempotència o transacció). Backpressure protegeix el motor deixant que Kafka absorbeixi els pics, i el lag és el senyal. Lambda i Kappa són dues maneres de conviure amb els lots: la primera duplica la lògica, la segona la unifica, i els motors moderns fan que la tria sigui de desplegament. A Quilòmetre Zero, simulacions/finestra_tumbling.py va fer visible el mecanisme, PyFlink va resoldre l'alerta d'estoc baix amb tres sentències SQL i un sink upsert-kafka, la session window detecta repartidors sense senyal, i Spark Structured Streaming va demostrar que el codi de 05-03 serveix gairebé sense canvis amb withWatermark i window.
Amb aquesta lliçó, analitica té les seves dues meitats: els lots de vendes_diaries.py i ALS, i els fluxos d'alertes i repartiment. Però els lots no es llancen sols. Algú ha d'esperar que el fitxer del dia sigui complet a HDFS, validar-lo, llançar spark-submit, carregar el resultat a la base de dades d'analitica, avisar si alguna cosa falla i reintentar, i fer-ho cada dia, i per als set dies de la Setmana de la Verema quan es va corregir el preu del vi. Fins ara això era un cron i una cadena d'scripts. L'última lliçó del mòdul tracta la planificació de treballs i els pipelines de dades: com expressar aquestes dependències com un DAG de tasques a Airflow, amb sensors, reintents, backfill i alertes, perquè la plataforma de dades de Quilòmetre Zero funcioni sense que ningú no la llanci a mà.
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
