La lliçó anterior va acabar amb tres preguntes incòmodes. Si inventari descompta l'estoc i mor abans de confirmar el missatge, el broker el reenvia i l'estoc es descompta dues vegades. Si comandes desa la comanda a PostgreSQL i cau abans de publicar comanda.creada, inventari i analitica no sabran mai que aquesta comanda existeix. I si un missatge malformat fa fallar el consumidor una vegada i una altra, tots els missatges que venen darrere es queden esperant. Cap de les tres no és un cas rar: són conseqüències directes de les fal·làcies de 01-04 aplicades a la missatgeria, i apareixeran en producció a la primera campanya.

Aquesta lliçó presenta els patrons amb què la indústria ha après a conviure-hi. Primer posarem nom al que un broker pot i no pot garantir (at-most-once, at-least-once i el mite d'exactly-once), i veurem com aquestes garanties depenen d'on es col·loca l'ack. Després tancarem, per fi, el problema dels duplicats que vam deixar obert a 01-04 i 02-02 amb consumidors idempotents recolzats en una taula a PostgreSQL. Resoldrem l'escriptura dual amb el patró Outbox transaccional i un relay, implementarem petició/resposta sobre missatgeria, escalarem amb consumidors competidors respectant l'ordre per partició, i tractarem els missatges enverinats amb reintents i cues de missatges morts. Acabarem amb el versionat d'esquemes d'esdeveniments i una taula que mapa cada patró al problema que resol. Les sagues (transaccions entre serveis) i el circuit breaker (resiliència de crides síncrones) queden per a 03-05 i 07-04.

Contingut

  1. Garanties de lliurament: at-most-once, at-least-once i el mite d'exactly-once
  2. Consumidors idempotents: tancar el problema dels duplicats
  3. El problema de l'escriptura dual i el patró Outbox transaccional
  4. Petició/resposta sobre missatgeria
  5. Consumidors competidors i ordre per partició
  6. Missatges enverinats: reintents amb backoff i cues de missatges morts
  7. Event-driven i event-sourcing: dues coses diferents
  8. Esquemes d'esdeveniments i versionat
  9. Taula de patrons: quin problema resol cadascun
  10. Errors comuns i consells
  11. Exercicis
  12. Conclusió

  1. Garanties de lliurament: at-most-once, at-least-once i el mite d'exactly-once

A 02-02 vam veure les semàntiques d'invocació d'RPC. Les garanties de lliurament de la missatgeria són la mateixa idea amb un broker al mig, i el factor que les decideix és quan es confirma (ack a RabbitMQ, commit d'offset a Kafka) respecte de quan es processa:

sequenceDiagram
    participant B as Broker
    participant C as Consumidor (inventari)
    participant DB as PostgreSQL
    Note over B,DB: Opció A: confirmar ABANS de processar → at-most-once
    B->>C: comanda.creada (offset 41)
    C->>B: commit(41)
    C-xDB: UPDATE estoc ... (el procés mor aquí)
    Note over B,DB: El broker creu que el 41 està fet: el missatge es PERD
    Note over B,DB: Opció B: confirmar DESPRÉS de processar → at-least-once
    B->>C: comanda.creada (offset 42)
    C->>DB: UPDATE estoc ... COMMIT
    C-xB: commit(42) (el procés mor aquí)
    Note over B,DB: El broker reenvia el 42: l'estoc es descompta DUES vegades
Garantia Com s'aconsegueix Risc Quan és acceptable
At-most-once Confirmar abans de processar (o no confirmar mai, amb auto_ack) Pèrdua de missatges Telemetria, mètriques, qualsevol dada que el missatge següent substitueix
At-least-once Confirmar després de processar; el productor reintenta fins a rebre confirmació del broker Duplicats Tot el que importi, sempre que el consumidor sigui idempotent
Exactly-once No existeix d'extrem a extrem S'emula amb at-least-once + idempotència

El mite, amb precisió. Kafka ofereix una funcionalitat anomenada exactly-once semantics (productors idempotents i transaccions de Kafka) que garanteix que un missatge no es duplica dins de Kafka: si el productor reintenta per un timeout, el broker deduplica; i una aplicació que llegeix d'un tòpic i escriu en un altre ho pot fer atòmicament. És valuós per a pipelines Kafka → Kafka (Mòdul 5). Però tan bon punt l'efecte del missatge surt de Kafka (un UPDATE a PostgreSQL, un correu, un cobrament), la finestra entre "efecte aplicat" i "offset confirmat" reapareix, i ningú no la pot tancar des de fora. El mateix passa amb el QoS 2 de MQTT, que garanteix un únic lliurament al client, no un únic efecte. La conclusió pràctica és la de 02-02: at-least-once al transport, idempotència al consumidor. Tota la resta d'aquesta lliçó es construeix sobre aquesta base.

  1. Consumidors idempotents: tancar el problema dels duplicats

Una operació és idempotent si executar-la N vegades produeix el mateix estat que executar-la una. Algunes ho són per naturalesa (SET estoc = 118, "marcar la comanda com a pagada", "inserir amb clau primària"); d'altres no (estoc = estoc - 2, "enviar un correu", "cobrar 14,50 €"). Per a les que no ho són, la tècnica és fer que el consumidor recordi quins missatges ja ha processat i descarti els repetits. Dos requisits:

  1. Una clau d'idempotència per missatge: un identificador únic generat pel productor i estable entre reintents. És l'id_esdeveniment de l'envoltura de 02-04, i era l'id_reserva de la petició gRPC de 02-03. Sense ella no hi ha manera de saber que dos missatges són "el mateix".
  2. Registrar la clau i aplicar l'efecte a la mateixa transacció. Si es registra la clau en una transacció i s'aplica l'efecte en una altra, reapareix la finestra: es registra, mor el procés, l'efecte no passa mai, i el reintent es descarta com a duplicat (pèrdua). Si és a l'inrevés, l'efecte s'aplica, mor, i el reintent l'aplica una altra vegada (duplicat). Amb tots dos a la mateixa transacció, o es fan els dos o cap.

inventari ja té la seva pròpia base de dades (un dels principis de 01-06), així que hi afegim la taula:

-- km0/sql/inventari/002_missatges_processats.sql
CREATE TABLE estoc (
    producte TEXT PRIMARY KEY,
    unitats  INTEGER NOT NULL CHECK (unitats >= 0)
);
INSERT INTO estoc VALUES ('tomaquet-rosa', 120), ('carbasso', 80), ('formatge-curat', 5),
                         ('formatge-fresc', 30), ('vi-crianca', 200);

CREATE TABLE missatges_processats (
    id_missatge  TEXT        NOT NULL,   -- id_esdeveniment de l'envoltura
    consumidor   TEXT        NOT NULL,   -- 'inventari.descompte_estoc'
    processat_en TIMESTAMPTZ NOT NULL DEFAULT now(),
    PRIMARY KEY (id_missatge, consumidor)
);

La clau primària composta permet que dues lògiques diferents dins d'inventari (descomptar estoc i, per exemple, avisar el productor) processin el mateix esdeveniment cadascuna una vegada. El consumidor de Kafka de 02-04, ara idempotent, amb psycopg (versió 3):

# km0/serveis/inventari/consumidor_idempotent.py
import json

import psycopg
from confluent_kafka import Consumer, KafkaError

CONSUMIDOR = "inventari.descompte_estoc"
DSN = "postgresql://km0:km0_dev@localhost:5432/km0_inventari"


class EstocInsuficient(Exception):
    pass


def processar_comanda_creada(conn, esdeveniment):
    """Aplica el descompte d'estoc UNA sola vegada per id_esdeveniment.
    Retorna 'processat' o 'duplicat'. Llança si l'efecte no es pot aplicar."""
    with conn.transaction():                              # BEGIN ... COMMIT/ROLLBACK
        with conn.cursor() as cur:
            # 1. Intentar registrar la clau. Si ja existeix, ON CONFLICT no insereix
            #    i RETURNING no retorna cap fila: és un duplicat.
            cur.execute(
                "INSERT INTO missatges_processats (id_missatge, consumidor) VALUES (%s, %s) "
                "ON CONFLICT DO NOTHING RETURNING id_missatge",
                (esdeveniment["id_esdeveniment"], CONSUMIDOR))
            if cur.fetchone() is None:
                return "duplicat"
            # 2. Aplicar l'efecte a LA MATEIXA transacció.
            for linia in esdeveniment["dades"]["linies"]:
                cur.execute(
                    "UPDATE estoc SET unitats = unitats - %s "
                    "WHERE producte = %s AND unitats >= %s",
                    (linia["quantitat"], linia["producte"], linia["quantitat"]))
                if cur.rowcount == 0:
                    # Desfà també el registre del pas 1: el missatge NO queda
                    # marcat com a processat i es podrà reintentar o anar a la DLQ (ap. 6)
                    raise EstocInsuficient(f"sense estoc de {linia['producte']}")
    return "processat"


def main():
    consumidor = Consumer({"bootstrap.servers": "localhost:9092", "group.id": "inventari",
                           "auto.offset.reset": "earliest", "enable.auto.commit": False})
    consumidor.subscribe(["comandes.esdeveniments"])
    with psycopg.connect(DSN) as conn:
        while True:
            msg = consumidor.poll(1.0)
            if msg is None:
                continue
            if msg.error():
                if msg.error().code() != KafkaError._PARTITION_EOF:
                    print("[inventari] error de Kafka:", msg.error())
                continue
            esdeveniment = json.loads(msg.value())
            if esdeveniment["tipus"] == "comanda.creada":
                resultat = processar_comanda_creada(conn, esdeveniment)
                print(f"[inventari] {esdeveniment['id_esdeveniment'][:8]} comanda {esdeveniment['dades']['id']}: {resultat}")
            consumidor.commit(message=msg)   # SEMPRE després de la transacció de BD


if __name__ == "__main__":
    main()

Recorrem les fallades possibles amb aquest codi:

  • Mor després del COMMIT de PostgreSQL i abans del commit de Kafka. Kafka reenvia el missatge; l'INSERT xoca amb la clau primària; ON CONFLICT DO NOTHING no insereix; fetchone() retorna None; es respon duplicat sense tocar l'estoc; es confirma l'offset. Sense duplicació.
  • Mor enmig de la transacció. PostgreSQL fa ROLLBACK: ni la clau ni el descompte no persisteixen. Kafka reenvia; es processa des de zero. Sense pèrdua.
  • Mor abans de començar. Kafka reenvia. Trivial.
  • El productor ha publicat el mateix esdeveniment dues vegades (ha reintentat per un timeout del broker). Mateix id_esdeveniment, mateix resultat: la segona còpia és un duplicat. És la situació de la simulació de 01-04, ara resolta: el reintent ja no converteix la pèrdua d'un missatge en duplicació d'efectes.

El mateix s'aplica a la crida síncrona de 02-03: ReservarEstoc pot registrar id_reserva a missatges_processats amb consumidor = 'inventari.reserva_grpc' i retornar la resposta desada si es repeteix, amb la qual cosa l'"estat desconegut" després d'un DEADLINE_EXCEEDED deixa de ser un problema: comandes pot reintentar amb el mateix id_reserva sense risc. És la semàntica at-most-once de 02-02 feta amb una taula. Dues consideracions operatives: la taula creix (una fila per missatge) i s'ha de purgar per antiguitat (més enllà de la retenció del tòpic ja no pot arribar cap reintent: DELETE ... WHERE processat_en < now() - interval '14 days'); i si el consumidor no té base de dades transaccional (per exemple, envia correus), la idempotència s'ha de recolzar en el sistema destí (un proveïdor de correu que accepti una clau d'idempotència) o acceptar at-most-once.

  1. El problema de l'escriptura dual i el patró Outbox transaccional

Ara el costat del productor. A 02-04, comandes feia dues coses en crear una comanda: escriure-la al seu PostgreSQL i publicar comanda.creada a Kafka. Són dos sistemes diferents i no hi ha cap transacció que els abasti (és l'escriptura dual, dual write):

  • Si escriu a la base de dades i cau abans de publicar, la comanda existeix però ningú no se n'assabenta: inventari no descompta, analitica no compta, repartiment no reparteix.
  • Si publica primer i després falla el COMMIT, tots reaccionen a una comanda que no existeix.
  • Si publica dins de la transacció "perquè es desfaci", no es desfà: Kafka no participa en el ROLLBACK de PostgreSQL.

El patró Outbox transaccional resol el problema convertint dues escriptures en una: l'esdeveniment s'escriu en una taula de la mateixa base de dades, a la mateixa transacció que la comanda. Un procés a part, el relay, llegeix aquesta taula i publica a Kafka.

sequenceDiagram
    participant A as App de l'Anna
    participant P as comandes
    participant DB as PostgreSQL (comandes)
    participant R as Relay outbox
    participant K as Kafka
    participant I as inventari
    A->>P: crear comanda
    P->>DB: BEGIN
    P->>DB: INSERT comandes, linies_comanda
    P->>DB: INSERT outbox (comanda.creada)
    P->>DB: COMMIT (atòmic: comanda + esdeveniment, o res)
    P-->>A: comanda P-2026-000123 creada
    loop cada 200 ms
        R->>DB: SELECT ... FROM outbox WHERE publicat_en IS NULL FOR UPDATE SKIP LOCKED
        R->>K: produce(comanda.creada)
        K-->>R: confirmat
        R->>DB: UPDATE outbox SET publicat_en = now()
    end
    K->>I: comanda.creada (at-least-once)

La taula, a la base de dades de comandes:

-- km0/sql/comandes/002_outbox.sql
CREATE TABLE outbox (
    id           UUID        PRIMARY KEY,           -- serà l'id_esdeveniment
    agregat      TEXT        NOT NULL,              -- 'comanda'
    agregat_id   TEXT        NOT NULL,              -- 'P-2026-000123': clau de partició
    tipus        TEXT        NOT NULL,              -- 'comanda.creada'
    versio       INTEGER     NOT NULL DEFAULT 1,
    carrega      JSONB       NOT NULL,
    creat_en     TIMESTAMPTZ NOT NULL DEFAULT now(),
    publicat_en  TIMESTAMPTZ
);
CREATE INDEX outbox_pendents ON outbox (creat_en) WHERE publicat_en IS NULL;

L'escriptura a comandes:

# km0/serveis/comandes/crear_comanda.py
import json
import uuid

import psycopg

DSN = "postgresql://km0:km0_dev@localhost:5432/km0_comandes"


def crear_comanda(conn, comanda):
    """Desa la comanda I el seu esdeveniment en una única transacció. Sense Kafka aquí."""
    id_esdeveniment = uuid.uuid4()
    with conn.transaction():
        with conn.cursor() as cur:
            cur.execute("INSERT INTO comandes (id, client, mercat, estat) VALUES (%s, %s, %s, 'creada')",
                        (comanda["id"], comanda["client"], comanda["mercat"]))
            for linia in comanda["linies"]:
                cur.execute("INSERT INTO linies_comanda (comanda_id, producte, quantitat, preu_centims) "
                            "VALUES (%s, %s, %s, %s)",
                            (comanda["id"], linia["producte"], linia["quantitat"], linia["preu_centims"]))
            cur.execute("INSERT INTO outbox (id, agregat, agregat_id, tipus, carrega) "
                        "VALUES (%s, 'comanda', %s, 'comanda.creada', %s)",
                        (id_esdeveniment, comanda["id"], json.dumps(comanda)))
    return id_esdeveniment


if __name__ == "__main__":
    with psycopg.connect(DSN) as conn:
        crear_comanda(conn, {"id": "P-2026-000123", "client": "anna", "mercat": "girona",
                             "linies": [{"producte": "formatge-curat", "quantitat": 2, "preu_centims": 1450}]})

I el relay, un procés independent que pot tenir diverses instàncies gràcies a FOR UPDATE SKIP LOCKED (cadascuna bloqueja files diferents):

# km0/serveis/comandes/outbox_relay.py
import json
import time

import psycopg
from confluent_kafka import Producer

DSN = "postgresql://km0:km0_dev@localhost:5432/km0_comandes"
TOPIC = "comandes.esdeveniments"
LOT = 100

productor = Producer({"bootstrap.servers": "localhost:9092", "acks": "all"})


def publicar_pendents(conn):
    """Retorna quants esdeveniments ha publicat en aquesta passada."""
    with conn.transaction():
        with conn.cursor() as cur:
            cur.execute(
                "SELECT id, agregat_id, tipus, versio, carrega, creat_en FROM outbox "
                "WHERE publicat_en IS NULL ORDER BY creat_en "
                "LIMIT %s FOR UPDATE SKIP LOCKED", (LOT,))
            files = cur.fetchall()
            for id_esdeveniment, agregat_id, tipus, versio, carrega, creat_en in files:
                embolcall = {"id_esdeveniment": str(id_esdeveniment), "tipus": tipus, "versio": versio,
                             "data_ms": int(creat_en.timestamp() * 1000),
                             "origen": "comandes", "dades": carrega}
                productor.produce(TOPIC, key=agregat_id, value=json.dumps(embolcall).encode(),
                                  headers=[("tipus", tipus), ("id_esdeveniment", str(id_esdeveniment))])
            productor.flush()                       # espera la confirmació de Kafka
            if files:
                cur.execute("UPDATE outbox SET publicat_en = now() WHERE id = ANY(%s)",
                            ([f[0] for f in files],))
    return len(files)


if __name__ == "__main__":
    with psycopg.connect(DSN) as conn:
        while True:
            n = publicar_pendents(conn)
            if n == 0:
                time.sleep(0.2)                     # sense feina: espera curta
            else:
                print(f"[relay] publicats {n} esdeveniments")

Analitzem el relay amb la mateixa lupa que el consumidor. Si mor després de flush() i abans de l'UPDATE, la transacció es desfà, les files continuen pendents i a la passada següent es tornen a publicar: el relay és at-least-once, i produeix duplicats amb el mateix id_esdeveniment (l'id de la fila d'outbox). Aquests duplicats els absorbeix el consumidor idempotent de l'apartat 2. Els dos patrons es necessiten mútuament: outbox garanteix que cap esdeveniment no es perd; idempotència garanteix que cap no s'aplica dues vegades. I ORDER BY creat_en amb la clau de partició agregat_id manté l'ordre dels esdeveniments de cada comanda.

CDC com a alternativa al relay. En lloc d'un procés que fa polling a la taula, eines de captura de canvis de dades (Change Data Capture, CDC) com Debezium llegeixen el write-ahead log de PostgreSQL i publiquen cada inserció a outbox a Kafka amb latència de mil·lisegons i sense consultes repetides. És la implementació d'outbox recomanada a escala; el relay per polling és més fàcil d'entendre i suficient per començar. Només ho anomenem: la mecànica del log de replicació pertany a 03-04.

  1. Petició/resposta sobre missatgeria

De vegades es vol la desconnexió temporal de la missatgeria però es necessita una resposta: un servei intern que demana a pagaments que autoritzi una targeta i en vol el resultat, encara que trigui segons. El patró fa servir dues propietats del missatge AMQP:

  • reply_to: el nom de la cua on el sol·licitant espera la resposta (sovint una cua exclusiva i temporal creada per ell).
  • correlation_id: un identificador que el sol·licitant posa a la petició i que el servidor copia a la resposta, perquè el sol·licitant aparelli respostes amb peticions quan en té diverses en vol.
# km0/serveis/comandes/client_pagaments_rpc.py  (fragment: només l'enviament i l'espera)
import json
import uuid

import pika

connexio = pika.BlockingConnection(pika.ConnectionParameters("localhost"))
canal = connexio.channel()
# Cua de respostes exclusiva d'aquest procés; el broker l'esborra en desconnectar
cua_respostes = canal.queue_declare(queue="", exclusive=True).method.queue
respostes = {}


def en_respondre(ch, metode, props, cos):
    respostes[props.correlation_id] = json.loads(cos)


canal.basic_consume(queue=cua_respostes, on_message_callback=en_respondre, auto_ack=True)


def autoritzar_targeta(comanda_id, import_centims, timeout_s=5.0):
    corr_id = str(uuid.uuid4())
    canal.basic_publish(
        exchange="", routing_key="pagaments.autoritzacions",   # cua de treball de pagaments
        properties=pika.BasicProperties(reply_to=cua_respostes, correlation_id=corr_id,
                                        content_type="application/json"),
        body=json.dumps({"comanda_id": comanda_id, "import_centims": import_centims}))
    # Espera activa amb límit: processa esdeveniments del canal fins que arribi la resposta
    import time
    fi = time.monotonic() + timeout_s
    while corr_id not in respostes:
        if time.monotonic() > fi:
            raise TimeoutError(f"pagaments no ha respost a {corr_id} en {timeout_s} s")
        connexio.process_data_events(time_limit=0.1)
    return respostes.pop(corr_id)

Al costat de pagaments, el consumidor processa la petició i publica la resposta a props.reply_to amb correlation_id=props.correlation_id. (Amb BlockingConnection, process_data_events(time_limit=...) és la manera d'atendre les respostes que arriben mentre s'espera; en un client asíncron el bucle seria diferent, però la idea és la mateixa.) Aquest patró conserva els avantatges de la missatgeria (si pagaments està saturat, les peticions esperen a la cua en lloc de saturar-lo més; si es reinicia, no es perden) a canvi de més complexitat que una crida gRPC. A Quilòmetre Zero es farà servir poc: per a consultes ràpides, gRPC (02-03) és més simple; per a conseqüències, esdeveniments purs. El seu lloc són les operacions lentes o amb pics on es vol resposta però no bloquejar el servidor.

  1. Consumidors competidors i ordre per partició

Quan inventari no dona l'abast durant la "Setmana del Formatge Artesà", la solució és engegar més instàncies que comparteixin la cua (RabbitMQ) o el grup (Kafka): consumidors competidors. Cada missatge el processa una sola instància, i el rendiment escala amb el nombre d'instàncies... fins a un límit, que és diferent a cada broker:

  • A RabbitMQ, qualsevol nombre d'instàncies competeix per la mateixa cua, però en repartir missatges entre elles es perd l'ordre: comanda.creada de l'Anna pot anar a la instància 1 i comanda.cancellada a la 2, que la pot processar abans. Si l'ordre importa, cal serialitzar per una altra via (una cua per clau, o el plugin de consistent hash exchange).
  • A Kafka, el límit és el nombre de particions (02-04) i l'ordre està garantit dins de cada partició: tots els esdeveniments de P-2026-000123 els processa la mateixa instància en ordre, mentre els d'altres comandes es reparteixen. És la raó principal per la qual els esdeveniments de domini de Quilòmetre Zero van a Kafka.
flowchart LR
    K[(comandes.esdeveniments<br/>6 particions)]
    K -- "p0, p1" --> I1[inventari 1]
    K -- "p2, p3" --> I2[inventari 2]
    K -- "p4, p5" --> I3[inventari 3]
    note["clau P-2026-000123 → sempre p2 → sempre inventari 2 → en ordre"]

El preu de l'ordre per partició és el bloqueig de cap de línia de tota la partició (el mateix fenomen que a TCP, 02-01): si el missatge de l'offset 41 de la partició 2 triga deu segons, els offsets 42 en endavant d'aquesta partició esperen, encara que siguin d'altres comandes. D'aquí la importància que el consumidor no es quedi encallat en un missatge, que és el tema de l'apartat següent.

  1. Missatges enverinats: reintents amb backoff i cues de missatges morts

Un missatge enverinat (poison message) és un que fa fallar el consumidor cada vegada que l'intenta: JSON malformat, un producte que no existeix a estoc, un bug al consumidor amb certes dades. Amb at-least-once i res més, el broker el reenvia indefinidament i el consumidor es queda en bucle, bloquejant la seva partició o la seva cua. Cal distingir dues classes de fallada:

  • Transitòria: la base de dades no respon, un servei remot dona timeout. Reintentar després d'una espera té sentit, i l'espera ha de créixer (backoff exponencial: 1 s, 2 s, 4 s, 8 s...) per no martellejar un sistema malalt.
  • Permanent: dades invàlides, violació d'una regla de negoci (EstocInsuficient de l'apartat 2), bug. Reintentar no canvia res; el missatge s'ha d'apartar perquè la resta avanci, i algú l'ha d'examinar.

El destí dels missatges apartats és la cua de missatges morts (dead-letter queue, DLQ). A RabbitMQ és nativa: es declara la cua amb un dead-letter exchange i, en rebutjar sense reencuar, el broker hi mou el missatge.

# RabbitMQ: cua amb DLQ nativa
canal.exchange_declare(exchange="km0.dlx", exchange_type="direct", durable=True)
canal.queue_declare(queue="inventari.comandes.dlq", durable=True)
canal.queue_bind(queue="inventari.comandes.dlq", exchange="km0.dlx", routing_key="inventari.comandes")
canal.queue_declare(queue="inventari.comandes", durable=True, arguments={
    "x-dead-letter-exchange": "km0.dlx",
    "x-dead-letter-routing-key": "inventari.comandes",
})

def en_rebre(ch, metode, props, cos):
    try:
        processar(cos)
        ch.basic_ack(delivery_tag=metode.delivery_tag)
    except ErrorPermanent:
        ch.basic_nack(delivery_tag=metode.delivery_tag, requeue=False)   # → DLQ
    except ErrorTransitori:
        ch.basic_nack(delivery_tag=metode.delivery_tag, requeue=True)    # torna a la cua

(El requeue=True sense límit pot degenerar en el bucle; a RabbitMQ la pràctica és comptar intents en una capçalera, x-death, i fer servir una cua d'espera amb TTL per al backoff. La idea és la mateixa que implementarem a Kafka.)

Kafka no té DLQ nativa: es construeix amb tòpics de reintent i un tòpic de morts, i el consumidor decideix a quin enviar cada fallada:

sequenceDiagram
    participant T as comandes.esdeveniments
    participant C as inventari
    participant R as comandes.esdeveniments.retry
    participant D as comandes.esdeveniments.dlq
    participant O as Operador
    T->>C: esdeveniment (intent 1)
    C--xC: ErrorTransitori (BD no respon)
    C->>R: esdeveniment + capçaleres {intents: 1, no_abans: t+1s}
    C->>T: commit (la partició avança: sense bloqueig)
    R->>C: esdeveniment (espera fins a no_abans; intent 2)
    C--xC: ErrorTransitori
    C->>R: esdeveniment {intents: 2, no_abans: t+2s}
    R->>C: esdeveniment (intent 3)
    C--xC: ErrorTransitori (màxim assolit)
    C->>D: esdeveniment {intents: 3, motiu: "..."}
    D->>O: alerta: 1 missatge a la DLQ
    O->>T: després de corregir la causa, republica des de la DLQ
# km0/serveis/inventari/consum_amb_reintents.py (fragment)
import json
import time

MAX_INTENTS = 3
ESPERA_BASE_S = 1.0
TOPIC, RETRY, DLQ = "comandes.esdeveniments", "comandes.esdeveniments.retry", "comandes.esdeveniments.dlq"


class ErrorTransitori(Exception): ...
class ErrorPermanent(Exception): ...


def capcalera(msg, nom, per_defecte):
    for k, v in (msg.headers() or []):
        if k == nom:
            return v.decode()
    return per_defecte


def consumir_amb_reintents(msg, gestor, productor):
    intents = int(capcalera(msg, "intents", "0"))
    try:
        gestor(json.loads(msg.value()))
    except ErrorTransitori as e:
        if intents + 1 < MAX_INTENTS:
            espera = ESPERA_BASE_S * (2 ** intents)                # 1, 2, 4 segons
            productor.produce(RETRY, key=msg.key(), value=msg.value(), headers=[
                ("intents", str(intents + 1)), ("no_abans", str(time.time() + espera)),
                ("ultim_error", str(e)[:200])])
        else:
            productor.produce(DLQ, key=msg.key(), value=msg.value(), headers=[
                ("intents", str(intents + 1)), ("motiu", f"transitori esgotat: {e}"[:200])])
    except ErrorPermanent as e:
        productor.produce(DLQ, key=msg.key(), value=msg.value(), headers=[
            ("intents", str(intents + 1)), ("motiu", f"permanent: {e}"[:200])])
    productor.flush()
    # En tots els casos es confirma l'offset del tòpic d'origen: el missatge ja
    # és fora de perill a retry o a dlq, i la partició no es bloqueja.

Un consumidor del tòpic .retry llegeix cada missatge, dorm fins a no_abans si cal, i el processa amb la mateixa funció. Consideracions que no es veuen al codi: el missatge que va a .retry perd la seva posició en l'ordre de la partició original (es processarà després d'altres de més recents), cosa que per a un descompte d'estoc és acceptable i per a una seqüència creada → cancellada pot no ser-ho; la DLQ necessita alertes (07-01) i un procediment per examinar, corregir i republicar; i EstocInsuficient de l'apartat 2 és un bon exemple d'error permanent la resolució del qual no és tècnica sinó de negoci (avisar el client, cancel·lar la comanda), que és el territori de les sagues de 03-05. Els reintents, backoff i circuit breakers de les crides síncrones es tracten a 07-04; aquí només hem vist els que pertanyen al consum de missatges.

  1. Event-driven i event-sourcing: dues coses diferents

Tot el que hem fet a 02-04 i en aquesta lliçó és arquitectura dirigida per esdeveniments (event-driven): els serveis comuniquen fets ocorreguts (comanda.creada) i d'altres hi reaccionen. L'estat de cada servei continua vivint a les seves taules (comandes, estoc); els esdeveniments són notificacions.

Event-sourcing és una altra cosa, que sovint es confon amb l'anterior: consisteix en el fet que l'estat no es desa; es desen els esdeveniments, i l'estat es reconstrueix reproduint-los. La comanda de l'Anna no seria una fila amb estat = 'pagada', sinó la seqüència ComandaCreada, LiniaAfegida, ComandaPagada, i la fila seria una projecció derivada. Dona auditoria completa i permet reconstruir qualsevol estat passat, a canvi d'una complexitat considerable (projeccions, instantànies, evolució d'esdeveniments històrics). Es pot fer event-driven sense event-sourcing (Quilòmetre Zero ho fa) i viceversa. Ho anomenem perquè el terme no es confongui; el seu desenvolupament excedeix aquest curs.

  1. Esquemes d'esdeveniments i versionat

