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

  1. Per què desacoblar en el temps
  2. Vocabulari de la missatgeria
  3. Dos models: punt a punt i publicació/subscripció
  4. RabbitMQ i AMQP: exchanges, cues, bindings i routing keys
  5. comanda.creada amb RabbitMQ i pika
  6. Apache Kafka: el log distribuït
  7. comanda.creada amb Kafka
  8. RabbitMQ contra Kafka: criteris de tria
  9. Ampliar el docker-compose.yml de km0/
  10. Errors comuns i consells
  11. Exercicis
  12. Conclusió

  1. 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 analitica s'està reiniciant quan comandes intenta notificar-li una venda, la notificació falla o comandes espera.
  • En l'espai: l'emissor necessita saber qui és el receptor i on és (adreça, port). Si demà repartiment també vol saber de les comandes, cal canviar comandes.
  • En el ritme: l'emissor no pot anar més ràpid que el receptor més lent. A la "Setmana de la Verema", comandes produeix 1.200 comandes per segon; si analitica només en processa 300, comandes s'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.

  1. 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

  1. 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.

  1. 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.

  1. comanda.creada amb RabbitMQ i pika

pika é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. inventari només vol comanda.creada; analitica vol comanda.# (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: inventari i analitica reben cadascun el seu exemplar (pub/sub). Si engegues dues instàncies del consumidor_rabbit.py d'inventari, compartiran la cua inventari.comandes i 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=2 só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=1 a inventari: 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.

  1. 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 N per 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.pagada i comanda.cancellada de P-2026-000123 arriben 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.id es 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 grup inventari ha 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; i analitica pot 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.

  1. comanda.creada amb Kafka

Farem 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 1

El 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 sortir

produce 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.

  1. 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ó:

  1. 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: analitica voldrà reprocessar el mes; un servei nou voldrà l'històric.
  2. Són tasques que cal executar una vegada i oblidar, amb encaminament fi, prioritats o resposta? RabbitMQ.
  3. Importa l'ordre per entitat (tots els esdeveniments d'una comanda en ordre)? Kafka amb clau de partició el dona de manera natural.
  4. Volum? Per sota d'uns pocs milers de missatges per segon, qualsevol dels dos; molt per sobre, Kafka.
  5. 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.

  1. Ampliar el 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. produce només encua; un procés que acaba sense flush perd 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.pagada es pot processar abans que comanda.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 inventari no 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:

  1. Generar el PDF de la factura de cada comanda pagada (tasca pesada, una vegada, sense ordre).
  2. Els esdeveniments d'estoc (inventari.esdeveniments) que cataleg consumeix per a l'indicador "en queden poques unitats" i que analitica vol reprocessar cada mes per estudiar trencaments d'estoc.
  3. Posicions de 400 repartidors, una per segon cadascun, per al mapa en viu i per a l'anàlisi posterior de rutes.
  4. Un servei intern que necessita preguntar a pagaments si 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:

  1. 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.
  2. 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.
  3. 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.
  4. RabbitMQ, si es decideix fer-ho per missatgeria: el patró petició/resposta amb reply_to i correlation_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 a pagaments (02-03) és més simple i més ràpida. La missatgeria només es justificaria si pagaments fos 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

Mòdul 2: Comunicació en Sistemes Distribuïts

Mòdul 3: Consistència i Replicació

Mòdul 4: Emmagatzematge Distribuït

Mòdul 5: Computació Distribuïda

Mòdul 6: Seguretat en Sistemes Distribuïts

Mòdul 7: Monitoratge i Manteniment

Mòdul 8: Casos d'Estudi i Aplicacions

© Copyright 2026. Tots els drets reservats