Fins ara, tota la comunicació de Quilòmetre Zero ha estat síncrona: comandes crida inventari per gRPC, espera i rep. És el correcte quan la resposta cal en aquell instant, i és l'incorrecte per a tota la resta. Recordem el segon símptoma del monòlit (lliçó 01-06): la passarel·la de pagaments externa es va tornar lenta i, com que el cobrament era dins de la mateixa transacció que la reserva d'estoc i la creació de la comanda, cada comanda bloquejada retenia connexions i fils fins a tombar la plataforma sencera. Convertir aquesta cadena en crides gRPC no ho arregla: continuaria sent una cadena de disponibilitats multiplicades i latències sumades. El que cal és que comandes pugui dir "ha passat això" i tirar endavant, i que qui hagi de reaccionar ho faci quan pugui.
Aquesta és la funció de la missatgeria: un intermediari (el broker) que desa els missatges i els lliura als seus destinataris, desacoblant qui envia de qui rep en el temps, en l'espai i en el ritme. En aquesta lliçó fixarem els conceptes (productor, consumidor, cua, tòpic, ack, persistència), els dos models (punt a punt i publicació/subscripció), i les dues tecnologies que dominen el sector, RabbitMQ i Apache Kafka, amb les seves diferències de fons. Construirem el primer flux asíncron de Quilòmetre Zero: comandes publica l'esdeveniment comanda.creada, i inventari i analitica el consumeixen cadascun al seu ritme. Les garanties de lliurament i els patrons que les fan segures (idempotència, outbox, cues de missatges morts) s'anomenen aquí i es desenvolupen a la lliçó 02-05.
Contingut
- Per què desacoblar en el temps
- Vocabulari de la missatgeria
- Dos models: punt a punt i publicació/subscripció
- RabbitMQ i AMQP: exchanges, cues, bindings i routing keys
comanda.creadaamb RabbitMQ ipika- Apache Kafka: el log distribuït
comanda.creadaamb Kafka- RabbitMQ contra Kafka: criteris de tria
- Ampliar el
docker-compose.ymldekm0/ - Errors comuns i consells
- Exercicis
- Conclusió
- Per què desacoblar en el temps
Una crida síncrona acobla els dos participants de tres maneres:
- En el temps: tots dos han de ser vius i disponibles en el mateix instant. Si
analiticas'està reiniciant quancomandesintenta notificar-li una venda, la notificació falla ocomandesespera. - En l'espai: l'emissor necessita saber qui és el receptor i on és (adreça, port). Si demà
repartimenttambé vol saber de les comandes, cal canviarcomandes. - En el ritme: l'emissor no pot anar més ràpid que el receptor més lent. A la "Setmana de la Verema",
comandesprodueix 1.200 comandes per segon; sianaliticanomés en processa 300,comandess'alenteix fins a 300.
Un broker trenca els tres acoblaments: comandes lliura el missatge al broker (que sí que està disponible, perquè és infraestructura replicada) i continua; el broker el desa; els consumidors el recullen quan volen i al ritme que poden; i afegir un consumidor nou no toca el productor.
flowchart LR
subgraph Abans["Síncron: cadena de dependències"]
P1[comandes] --> I1[inventari]
P1 --> PA1[pagaments]
P1 --> A1[analitica]
P1 --> R1[repartiment]
end
subgraph Despres["Asíncron: el broker al mig"]
P2[comandes] -- comanda.creada --> B[(broker)]
B --> I2[inventari]
B --> A2[analitica]
B --> R2[repartiment]
end
El cost és igual de clar i cal acceptar-lo amb els ulls oberts: comandes ja no sap si el missatge s'ha processat, ni quan. Les decisions que requereixen resposta immediata (hi ha estoc?) continuen sent síncrones; l'asíncron és per a les conseqüències d'una decisió ja presa. I apareix una peça d'infraestructura crítica més, amb la seva pròpia disponibilitat i el seu propi model de fallades.
- Vocabulari de la missatgeria
| Terme | Què és | A Quilòmetre Zero |
|---|---|---|
| Missatge | Unitat de dades que viatja: capçaleres (metadades) + cos (bytes serialitzats, 02-03) | Un esdeveniment comanda.creada amb la comanda de l'Anna al cos |
| Productor (producer, publisher) | Qui envia missatges al broker | comandes |
| Consumidor (consumer, subscriber) | Qui rep missatges del broker | inventari, analitica, repartiment |
| Broker | El servidor intermediari que rep, emmagatzema i lliura | RabbitMQ o Kafka |
| Cua (queue) | Magatzem ordenat del qual els consumidors extreuen missatges; cada missatge es lliura normalment a un consumidor | inventari.comandes a RabbitMQ |
| Tòpic (topic) | Canal amb nom al qual es publiquen missatges i al qual se subscriuen els interessats; cada missatge pot arribar a diversos | comandes.esdeveniments a Kafka |
| Ack (acknowledgement) | Confirmació del consumidor al broker que ha processat el missatge; fins llavors el broker el conserva | basic_ack a RabbitMQ, commit d'offset a Kafka |
| Persistència | El broker escriu el missatge al disc, no només en memòria | Sobreviu a un reinici del broker |
| Durabilitat | La cua o el tòpic sobreviu al reinici del broker (a més dels seus missatges, si són persistents) | Les cues d'inventari han de ser durables |
| Retenció | Quant de temps (o quants bytes) conserva el broker els missatges | RabbitMQ: fins a l'ack; Kafka: dies, encara que ja s'hagin llegit |
| Prefetch | Quants missatges pot tenir un consumidor pendents d'ack alhora | 1 per a processament segur, més per a rendiment |
- Dos models: punt a punt i publicació/subscripció
Punt a punt (cua de treball): els productors deixen missatges en una cua i un o diversos consumidors els extreuen; cada missatge el processa exactament un consumidor. És el model per repartir feina: generar les factures en PDF de les comandes, enviar correus, calcular rutes. Afegir consumidors augmenta el ritme de processament (consumidors competidors, que tractarem a 02-05).
Publicació/subscripció (pub/sub): els productors publiquen en un tòpic sense saber qui escolta; cada subscriptor rep la seva pròpia còpia de cada missatge. És el model per a esdeveniments: "ha passat una comanda" interessa a inventari, a analitica i a repartiment, i cadascun en fa una cosa diferent.
flowchart LR
subgraph PP["Punt a punt"]
Pr1[productor] --> Q[(cua factures)]
Q --> C1[consumidor A]
Q --> C2[consumidor B]
note1[cada missatge va a UN dels dos]
end
subgraph PS["Publicació/subscripció"]
Pr2[comandes] --> T[(tòpic comanda.creada)]
T --> S1[inventari]
T --> S2[analitica]
T --> S3[repartiment]
note2[cada missatge va a TOTS]
end
A la pràctica els dos models es combinen: cada subscriptor d'un tòpic sol ser un grup d'instàncies que competeixen entre elles pels missatges d'aquest subscriptor (tres rèpliques d'analitica es reparteixen els esdeveniments, però analitica com a conjunt els rep tots). Tant RabbitMQ com Kafka suporten aquesta combinació, amb mecanismes diferents.
- RabbitMQ i AMQP: exchanges, cues, bindings i routing keys
RabbitMQ és un broker que implementa AMQP 0-9-1 (Advanced Message Queuing Protocol), un protocol binari sobre TCP amb un model d'encaminament molt flexible. La clau per entendre'l és que els productors mai no publiquen en cues: publiquen en un exchange, i l'exchange decideix a quines cues copiar el missatge segons unes regles (bindings) i una etiqueta del missatge (routing key).
flowchart LR
P[comandes] -- "routing key: comanda.creada" --> X{{exchange km0.comandes<br/>tipus topic}}
X -- "binding: comanda.creada" --> Q1[(cua inventari.comandes)]
X -- "binding: comanda.#" --> Q2[(cua analitica.comandes)]
X -- "binding: comanda.pagada" --> Q3[(cua repartiment.comandes)]
Q1 --> I[inventari]
Q2 --> A[analitica]
Q3 --> R[repartiment]
Els quatre tipus d'exchange:
| Tipus | Regla d'encaminament | Ús |
|---|---|---|
direct |
La routing key del missatge ha de ser igual a la del binding | Cues de treball amb destí explícit |
fanout |
Ignora la routing key: copia a totes les cues enllaçades | Difusió pura |
topic |
Routing key amb punts (comanda.creada); els bindings fan servir comodins: * (una paraula), # (zero o més) |
Esdeveniments classificats: comanda.*, repartiment.furgoneta-3.# |
headers |
Encamina per capçaleres del missatge en lloc de per routing key | Casos especials |
Altres conceptes de RabbitMQ: cada consumidor manté una connexió TCP i dins d'ella un o més canals (multiplexació lleugera); els missatges es lliuren al consumidor mitjançant push i queden "sense confirmar" fins a l'ack; si el consumidor mor sense confirmar, el broker reencua el missatge per a un altre consumidor; i un cop confirmat, el missatge desapareix de la cua. Aquesta última propietat és la diferència de fons amb Kafka.
comanda.creada amb RabbitMQ i pika
comanda.creada amb RabbitMQ i pikapika és el client Python oficial de RabbitMQ (pip install pika). Primer, el productor a comandes:
# km0/serveis/comandes/publicador_rabbit.py
import json
import time
import uuid
import pika
EXCHANGE = "km0.comandes"
def connectar():
connexio = pika.BlockingConnection(pika.ConnectionParameters("localhost"))
canal = connexio.channel()
# Declarar és idempotent: si l'exchange existeix amb els mateixos paràmetres, no passa res.
# durable=True: l'exchange sobreviu a reinicis del broker.
canal.exchange_declare(exchange=EXCHANGE, exchange_type="topic", durable=True)
return connexio, canal
def publicar_comanda_creada(canal, comanda):
esdeveniment = {
"id_esdeveniment": str(uuid.uuid4()), # identitat del missatge (clau a 02-05)
"tipus": "comanda.creada",
"versio": 1, # versió de l'esquema de l'esdeveniment
"data_ms": int(time.time() * 1000),
"origen": "comandes",
"dades": comanda,
}
canal.basic_publish(
exchange=EXCHANGE,
routing_key="comanda.creada",
body=json.dumps(esdeveniment).encode("utf-8"),
properties=pika.BasicProperties(
content_type="application/json",
message_id=esdeveniment["id_esdeveniment"],
delivery_mode=2, # 2 = persistent: s'escriu al disc
),
)
print(f"[comandes] publicat comanda.creada {comanda['id']}")
if __name__ == "__main__":
connexio, canal = connectar()
comanda_anna = {
"id": "P-2026-000123", "client": "anna", "mercat": "girona",
"linies": [{"producte": "formatge-curat", "quantitat": 2, "preu_centims": 1450},
{"producte": "tomaquet-rosa", "quantitat": 3, "preu_centims": 320}],
}
publicar_comanda_creada(canal, comanda_anna)
connexio.close()Fixa't en l'envoltura de l'esdeveniment: a més de les dades de la comanda porta un identificador únic, un tipus, una versió, una data i un origen. És un conveni que mantindrem en tots els esdeveniments de Quilòmetre Zero, i cada camp té un ús que apareixerà a 02-05 (l'id_esdeveniment per detectar duplicats, la versio per evolucionar l'esquema).
Ara el consumidor d'inventari. La seva feina, en aquesta primera versió, és descomptar l'estoc reservat (més endavant decidirem si la reserva síncrona per gRPC i el descompte per esdeveniment conviuen o se substitueixen; de moment, l'important és el mecanisme):
# km0/serveis/inventari/consumidor_rabbit.py
import json
import pika
EXCHANGE = "km0.comandes"
CUA = "inventari.comandes"
estoc = {"tomaquet-rosa": 120, "carbasso": 80, "formatge-curat": 5,
"formatge-fresc": 30, "vi-crianca": 200}
def en_rebre(canal, metode, propietats, cos):
esdeveniment = json.loads(cos)
comanda = esdeveniment["dades"]
for linia in comanda["linies"]:
estoc[linia["producte"]] -= linia["quantitat"]
print(f"[inventari] processat {esdeveniment['tipus']} {comanda['id']} "
f"(esdeveniment {esdeveniment['id_esdeveniment'][:8]}); formatge-curat={estoc['formatge-curat']}")
# Ack DESPRÉS de processar: si morim abans, el broker reencua el missatge
canal.basic_ack(delivery_tag=metode.delivery_tag)
connexio = pika.BlockingConnection(pika.ConnectionParameters("localhost"))
canal = connexio.channel()
canal.exchange_declare(exchange=EXCHANGE, exchange_type="topic", durable=True)
canal.queue_declare(queue=CUA, durable=True) # la cua sobreviu al reinici
canal.queue_bind(queue=CUA, exchange=EXCHANGE, routing_key="comanda.creada")
canal.basic_qos(prefetch_count=1) # un missatge sense confirmar alhora per consumidor
canal.basic_consume(queue=CUA, on_message_callback=en_rebre)
print("[inventari] esperant comanda.creada a", CUA)
canal.start_consuming()I el d'analitica, que vol tots els esdeveniments de comanda, no només els de creació, per acumular estadístiques:
# km0/serveis/analitica/consumidor_rabbit.py
import json
from collections import Counter
import pika
vendes_per_producte = Counter()
def en_rebre(canal, metode, propietats, cos):
esdeveniment = json.loads(cos)
if esdeveniment["tipus"] == "comanda.creada":
for linia in esdeveniment["dades"]["linies"]:
vendes_per_producte[linia["producte"]] += linia["quantitat"]
print(f"[analitica] {esdeveniment['tipus']} -> {dict(vendes_per_producte)}")
canal.basic_ack(delivery_tag=metode.delivery_tag)
connexio = pika.BlockingConnection(pika.ConnectionParameters("localhost"))
canal = connexio.channel()
canal.exchange_declare(exchange="km0.comandes", exchange_type="topic", durable=True)
canal.queue_declare(queue="analitica.comandes", durable=True)
canal.queue_bind(queue="analitica.comandes", exchange="km0.comandes", routing_key="comanda.#")
canal.basic_qos(prefetch_count=50) # analitica tolera lots: més rendiment
canal.basic_consume(queue="analitica.comandes", on_message_callback=en_rebre)
canal.start_consuming()Punts importants dels tres programes:
- Cada consumidor declara la seva pròpia cua i l'enllaça a l'exchange amb el patró que li interessa.
inventarinomés volcomanda.creada;analiticavolcomanda.#(creada, pagada, cancel·lada...). El productor no sap res de cap de les dues cues: això és el desacoblament en l'espai. - El missatge es copia a cada cua enllaçada:
inventariianaliticareben cadascun el seu exemplar (pub/sub). Si engegues dues instàncies delconsumidor_rabbit.pyd'inventari, compartiran la cuainventari.comandesi cada missatge anirà a una d'elles (punt a punt dins del subscriptor). - Engega els consumidors després de publicar i veuràs que reben el missatge igualment: estava desat a la cua. Però si la cua no existia quan es va publicar (perquè el consumidor no havia arrencat mai), el missatge es va descartar: l'exchange no tenia on copiar-lo. Per això les cues dels consumidors crítics s'han de declarar al desplegament, no a l'arrencada del primer consumidor.
durable=True+delivery_mode=2són les dues meitats de la supervivència a un reinici del broker: la cua i el missatge. L'una sense l'altra no serveix.prefetch_count=1ainventari: el broker no li envia el missatge següent fins que confirmi l'actual, cosa que evita que una instància acapari missatges que no pot processar si mor. Costa rendiment;analitica, que es pot permetre reprocessar lots, en fa servir 50.
Amb la interfície web d'administració (http://localhost:15672, usuari i contrasenya km0) es veuen exchanges, cues, missatges pendents i consumidors connectats: és la primera eina de diagnòstic.
- Apache Kafka: el log distribuït
Kafka va néixer a LinkedIn (2011) amb una idea diferent: en lloc d'una cua de la qual els missatges desapareixen en consumir-se, un log, és a dir, un fitxer de només afegir (append-only) en què cada missatge ocupa una posició fixa (offset) i hi roman durant un temps de retenció (per defecte 7 dies) independentment que algú l'hagi llegit. Els consumidors no extreuen missatges: llegeixen el log des d'una posició i recorden per on van.
flowchart LR
subgraph T["tòpic comandes.esdeveniments (3 particions)"]
P0["partició 0: [0][1][2][3][4] →"]
P1["partició 1: [0][1][2] →"]
P2["partició 2: [0][1][2][3] →"]
end
Pr[comandes<br/>clau = id de la comanda] --> T
subgraph G1["grup inventari"]
C1[instància 1] -.- P0
C1 -.- P1
C2[instància 2] -.- P2
end
subgraph G2["grup analitica"]
C3[instància única] -.- P0
C3 -.- P1
C3 -.- P2
end
Els conceptes que cal dominar:
- Tòpic: el nom lògic del flux (
comandes.esdeveniments). - Partició: cada tòpic es divideix en N logs independents. És la unitat de paral·lelisme (cada partició la llegeix una sola instància de cada grup) i d'ordre (els missatges estan ordenats dins d'una partició, no entre particions).
- Clau de partició: el productor pot assignar una clau a cada missatge; Kafka calcula
hash(clau) mod Nper triar partició. Tots els missatges amb la mateixa clau van a la mateixa partició i, per tant, es llegeixen en ordre. Fent servir l'identificador de la comanda com a clau,comanda.creada,comanda.pagadaicomanda.cancelladadeP-2026-000123arriben sempre en aquest ordre al mateix consumidor. Sense clau, es reparteixen i l'ordre es perd. (El hashing per repartir claus entre particions és el mateix problema del particionament de dades, lliçó 04-01.) - Offset: posició d'un missatge dins de la seva partició. És un enter creixent; els consumidors el fan servir com a marcador.
- Grup de consumidors: instàncies que comparteixen un
group.ides reparteixen les particions del tòpic (punt a punt dins del grup); grups diferents llegeixen el tòpic complet cadascun (pub/sub entre grups). El commit de l'offset és l'equivalent de l'ack: "el grupinventariha processat fins a l'offset 4 de la partició 0". - Retenció: els missatges s'esborren per edat o per mida, no per consum. Un consumidor nou pot llegir l'històric complet (
auto.offset.reset=earliest); un que ha estat caigut tres dies recupera el que s'ha perdut; ianaliticapot rellegir tot el mes si canvia la seva lògica. És una capacitat que RabbitMQ no té. - Rèpliques: cada partició es replica en diversos brokers (factor de replicació 3 en producció); un és el líder i atén lectures i escriptures. La replicació i les seves garanties s'estudien a 03-04.
comanda.creada amb Kafka
comanda.creada amb KafkaFarem servir confluent-kafka (pip install confluent-kafka), el client Python més eficient (embolcalla la biblioteca C librdkafka); kafka-python és una alternativa en Python pur amb una API semblant. Primer creem el tòpic amb 6 particions (amb un sol broker de desenvolupament, factor de replicació 1):
docker compose exec kafka /opt/kafka/bin/kafka-topics.sh --bootstrap-server localhost:9092 \
--create --topic comandes.esdeveniments --partitions 6 --replication-factor 1El productor:
# km0/serveis/comandes/publicador_kafka.py
import json
import time
import uuid
from confluent_kafka import Producer
TOPIC = "comandes.esdeveniments"
productor = Producer({
"bootstrap.servers": "localhost:9092",
"acks": "all", # el líder i les rèpliques sincronitzades confirmen abans de respondre
"client.id": "comandes",
})
def en_confirmar(err, msg):
"""Callback asíncron: Kafka acumula missatges en lots i confirma després."""
if err is not None:
print(f"[comandes] ERROR en publicar: {err}")
else:
print(f"[comandes] confirmat a {msg.topic()}[{msg.partition()}] offset {msg.offset()}")
def publicar_comanda_creada(comanda):
esdeveniment = {"id_esdeveniment": str(uuid.uuid4()), "tipus": "comanda.creada", "versio": 1,
"data_ms": int(time.time() * 1000), "origen": "comandes", "dades": comanda}
productor.produce(
TOPIC,
key=comanda["id"], # mateixa clau -> mateixa partició -> ordre
value=json.dumps(esdeveniment).encode("utf-8"),
headers=[("tipus", "comanda.creada"), ("id_esdeveniment", esdeveniment["id_esdeveniment"])],
callback=en_confirmar,
)
productor.poll(0) # atén callbacks pendents sense bloquejar
if __name__ == "__main__":
for i, (client, producte, quantitat) in enumerate(
[("anna", "formatge-curat", 2), ("marc", "vi-crianca", 6), ("llucia", "tomaquet-rosa", 3)]):
publicar_comanda_creada({"id": f"P-2026-00012{3 + i}", "client": client, "mercat": "girona",
"linies": [{"producte": producte, "quantitat": quantitat}]})
productor.flush() # espera que tot estigui confirmat abans de sortirproduce no envia: encua el missatge en un buffer intern que s'envia en lots per un fil en segon pla (per això Kafka aconsegueix centenars de milers de missatges per segon). flush() bloqueja fins que tots estan confirmats; acks=all fa que "confirmat" vulgui dir escrit al líder i a les rèpliques sincronitzades. Sense flush en acabar el programa, els missatges del buffer es perdrien.
El consumidor d'inventari:
# km0/serveis/inventari/consumidor_kafka.py
import json
from confluent_kafka import Consumer, KafkaError
consumidor = Consumer({
"bootstrap.servers": "localhost:9092",
"group.id": "inventari", # totes les instàncies d'inventari comparteixen grup
"auto.offset.reset": "earliest", # un grup nou comença pel principi del log
"enable.auto.commit": False, # confirmarem nosaltres, després de processar
})
consumidor.subscribe(["comandes.esdeveniments"])
estoc = {"tomaquet-rosa": 120, "carbasso": 80, "formatge-curat": 5, "formatge-fresc": 30, "vi-crianca": 200}
try:
while True:
msg = consumidor.poll(timeout=1.0) # pull: el consumidor demana, el broker no empeny
if msg is None:
continue
if msg.error():
if msg.error().code() != KafkaError._PARTITION_EOF:
print("[inventari] error:", msg.error())
continue
esdeveniment = json.loads(msg.value())
if esdeveniment["tipus"] == "comanda.creada":
for linia in esdeveniment["dades"]["linies"]:
estoc[linia["producte"]] -= linia["quantitat"]
print(f"[inventari] partició {msg.partition()} offset {msg.offset()} "
f"clau {msg.key().decode()}: formatge-curat={estoc['formatge-curat']}")
consumidor.commit(message=msg) # equivalent a l'ack: "fins aquí processat"
finally:
consumidor.close()Per a analitica, el mateix codi amb "group.id": "analitica" i la seva pròpia lògica: com que és un altre grup, rep tots els missatges de nou, amb independència del que hagi confirmat inventari. Prova d'engegar dues instàncies del consumidor d'inventari: veuràs als logs com Kafka reassigna les 6 particions (3 i 3) entre elles, i com cada comanda va sempre a la mateixa instància per la seva clau. Atura'n una i veuràs les particions tornar a l'altra: és el rebalanceig de grups, que dona tolerància a fallades als consumidors sense que el productor se n'assabenti.
Un detall sobre el commit: commit(message=msg) marca fins a aquest offset en aquesta partició, no "aquest missatge". Si es processen els missatges 3, 4 i 5 i es confirma el 5, en reiniciar es continua des del 6. Si es confirma el 5 sense haver processat el 4 (per exemple, processant en paral·lel), el 4 es perd. El commit és un marcador de posició, i les conseqüències de confirmar abans o després de processar (perdre o duplicar) són exactament les garanties de lliurament que analitza 02-05.
- RabbitMQ contra Kafka: criteris de tria
| Criteri | RabbitMQ | Apache Kafka |
|---|---|---|
| Model | Cua intel·ligent: el broker encamina i fa seguiment de cada missatge | Log ximple: el broker desa; el consumidor recorda per on va |
| El missatge després de consumir-se | S'esborra | Roman fins a la retenció |
| Encaminament | Molt flexible (exchanges, comodins, capçaleres) | Per tòpic i partició; el filtratge el fa el consumidor |
| Ordre | Per cua, amb un consumidor | Per partició (per clau) |
| Lliurament | Push al consumidor | Pull del consumidor |
| Rendiment típic | Desenes de milers de missatges/s | Centenars de milers a milions de missatges/s |
| Rellegir l'històric | No | Sí (per disseny) |
| Consumidors lents | Acumulen missatges a la cua; pot degradar el broker | No afecten el broker ni altres grups |
| Prioritats, TTL, cues de morts | Natives | No natives (es construeixen amb tòpics, 02-05) |
| Petició/resposta | Còmode (reply_to, 02-05) |
Incòmode |
| Complexitat operativa | Baixa-mitjana | Mitjana-alta (particions, rebalanceigs, retenció), encara que KRaft l'ha reduïda |
| Protocol | AMQP (també MQTT i STOMP amb plugins) | Propi, binari sobre TCP |
| Ecosistema de dades | Poc | Enorme: Kafka Connect, Streams, integració amb Spark i Flink (Mòdul 5) |
| Paper a Quilòmetre Zero | Cues de treball (factures, correus), petició/resposta interna | Esdeveniments de domini (comandes.esdeveniments, inventari.esdeveniments, telemetria de repartiment) |
Criteris de decisió:
- Els missatges són esdeveniments que diversos sistemes voldran, potser en el futur, potser rellegits? Kafka. Aquesta és la raó per la qual l'arquitectura objectiu de 01-06 tria Kafka per als esdeveniments de domini:
analiticavoldrà reprocessar el mes; un servei nou voldrà l'històric. - Són tasques que cal executar una vegada i oblidar, amb encaminament fi, prioritats o resposta? RabbitMQ.
- Importa l'ordre per entitat (tots els esdeveniments d'una comanda en ordre)? Kafka amb clau de partició el dona de manera natural.
- Volum? Per sota d'uns pocs milers de missatges per segon, qualsevol dels dos; molt per sobre, Kafka.
- Qui l'operarà? Un broker mal operat és pitjor que cap. Si l'equip és petit, comença amb un de sol i amb un servei gestionat.
Quilòmetre Zero, seguint l'arquitectura objectiu, farà servir Kafka per als esdeveniments de domini i es reserva RabbitMQ per a cues de treball i per al patró petició/resposta que veurem a 02-05. Tenir tots dos no és obligatori; molts sistemes viuen bé amb un de sol.
- Ampliar el
docker-compose.yml de km0/
docker-compose.yml de km0/Afegim els dos brokers a l'entorn de 01-06. Kafka en mode KRaft (sense ZooKeeper, el mode estàndard des de la versió 3.x), amb dos listeners: un d'intern per als serveis de la xarxa de Compose i un altre per als scripts que executem des del host.
rabbitmq:
image: rabbitmq:3.13-management
environment:
RABBITMQ_DEFAULT_USER: km0
RABBITMQ_DEFAULT_PASS: km0_dev
ports:
- "5672:5672" # AMQP
- "15672:15672" # interfície web d'administració
volumes:
- rabbitmq_data:/var/lib/rabbitmq
healthcheck:
test: ["CMD", "rabbitmq-diagnostics", "-q", "ping"]
interval: 10s
timeout: 5s
retries: 5
kafka:
image: apache/kafka:3.8.0
environment:
KAFKA_NODE_ID: 1
KAFKA_PROCESS_ROLES: broker,controller
KAFKA_LISTENERS: INTERN://:19092,EXTERN://:9092,CONTROLLER://:9093
KAFKA_ADVERTISED_LISTENERS: INTERN://kafka:19092,EXTERN://localhost:9092
KAFKA_LISTENER_SECURITY_PROTOCOL_MAP: CONTROLLER:PLAINTEXT,INTERN:PLAINTEXT,EXTERN:PLAINTEXT
KAFKA_INTER_BROKER_LISTENER_NAME: INTERN
KAFKA_CONTROLLER_LISTENER_NAMES: CONTROLLER
KAFKA_CONTROLLER_QUORUM_VOTERS: 1@kafka:9093
KAFKA_OFFSETS_TOPIC_REPLICATION_FACTOR: 1 # un sol broker en desenvolupament
KAFKA_TRANSACTION_STATE_LOG_REPLICATION_FACTOR: 1
KAFKA_TRANSACTION_STATE_LOG_MIN_ISR: 1
KAFKA_LOG_RETENTION_HOURS: 168 # 7 dies
ports:
- "9092:9092"
volumes:
- kafka_data:/var/lib/kafka/data
volumes:
postgres_data:
rabbitmq_data:
kafka_data:Els serveis que corrin dins de Compose (com inventari a partir de 02-03) faran servir kafka:19092 i rabbitmq:5672; els scripts llançats des del host, localhost:9092 i localhost:5672. Les contrasenyes en clar són només per a desenvolupament: la gestió de secrets és tema de 06-04. Amb docker compose up -d rabbitmq kafka i els tres scripts d'aquest capítol, tens el primer flux asíncron de Quilòmetre Zero funcionant.
Errors Comuns i Consells
- Fer servir missatgeria per al que necessita resposta immediata. "Hi ha estoc?" no pot esperar que un consumidor processi un esdeveniment. Síncron per a decisions, asíncron per a conseqüències.
- Cues no durables o missatges no persistents en fluxos crítics. Un reinici del broker s'emporta les comandes. Durable + persistent, sempre, llevat de la telemetria d'un sol ús.
- Confirmar (ack/commit) abans de processar. Si el consumidor mor després de confirmar i abans d'acabar, el missatge es perd. Confirmar després de processar duplica en el cas contrari, i 02-05 explica com conviure-hi.
- Oblidar
flush()al productor de Kafka.producenomés encua; un procés que acaba senseflushperd el buffer, sense error. - Publicar sense clau de partició i esperar ordre. Sense clau, els esdeveniments d'una mateixa comanda es reparteixen entre particions i
comanda.pagadaes pot processar abans quecomanda.creada. - Un tòpic per a cada tipus d'esdeveniment a Kafka. Trenca l'ordre entre esdeveniments de la mateixa entitat (són en tòpics diferents) i multiplica les particions. Un tòpic per entitat (
comandes.esdeveniments) amb el tipus en una capçalera sol ser millor. - Consumidors que declaren cues que "haurien d'existir". Si
inventarino ha arrencat mai, la seva cua de RabbitMQ no existeix i els esdeveniments publicats entremig es perden. Declara la topologia (exchanges, cues, bindings, tòpics) al desplegament. - Missatges enormes. Ni RabbitMQ ni Kafka estan fets per transportar fotos de productes de 5 MB. Al missatge hi va una referència al magatzem d'objectes (04-03).
- Consell: posa sempre una envoltura estàndard als esdeveniments (
id_esdeveniment,tipus,versio,data_ms,origen,dades). Costa cinc línies i ho agrairàs a cada patró de 02-05 i a cada investigació de 07-02. - Consell: mira la interfície de RabbitMQ i els lag dels grups de Kafka (
kafka-consumer-groups.sh --describe) des del primer dia. Una cua que creix o un lag que no baixa són el primer símptoma d'un consumidor malalt, molt abans que ningú no es queixi.
Exercicis
Exercici 1: Encaminament amb topic exchange
repartiment vol rebre només les comandes ja pagades (comanda.pagada) i atencio-client vol tots els esdeveniments de cancel·lació de qualsevol entitat (comanda.cancellada, repartiment.cancellat, ...). Escriu les declaracions de cua i binding de tots dos sobre l'exchange km0.comandes (assumeix que els esdeveniments de repartiment també s'hi publiquen amb routing keys repartiment.<tipus>), i indica què rebria cadascun si es publiquessin, en ordre: comanda.creada, comanda.pagada, repartiment.assignat, repartiment.cancellat, comanda.cancellada.
Exercici 2: Particions, claus i ordre
El tòpic comandes.esdeveniments té 6 particions i inventari té 3 instàncies al mateix grup. (a) Quantes particions llegeix cada instància? (b) Si s'hi afegeixen 4 instàncies més (7 en total), què passa? (c) comandes publica comanda.creada i comanda.cancellada per a P-2026-000123 amb clau P-2026-000123, i comanda.creada per a P-2026-000124. Està garantit que inventari processa la creació de la 123 abans que la seva cancel·lació? I abans que la creació de la 124? (d) Un desenvolupador proposa fer servir client com a clau en lloc de l'id de la comanda, "perquè totes les comandes de l'Anna vagin al mateix consumidor". Què hi guanya i què hi arrisca?
Exercici 3: Triar broker
Per a cada flux de Quilòmetre Zero, indica RabbitMQ o Kafka i justifica-ho amb la taula de l'apartat 8:
- Generar el PDF de la factura de cada comanda pagada (tasca pesada, una vegada, sense ordre).
- Els esdeveniments d'estoc (
inventari.esdeveniments) quecatalegconsumeix per a l'indicador "en queden poques unitats" i queanaliticavol reprocessar cada mes per estudiar trencaments d'estoc. - Posicions de 400 repartidors, una per segon cadascun, per al mapa en viu i per a l'anàlisi posterior de rutes.
- Un servei intern que necessita preguntar a
pagamentssi una targeta està bloquejada, i esperar-ne la resposta.
Solucions
Solució 1:
# repartiment: només comandes pagades
canal.queue_declare(queue="repartiment.comandes", durable=True)
canal.queue_bind(queue="repartiment.comandes", exchange="km0.comandes", routing_key="comanda.pagada")
# atenció al client: qualsevol cancel·lació de qualsevol entitat
canal.queue_declare(queue="atencio.cancellacions", durable=True)
canal.queue_bind(queue="atencio.cancellacions", exchange="km0.comandes", routing_key="*.cancellat")
canal.queue_bind(queue="atencio.cancellacions", exchange="km0.comandes", routing_key="*.cancellada")Amb els cinc esdeveniments publicats: repartiment.comandes rep només comanda.pagada. atencio.cancellacions rep repartiment.cancellat i comanda.cancellada (el comodí * casa exactament una paraula: comanda o repartiment; com que el sufix concorda en gènere amb l'entitat, calen dos bindings, un per cancellat i un per cancellada). Cap de les dues no rep comanda.creada ni repartiment.assignat. I les cues ja existents continuen rebent el que és seu: inventari.comandes només comanda.creada; analitica.comandes (comanda.#) els tres esdeveniments de comanda, però no els de repartiment. Cada consumidor decideix què vol; el productor no canvia.
Solució 2:
(a) Kafka reparteix les 6 particions entre les 3 instàncies: 2 cadascuna. (b) Amb 7 instàncies i 6 particions, 6 instàncies llegeixen una partició cadascuna i la setena queda ociosa: el nombre de particions és el límit superior del paral·lelisme d'un grup. Per aprofitar més instàncies caldria crear el tòpic amb més particions (i no se'n poden afegir sense canviar l'assignació clau→partició dels missatges futurs, cosa que trenca temporalment l'ordre per clau). (c) Sí per a la 123: tots dos missatges tenen la mateixa clau, van a la mateixa partició, i una partició la llegeix una sola instància en ordre. No respecte de la 124: pot ser en una altra partició llegida per una altra instància; el seu ordre relatiu no està definit, i no importa, perquè són comandes independents. (d) Hi guanya que tots els esdeveniments d'un mateix client es processen en ordre i a la mateixa instància (útil si hi hagués regles per client, com un límit de comandes per hora). Hi arrisca dues coses: particions desequilibrades (un client empresarial que fa el 30 % de les comandes concentra el 30 % de la càrrega en una partició i una instància) i, en cert sentit, ordre innecessari: les comandes de l'Anna són independents entre elles, així que serialitzar-les no aporta res i redueix el paral·lelisme. La clau ha de ser l'entitat l'ordre de la qual importa, ni més fina ni més gruixuda.
Solució 3:
- RabbitMQ: cua de treball clàssica; cada factura la genera una instància i desapareix; sense ordre ni relectura; amb prioritats si calgués (factures de campanya primer). Kafka també podria, però no hi aporta res.
- Kafka: esdeveniments de domini amb dos consumidors amb necessitats diferents (un en temps real, l'altre rellegint un mes), ordre per producte (clau = slug del producte) i retenció llarga. És exactament el cas d'ús per al qual es va dissenyar.
- Kafka (amb un pont MQTT al davant per al tram mòbil, 08-02): volum alt (400 missatges/s sostinguts, molt més en campanya), dos consumidors (mapa en viu i anàlisi de rutes), ordre per repartidor (clau =
furgoneta-3), i retenció per a l'anàlisi posterior, que és una feina per lots del Mòdul 5. RabbitMQ acumularia les posicions i no permetria rellegir-les. - RabbitMQ, si es decideix fer-ho per missatgeria: el patró petició/resposta amb
reply_toicorrelation_id(02-05) és natural a AMQP i incòmode a Kafka. Però la pregunta prèvia és si ha de ser missatgeria en absolut: és una consulta síncrona que necessita resposta immediata, i una crida gRPC apagaments(02-03) és més simple i més ràpida. La missatgeria només es justificaria sipagamentsfos lent o intermitent i es volguessin absorbir pics, o si la resposta pogués trigar i el sol·licitant pogués esperar sense bloquejar.
Conclusió
Aquesta lliçó ha introduït la segona meitat de la comunicació entre serveis. Una crida síncrona acobla en el temps, en l'espai i en el ritme; un broker trenca els tres acoblaments a canvi que l'emissor deixi de saber quan, i si, es processarà el seu missatge. Amb aquest vocabulari (productor, consumidor, cua, tòpic, ack, persistència, durabilitat, retenció) hem distingit el model punt a punt, per repartir feina, del de publicació/subscripció, per difondre esdeveniments, i hem vist com cada broker els combina. RabbitMQ és una cua intel·ligent amb encaminament flexible (exchanges, bindings, routing keys amb comodins) que esborra els missatges en confirmar-se; Kafka és un log distribuït en què els missatges romanen, els consumidors recorden el seu offset, les particions donen paral·lelisme i les claus donen ordre per entitat, i qualsevol grup pot rellegir l'històric. Quilòmetre Zero té ara el seu primer flux asíncron: comandes publica comanda.creada amb una envoltura estàndard, inventari i analitica el consumeixen al seu ritme, i tots dos brokers són al docker-compose.yml.
Però hem deixat diverses preguntes obertes a propòsit. Què passa si inventari mor després de descomptar l'estoc i abans de confirmar el missatge? El broker el reenviarà i l'estoc es descomptarà dues vegades, el mateix problema que vam deixar pendent a 01-04 i a 02-02. Què passa si comandes desa la comanda a PostgreSQL i cau abans de publicar l'esdeveniment, o publica l'esdeveniment i després falla el COMMIT? I amb un missatge malformat que fa fallar el consumidor una vegada i una altra, bloquejant tots els que venen darrere? Aquestes preguntes són les garanties de lliurament i els patrons que les fan manejables (consumidors idempotents, outbox transaccional, cues de missatges morts), i són el tema de l'última lliçó del mòdul: Patrons de Comunicació Asíncrona.
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