Un esdeveniment és un contracte entre el productor i tots els consumidors presents i futurs, inclosos els que llegiran l'històric de Kafka d'aquí a sis dies. S'hi apliquen les regles d'evolució d'esquemes de 02-03 amb més severitat, perquè no es pot saber quants consumidors hi ha ni obligar-los a actualitzar-se. L'envoltura estàndard que hem fet servir ja preveu el mecanisme:

{
  "id_esdeveniment": "6f1c9e2a-...",
  "tipus": "comanda.creada",
  "versio": 2,
  "data_ms": 1789000000000,
  "origen": "comandes",
  "dades": {
    "id": "P-2026-000123",
    "client": "anna",
    "linies": [{"producte": "formatge-curat", "quantitat": 2, "preu_centims": 1450}],
    "adreca_lliurament": {"ciutat": "Girona", "codi_postal": "17001"}
  }
}

Les regles, adaptades de 02-03:

  • Canvis compatibles (no incrementen versio): afegir camps opcionals a dades; afegir tipus d'esdeveniment nous. Els consumidors han d'ignorar camps desconeguts i tolerar l'absència dels nous.
  • Canvis incompatibles (incrementen versio): eliminar o reanomenar un camp, canviar-ne el tipus o el significat. Durant la transició, el productor pot publicar totes dues versions o els consumidors acceptar-les totes dues (if esdeveniment["versio"] == 1: ...). Mai no es canvia el significat d'un camp conservant-ne el nom.
  • Registre d'esquemes (Schema Registry): amb Avro o protobuf a Kafka, un servei central desa cada versió de l'esquema de cada tòpic, valida al productor que un canvi és compatible (cap enrere, cap endavant o totes dues, segons la política) i permet als consumidors deserialitzar missatges antics amb l'esquema amb què es van escriure. És la manera de convertir les regles anteriors en una comprovació automàtica, i el pas natural quan JSON es queda curt. Quilòmetre Zero comença amb JSON i l'envoltura; el registre d'esquemes s'introdueix amb el pipeline de dades del Mòdul 5.

  1. Taula de patrons: quin problema resol cadascun

Patró Problema que resol Cost On es fa servir a Quilòmetre Zero
At-least-once + ack després de processar Pèrdua de missatges en morir el consumidor Duplicats Tots els consumidors d'esdeveniments de domini
Consumidor idempotent (taula missatges_processats) Duplicats per reintent del productor, del relay o del broker Una fila per missatge; purga periòdica inventari (descompte d'estoc i ReservarEstoc gRPC), pagaments (cobraments)
Outbox transaccional + relay Escriptura dual: estat desat sense esdeveniment, o esdeveniment sense estat Un procés més; latència de mil·lisegons a segons comandes (tots els seus esdeveniments); després, cada servei que publiqui
CDC (Debezium) Relay per polling a escala Infraestructura addicional Quan el polling no n'hi hagi prou
Petició/resposta (reply_to, correlation_id) Necessitar resposta amb l'amortiment d'una cua Complexitat davant de gRPC Operacions lentes amb pics (autoritzacions de pagaments)
Consumidors competidors Un consumidor no dona l'abast Pèrdua d'ordre (RabbitMQ) o límit per particions (Kafka) Tots els serveis en campanya
Clau de partició Ordre entre esdeveniments d'una mateixa entitat Bloqueig de cap de línia per partició; particions desequilibrades comandes.esdeveniments per id de comanda; telemetria per repartidor
Reintents amb backoff en el consum Fallades transitòries del consumidor Pèrdua d'ordre del missatge reintentat Tots els consumidors
Dead-letter queue Missatges enverinats que bloquegen la cua o la partició Necessita alertes i procediment de reprocessament Tots els consumidors
Envoltura + versionat d'esdeveniments Evolució del contracte sense coordinar desplegaments Disciplina i revisió de canvis Tots els esdeveniments

Errors Comuns i Consells

  • Buscar exactly-once a la configuració del broker. No hi és. És al consumidor idempotent, i no hi ha drecera.
  • Registrar la clau d'idempotència fora de la transacció de l'efecte. Reobre la finestra que es volia tancar. Mateixa transacció, sempre.
  • Fer servir com a clau d'idempotència alguna cosa que canvia entre reintents. Si el productor genera un id_esdeveniment nou a cada reintent, el consumidor no pot reconèixer el duplicat. La clau es genera una vegada i es reutilitza.
  • Publicar a Kafka "dins" de la transacció de PostgreSQL. El ROLLBACK no la desfà. Outbox o res.
  • Un relay que marca com a publicat abans del flush. Si Kafka no confirma, l'esdeveniment es perd per sempre amb la fila marcada. Confirmació primer, marca després.
  • Reintentar errors permanents. Un JSON malformat no s'arregla esperant 4 segons; bloqueja la partició tres vegades més. Classifica els errors i envia els permanents a la DLQ a la primera.
  • Una DLQ sense alertes ni propietari. És un forat negre on les comandes desapareixen en silenci. Cada DLQ necessita una alerta, una persona i un procediment.
  • Canviar el significat d'un camp sense canviar la versió. Un consumidor que rellegirà l'històric llegirà dades amb dos significats sota el mateix nom. Camp nou o versió nova.
  • Consell: prova la idempotència de manera explícita: a l'entorn de proves, publica cada esdeveniment dues vegades i comprova que l'estat final és el mateix. És la prova més barata i la que més incidents evita.
  • Consell: mesura el lag del relay (files d'outbox amb publicat_en IS NULL i la seva antiguitat) i de cada grup de consumidors. Són els dos indicadors que diuen si l'asincronia està funcionant o acumulant deute invisible.

Exercicis

Exercici 1: Auditar un consumidor

Aquest consumidor de pagaments processa comanda.creada i executa el cobrament contra la passarel·la externa. Identifica tots els problemes de garanties de lliurament que té, indica què pot sortir malament en cadascun (pèrdua, duplicat, bloqueig) i reescriu-lo aplicant els patrons de la lliçó. Assumeix que la passarel·la accepta una capçalera Idempotency-Key i que pagaments té la seva pròpia base de dades PostgreSQL.

consumidor = Consumer({"bootstrap.servers": "localhost:9092", "group.id": "pagaments",
                       "enable.auto.commit": True})
consumidor.subscribe(["comandes.esdeveniments"])
while True:
    msg = consumidor.poll(1.0)
    if msg is None or msg.error():
        continue
    esdeveniment = json.loads(msg.value())
    if esdeveniment["tipus"] == "comanda.creada":
        import_centims = sum(l["quantitat"] * l["preu_centims"] for l in esdeveniment["dades"]["linies"])
        resposta = passarella.cobrar(targeta=esdeveniment["dades"]["targeta"], import_centims=import_centims)
        with psycopg.connect(DSN) as conn:
            conn.execute("INSERT INTO pagaments (comanda_id, import_centims, ref_passarella) VALUES (%s, %s, %s)",
                         (esdeveniment["dades"]["id"], import_centims, resposta["ref"]))

Exercici 2: Outbox per a inventari

inventari ha de publicar estoc.actualitzat (amb producte, abans, despres, motiu) cada vegada que canvia l'estoc, tant per la reserva gRPC de 02-03 com pel consum de comanda.creada de l'apartat 2. Dissenya la taula outbox d'inventari, modifica processar_comanda_creada per escriure l'esdeveniment a la mateixa transacció, i explica què garanteix el conjunt (idempotència + outbox) davant de cadascuna d'aquestes fallades: (a) el consumidor mor entre el COMMIT i el commit d'offset; (b) el relay mor després de publicar i abans de marcar; (c) Kafka no està disponible durant 10 minuts.

Exercici 3: Dissenyar la política de reintents i DLQ

Per al consumidor de repartiment que assigna un repartidor a cada comanda.pagada, classifica cadascuna d'aquestes fallades com a transitòria o permanent, indica on va el missatge (reintent amb backoff, DLQ) i què fa l'operador en cada cas: (1) el servei de mapes extern retorna 503; (2) la comanda té una adreça de lliurament en una ciutat on Quilòmetre Zero no reparteix; (3) KeyError: 'adreca_lliurament' en un esdeveniment de versio: 1; (4) la base de dades de repartiment rebutja la connexió durant un reinici; (5) no hi ha cap repartidor lliure a Tarragona ara mateix.

Solucions

Solució 1:

Problemes: (1) enable.auto.commit=True confirma offsets periòdicament amb independència de si s'ha processat: si el procés mor després de l'auto-commit i abans de cobrar, el cobrament es perd (at-most-once); si mor després de cobrar i abans de l'auto-commit, es cobra dues vegades (at-least-once sense idempotència). (2) El cobrament a la passarel·la i l'INSERT són dues escriptures sense transacció comuna: si la passarel·la cobra i l'INSERT falla, hi ha un cobrament sense registre; en reintentar-se, un altre cobrament. (3) No hi ha clau d'idempotència ni cap a la passarel·la ni a la base de dades: el mateix comanda.creada lliurat dues vegades cobra dues vegades a l'Anna. (4) msg.error() s'ignora en silenci. (5) Qualsevol excepció (JSON malformat, targeta invàlida) mata el bucle o, si es capturés, bloquejaria la partició. (6) No publica pagament.confirmat, així que ningú no s'assabenta del cobrament (escriptura dual pendent). Reescriptura:

consumidor = Consumer({"bootstrap.servers": "localhost:9092", "group.id": "pagaments",
                       "auto.offset.reset": "earliest", "enable.auto.commit": False})
consumidor.subscribe(["comandes.esdeveniments"])

def cobrar_comanda(conn, esdeveniment):
    comanda = esdeveniment["dades"]
    with conn.transaction():
        with conn.cursor() as cur:
            cur.execute("INSERT INTO missatges_processats (id_missatge, consumidor) VALUES (%s, 'pagaments.cobrament') "
                        "ON CONFLICT DO NOTHING RETURNING id_missatge", (esdeveniment["id_esdeveniment"],))
            if cur.fetchone() is None:
                return "duplicat"
            import_centims = sum(l["quantitat"] * l["preu_centims"] for l in comanda["linies"])
            # La passarel·la deduplica per Idempotency-Key: fer servir l'id de la comanda fa que
            # un reintent (fins i tot amb un altre id_esdeveniment) no cobri dues vegades.
            resposta = passarella.cobrar(targeta=comanda["targeta"], import_centims=import_centims,
                                         idempotency_key=f"cobrament-{comanda['id']}")
            cur.execute("INSERT INTO pagaments (comanda_id, import_centims, ref_passarella) VALUES (%s, %s, %s)",
                        (comanda["id"], import_centims, resposta["ref"]))
            cur.execute("INSERT INTO outbox (id, agregat, agregat_id, tipus, carrega) "
                        "VALUES (%s, 'pagament', %s, 'pagament.confirmat', %s)",
                        (uuid.uuid4(), comanda["id"], json.dumps({"comanda_id": comanda["id"], "import_centims": import_centims})))
    return "processat"

with psycopg.connect(DSN) as conn:
    while True:
        msg = consumidor.poll(1.0)
        if msg is None:
            continue
        if msg.error():
            print("error de Kafka:", msg.error()); continue
        try:
            esdeveniment = json.loads(msg.value())
        except json.JSONDecodeError as e:
            a_dlq(msg, f"permanent: {e}"); consumidor.commit(message=msg); continue
        if esdeveniment["tipus"] == "comanda.creada":
            try:
                cobrar_comanda(conn, esdeveniment)
            except PassarellaNoDisponible as e:
                a_retry(msg, e)                  # transitori: backoff
            except TargetaRebutjada as e:
                a_dlq(msg, f"permanent: {e}")    # de negoci: comanda a cancel·lar (03-05)
        consumidor.commit(message=msg)

Queda una finestra inevitable: si el procés mor entre passarella.cobrar i el COMMIT, la transacció es desfà (ni clau ni registre), el missatge es reintenta, i el segon cobrar amb la mateixa Idempotency-Key fa que la passarel·la retorni el cobrament ja fet en lloc de repetir-lo. La idempotència del sistema extern tanca el que la transacció local no pot cobrir. Sense aquesta capçalera, caldria registrar "cobrament iniciat" abans de cridar i reconciliar després: és la situació en què les sagues de 03-05 es fan necessàries.

Solució 2:

CREATE TABLE outbox (
    id UUID PRIMARY KEY, agregat TEXT NOT NULL, agregat_id TEXT NOT NULL,
    tipus TEXT NOT NULL, versio INTEGER NOT NULL DEFAULT 1, carrega JSONB NOT NULL,
    creat_en TIMESTAMPTZ NOT NULL DEFAULT now(), publicat_en TIMESTAMPTZ);

Amb agregat = 'producte' i agregat_id = producte, perquè els canvis d'un mateix producte vagin en ordre a la mateixa partició d'inventari.esdeveniments. A processar_comanda_creada, dins del mateix with conn.transaction(), després de cada UPDATE amb èxit:

cur.execute("SELECT unitats FROM estoc WHERE producte = %s", (linia["producte"],))
despres = cur.fetchone()[0]
cur.execute("INSERT INTO outbox (id, agregat, agregat_id, tipus, carrega) VALUES (%s, 'producte', %s, 'estoc.actualitzat', %s)",
            (uuid.uuid4(), linia["producte"],
             json.dumps({"producte": linia["producte"], "abans": despres + linia["quantitat"],
                         "despres": despres, "motiu": "comanda", "comanda_id": esdeveniment["dades"]["id"]})))

(I el mateix a ReservarEstoc del servidor gRPC, que deixaria de fer servir el diccionari en memòria per fer servir la mateixa base de dades.) Garanties: (a) el consumidor mor entre COMMIT i commit d'offset: l'estoc està descomptat i l'esdeveniment és a outbox; el relay el publicarà; el missatge reenviat es detecta com a duplicat i no genera ni descompte ni segon esdeveniment. (b) El relay mor després de publicar i abans de marcar: l'esdeveniment es publica dues vegades amb el mateix id_esdeveniment; cataleg i analitica, si són idempotents, l'ignoren la segona vegada; si cataleg només fa SET indicador = (despres < 5), és idempotent per naturalesa i ni tan sols necessita la taula. (c) Kafka caigut 10 minuts: inventari continua processant (si llegeix de RabbitMQ) o s'atura (si llegeix de Kafka), però en tots dos casos no perd res: els esdeveniments s'acumulen a outbox amb publicat_en IS NULL i el relay els publica en ordre quan Kafka torna. És exactament el comportament que es demanava a la solució 3.3 de la lliçó 01-06, ara implementat.

Solució 3:

  1. Transitòria: 503 és "torna més tard". Reintent amb backoff (1, 2, 4 s...) fins al màxim; després, DLQ. L'operador comprova l'estat del proveïdor de mapes; en recuperar-se, republica la DLQ.
  2. Permanent, de negoci: no hi ha reintent que ho arregli. DLQ amb motiu ciutat_no_coberta, i probablement un esdeveniment repartiment.rebutjat perquè comandes cancel·li i pagaments reemborsi (una saga, 03-05). L'operador o, millor, la validació a comandes en crear la comanda, ho hauria d'haver impedit abans: la DLQ revela un forat a la validació.
  3. Permanent, de contracte: el consumidor assumeix l'esquema versio: 2 i en rep un de versio: 1 (probablement de l'històric o d'un comandes encara no actualitzat). DLQ amb motiu esquema, alerta a l'equip, i la correcció és al consumidor (acceptar totes dues versions, apartat 8); després de desplegar-la, es republiquen els missatges de la DLQ. No és una fallada de dades sinó d'un consumidor que no ha seguit les regles de compatibilitat.
  4. Transitòria: reinici de la base de dades. Reintent amb backoff. Aquí és especialment important no confirmar l'offset sense haver apartat el missatge a .retry; amb la base de dades caiguda, ni tan sols es pot registrar la idempotència, així que el missatge ha de quedar íntegre per reintentar-se.
  5. Transitòria, però de negoci i de llarga durada: pot trigar 20 minuts a haver-hi un repartidor lliure. Un backoff de segons no hi encaixa; és millor que el consumidor processi el missatge registrant la comanda com a "pendent d'assignació" a la seva base de dades (idempotentment) i confirmi, i que un procés periòdic intenti assignar les pendents. Convertir una espera llarga en estat persistit és preferible a mantenir missatges rebotant entre tòpics de reintent, que a més perdrien l'ordre respecte d'una possible comanda.cancellada posterior.

Conclusió

Aquesta lliçó ha tancat el mòdul resolent els problemes que les anteriors van anar deixant oberts. Les garanties de lliurament depenen d'on es confirma respecte d'on es processa: confirmar abans perd missatges (at-most-once), confirmar després els duplica (at-least-once), i exactly-once no existeix d'extrem a extrem, així que l'estratègia és at-least-once al transport i idempotència al consumidor. El consumidor idempotent d'inventari, amb la seva taula missatges_processats a la mateixa transacció que l'efecte, ha tancat el problema dels duplicats que arrossegàvem des de la simulació de 01-04 i les semàntiques de 02-02, i la mateixa tècnica fa segurs els reintents de la crida gRPC de 02-03. El patró Outbox, amb el seu relay FOR UPDATE SKIP LOCKED, ha eliminat l'escriptura dual a comandes: estat i esdeveniment es desen junts o no es desen. Hem vist a més petició/resposta sobre cues amb reply_to i correlation_id, consumidors competidors amb ordre per partició, reintents amb backoff, cues de missatges morts amb les seves alertes i el seu procediment, la diferència entre event-driven i event-sourcing, i el versionat d'esdeveniments amb l'envoltura estàndard i les regles heretades de 02-03. La taula de l'apartat 9 resumeix quin patró resol quin problema i on el fa servir Quilòmetre Zero.

Amb això, el Mòdul 2 ha complert la seva promesa: comandes i inventari són dos processos que parlen de manera fiable, síncronament per gRPC quan cal resposta i asíncronament per Kafka per a les conseqüències, i analitica escolta sense que ningú no l'hagi d'esperar. Però fixa't en el que ha passat pel camí: inventari té ara la seva pròpia taula estoc i comandes la seva pròpia taula comandes, i la veritat sobre l'últim formatge curat de la Formatgeria Montblanc ja no és en un únic lloc. Quan l'Anna el reserva i cataleg rep l'esdeveniment estoc.actualitzat mig segon després, durant aquest mig segon en Marc veu al catàleg un formatge que ja no existeix. Quan inventari tingui les rèpliques inv-bcn i inv-vlc, quina de les dues té raó si difereixen? Què vol dir exactament "consistent" quan les dades viuen en diversos llocs, i què es pot prometre a un client davant d'una fallada de xarxa entre ells? Aquestes són les preguntes del Mòdul 3, Consistència i Replicació, que comença amb els models de consistència: el vocabulari precís per dir què garanteix un sistema distribuït sobre les seves dades i què no.

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