La lección anterior terminó con tres preguntas incómodas. Si inventario descuenta el stock y muere antes de confirmar el mensaje, el broker lo reenvía y el stock se descuenta dos veces. Si pedidos guarda el pedido en PostgreSQL y se cae antes de publicar pedido.creado, inventario y analitica nunca sabrán que ese pedido existe. Y si un mensaje malformado hace fallar al consumidor una y otra vez, todos los mensajes que vienen detrás se quedan esperando. Ninguna de las tres es un caso raro: son consecuencias directas de las falacias de 01-04 aplicadas a la mensajería, y aparecerán en producción en la primera campaña.

Esta lección presenta los patrones con los que la industria ha aprendido a convivir con ellas. Primero pondremos nombre a lo que un broker puede y no puede garantizar (at-most-once, at-least-once y el mito de exactly-once), y veremos cómo esas garantías dependen de dónde se coloca el ack. Después cerraremos, por fin, el problema de los duplicados que dejamos abierto en 01-04 y 02-02 con consumidores idempotentes respaldados por una tabla en PostgreSQL. Resolveremos la escritura dual con el patrón Outbox transaccional y un relay, implementaremos petición/respuesta sobre mensajería, escalaremos con consumidores competidores respetando el orden por partición, y trataremos los mensajes envenenados con reintentos y colas de mensajes muertos. Terminaremos con el versionado de esquemas de eventos y una tabla que mapea cada patrón al problema que resuelve. Las sagas (transacciones entre servicios) y el circuit breaker (resiliencia de llamadas síncronas) quedan para 03-05 y 07-04.

Contenido

  1. Garantías de entrega: at-most-once, at-least-once y el mito de exactly-once
  2. Consumidores idempotentes: cerrar el problema de los duplicados
  3. El problema de la escritura dual y el patrón Outbox transaccional
  4. Petición/respuesta sobre mensajería
  5. Consumidores competidores y orden por partición
  6. Mensajes envenenados: reintentos con backoff y colas de mensajes muertos
  7. Event-driven y event-sourcing: dos cosas distintas
  8. Esquemas de eventos y versionado
  9. Tabla de patrones: qué problema resuelve cada uno
  10. Errores comunes y consejos
  11. Ejercicios
  12. Conclusión

  1. Garantías de entrega: at-most-once, at-least-once y el mito de exactly-once

En 02-02 vimos las semánticas de invocación de RPC. Las garantías de entrega de la mensajería son la misma idea con un broker en medio, y el factor que las decide es cuándo se confirma (ack en RabbitMQ, commit de offset en Kafka) respecto a cuándo se procesa:

sequenceDiagram
    participant B as Broker
    participant C as Consumidor (inventario)
    participant DB as PostgreSQL
    Note over B,DB: Opción A: confirmar ANTES de procesar → at-most-once
    B->>C: pedido.creado (offset 41)
    C->>B: commit(41)
    C-xDB: UPDATE stock ... (el proceso muere aquí)
    Note over B,DB: El broker cree que 41 está hecho: el mensaje se PIERDE
    Note over B,DB: Opción B: confirmar DESPUÉS de procesar → at-least-once
    B->>C: pedido.creado (offset 42)
    C->>DB: UPDATE stock ... COMMIT
    C-xB: commit(42) (el proceso muere aquí)
    Note over B,DB: El broker reenvía 42: el stock se descuenta DOS veces
Garantía Cómo se consigue Riesgo Cuándo es aceptable
At-most-once Confirmar antes de procesar (o no confirmar nunca, con auto_ack) Pérdida de mensajes Telemetría, métricas, cualquier dato que el siguiente mensaje reemplaza
At-least-once Confirmar después de procesar; el productor reintenta hasta recibir confirmación del broker Duplicados Todo lo que importe, siempre que el consumidor sea idempotente
Exactly-once No existe de extremo a extremo Se emula con at-least-once + idempotencia

El mito, con precisión. Kafka ofrece una funcionalidad llamada exactly-once semantics (productores idempotentes y transacciones de Kafka) que garantiza que un mensaje no se duplica dentro de Kafka: si el productor reintenta por un timeout, el broker deduplica; y una aplicación que lee de un tópico y escribe en otro puede hacerlo atómicamente. Es valioso para pipelines Kafka → Kafka (Módulo 5). Pero en cuanto el efecto del mensaje sale de Kafka (un UPDATE en PostgreSQL, un correo, un cobro), la ventana entre "efecto aplicado" y "offset confirmado" reaparece, y nadie puede cerrarla desde fuera. Lo mismo ocurre con el QoS 2 de MQTT, que garantiza una única entrega al cliente, no un único efecto. La conclusión práctica es la de 02-02: at-least-once en el transporte, idempotencia en el consumidor. Todo lo demás en esta lección se construye sobre esa base.

  1. Consumidores idempotentes: cerrar el problema de los duplicados

Una operación es idempotente si ejecutarla N veces produce el mismo estado que ejecutarla una. Algunas lo son por naturaleza (SET stock = 118, "marcar el pedido como pagado", "insertar con clave primaria"); otras no (stock = stock - 2, "enviar un correo", "cobrar 14,50 €"). Para las que no lo son, la técnica es hacer que el consumidor recuerde qué mensajes ya ha procesado y descarte los repetidos. Dos requisitos:

  1. Una clave de idempotencia por mensaje: un identificador único generado por el productor y estable entre reintentos. Es el id_evento de la envoltura de 02-04, y era el id_reserva de la petición gRPC de 02-03. Sin ella no hay forma de saber que dos mensajes son "el mismo".
  2. Registrar la clave y aplicar el efecto en la misma transacción. Si se registra la clave en una transacción y se aplica el efecto en otra, reaparece la ventana: se registra, muere el proceso, el efecto nunca ocurre, y el reintento se descarta como duplicado (pérdida). Si es al revés, el efecto se aplica, muere, y el reintento lo aplica otra vez (duplicado). Con ambos en la misma transacción, o se hacen los dos o ninguno.

inventario ya tiene su propia base de datos (uno de los principios de 01-06), así que añadimos la tabla allí:

-- km0/sql/inventario/002_mensajes_procesados.sql
CREATE TABLE stock (
    producto TEXT PRIMARY KEY,
    unidades INTEGER NOT NULL CHECK (unidades >= 0)
);
INSERT INTO stock VALUES ('tomate-rosa', 120), ('calabacin', 80), ('queso-curado', 5),
                         ('queso-fresco', 30), ('vino-crianza', 200);

CREATE TABLE mensajes_procesados (
    id_mensaje   TEXT        NOT NULL,   -- id_evento de la envoltura
    consumidor   TEXT        NOT NULL,   -- 'inventario.descuento_stock'
    procesado_en TIMESTAMPTZ NOT NULL DEFAULT now(),
    PRIMARY KEY (id_mensaje, consumidor)
);

La clave primaria compuesta permite que dos lógicas distintas dentro de inventario (descontar stock y, por ejemplo, avisar al productor) procesen el mismo evento cada una una vez. El consumidor de Kafka de 02-04, ahora idempotente, con psycopg (versión 3):

# km0/servicios/inventario/consumidor_idempotente.py
import json

import psycopg
from confluent_kafka import Consumer, KafkaError

CONSUMIDOR = "inventario.descuento_stock"
DSN = "postgresql://km0:km0_dev@localhost:5432/km0_inventario"


class StockInsuficiente(Exception):
    pass


def procesar_pedido_creado(conn, evento):
    """Aplica el descuento de stock UNA sola vez por id_evento.
    Devuelve 'procesado' o 'duplicado'. Lanza si el efecto no puede aplicarse."""
    with conn.transaction():                              # BEGIN ... COMMIT/ROLLBACK
        with conn.cursor() as cur:
            # 1. Intentar registrar la clave. Si ya existe, ON CONFLICT no inserta
            #    y RETURNING no devuelve fila: es un duplicado.
            cur.execute(
                "INSERT INTO mensajes_procesados (id_mensaje, consumidor) VALUES (%s, %s) "
                "ON CONFLICT DO NOTHING RETURNING id_mensaje",
                (evento["id_evento"], CONSUMIDOR))
            if cur.fetchone() is None:
                return "duplicado"
            # 2. Aplicar el efecto en LA MISMA transacción.
            for linea in evento["datos"]["lineas"]:
                cur.execute(
                    "UPDATE stock SET unidades = unidades - %s "
                    "WHERE producto = %s AND unidades >= %s",
                    (linea["cantidad"], linea["producto"], linea["cantidad"]))
                if cur.rowcount == 0:
                    # Deshace también el registro del paso 1: el mensaje NO queda
                    # marcado como procesado y podrá reintentarse o ir a la DLQ (ap. 6)
                    raise StockInsuficiente(f"sin stock de {linea['producto']}")
    return "procesado"


def main():
    consumidor = Consumer({"bootstrap.servers": "localhost:9092", "group.id": "inventario",
                           "auto.offset.reset": "earliest", "enable.auto.commit": False})
    consumidor.subscribe(["pedidos.eventos"])
    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("[inventario] error de Kafka:", msg.error())
                continue
            evento = json.loads(msg.value())
            if evento["tipo"] == "pedido.creado":
                resultado = procesar_pedido_creado(conn, evento)
                print(f"[inventario] {evento['id_evento'][:8]} pedido {evento['datos']['id']}: {resultado}")
            consumidor.commit(message=msg)   # SIEMPRE después de la transacción de BD


if __name__ == "__main__":
    main()

Recorramos los fallos posibles con este código:

  • Muere después del COMMIT de PostgreSQL y antes del commit de Kafka. Kafka reenvía el mensaje; el INSERT choca con la clave primaria; ON CONFLICT DO NOTHING no inserta; fetchone() devuelve None; se responde duplicado sin tocar el stock; se confirma el offset. Sin duplicación.
  • Muere en medio de la transacción. PostgreSQL hace ROLLBACK: ni la clave ni el descuento persisten. Kafka reenvía; se procesa desde cero. Sin pérdida.
  • Muere antes de empezar. Kafka reenvía. Trivial.
  • El productor publicó el mismo evento dos veces (reintentó por un timeout del broker). Mismo id_evento, mismo resultado: la segunda copia es un duplicado. Es la situación de la simulación de 01-04, ahora resuelta: el reintento ya no convierte la pérdida de un mensaje en duplicación de efectos.

Lo mismo se aplica a la llamada síncrona de 02-03: ReservarStock puede registrar id_reserva en mensajes_procesados con consumidor = 'inventario.reserva_grpc' y devolver la respuesta guardada si se repite, con lo que el "estado desconocido" tras un DEADLINE_EXCEEDED deja de ser un problema: pedidos puede reintentar con el mismo id_reserva sin riesgo. Es la semántica at-most-once de 02-02 hecha con una tabla. Dos consideraciones operativas: la tabla crece (una fila por mensaje) y debe purgarse por antigüedad (más allá de la retención del tópico ya no puede llegar un reintento: DELETE ... WHERE procesado_en < now() - interval '14 days'); y si el consumidor no tiene base de datos transaccional (por ejemplo, envía correos), la idempotencia tiene que apoyarse en el sistema destino (un proveedor de correo que acepte una clave de idempotencia) o aceptar at-most-once.

  1. El problema de la escritura dual y el patrón Outbox transaccional

Ahora el lado del productor. En 02-04, pedidos hacía dos cosas al crear un pedido: escribirlo en su PostgreSQL y publicar pedido.creado en Kafka. Son dos sistemas distintos y no hay transacción que los abarque (es la escritura dual, dual write):

  • Si escribe en la base de datos y se cae antes de publicar, el pedido existe pero nadie se entera: inventario no descuenta, analitica no cuenta, reparto no reparte.
  • Si publica primero y luego falla el COMMIT, todos reaccionan a un pedido que no existe.
  • Si publica dentro de la transacción "para que se deshaga", no se deshace: Kafka no participa en el ROLLBACK de PostgreSQL.

El patrón Outbox transaccional resuelve el problema convirtiendo dos escrituras en una: el evento se escribe en una tabla de la misma base de datos, en la misma transacción que el pedido. Un proceso aparte, el relay, lee esa tabla y publica en Kafka.

sequenceDiagram
    participant A as App de Ana
    participant P as pedidos
    participant DB as PostgreSQL (pedidos)
    participant R as Relay outbox
    participant K as Kafka
    participant I as inventario
    A->>P: crear pedido
    P->>DB: BEGIN
    P->>DB: INSERT pedidos, lineas_pedido
    P->>DB: INSERT outbox (pedido.creado)
    P->>DB: COMMIT (atómico: pedido + evento, o nada)
    P-->>A: pedido P-2026-000123 creado
    loop cada 200 ms
        R->>DB: SELECT ... FROM outbox WHERE publicado_en IS NULL FOR UPDATE SKIP LOCKED
        R->>K: produce(pedido.creado)
        K-->>R: confirmado
        R->>DB: UPDATE outbox SET publicado_en = now()
    end
    K->>I: pedido.creado (at-least-once)

La tabla, en la base de datos de pedidos:

-- km0/sql/pedidos/002_outbox.sql
CREATE TABLE outbox (
    id           UUID        PRIMARY KEY,           -- será el id_evento
    agregado     TEXT        NOT NULL,              -- 'pedido'
    agregado_id  TEXT        NOT NULL,              -- 'P-2026-000123': clave de partición
    tipo         TEXT        NOT NULL,              -- 'pedido.creado'
    version      INTEGER     NOT NULL DEFAULT 1,
    carga        JSONB       NOT NULL,
    creado_en    TIMESTAMPTZ NOT NULL DEFAULT now(),
    publicado_en TIMESTAMPTZ
);
CREATE INDEX outbox_pendientes ON outbox (creado_en) WHERE publicado_en IS NULL;

La escritura en pedidos:

# km0/servicios/pedidos/crear_pedido.py
import json
import uuid

import psycopg

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


def crear_pedido(conn, pedido):
    """Guarda el pedido Y su evento en una única transacción. Sin Kafka aquí."""
    id_evento = uuid.uuid4()
    with conn.transaction():
        with conn.cursor() as cur:
            cur.execute("INSERT INTO pedidos (id, cliente, mercado, estado) VALUES (%s, %s, %s, 'creado')",
                        (pedido["id"], pedido["cliente"], pedido["mercado"]))
            for linea in pedido["lineas"]:
                cur.execute("INSERT INTO lineas_pedido (pedido_id, producto, cantidad, precio_centimos) "
                            "VALUES (%s, %s, %s, %s)",
                            (pedido["id"], linea["producto"], linea["cantidad"], linea["precio_centimos"]))
            cur.execute("INSERT INTO outbox (id, agregado, agregado_id, tipo, carga) "
                        "VALUES (%s, 'pedido', %s, 'pedido.creado', %s)",
                        (id_evento, pedido["id"], json.dumps(pedido)))
    return id_evento


if __name__ == "__main__":
    with psycopg.connect(DSN) as conn:
        crear_pedido(conn, {"id": "P-2026-000123", "cliente": "ana", "mercado": "girona",
                            "lineas": [{"producto": "queso-curado", "cantidad": 2, "precio_centimos": 1450}]})

Y el relay, un proceso independiente que puede tener varias instancias gracias a FOR UPDATE SKIP LOCKED (cada una bloquea filas distintas):

# km0/servicios/pedidos/outbox_relay.py
import json
import time

import psycopg
from confluent_kafka import Producer

DSN = "postgresql://km0:km0_dev@localhost:5432/km0_pedidos"
TOPICO = "pedidos.eventos"
LOTE = 100

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


def publicar_pendientes(conn):
    """Devuelve cuántos eventos ha publicado en esta pasada."""
    with conn.transaction():
        with conn.cursor() as cur:
            cur.execute(
                "SELECT id, agregado_id, tipo, version, carga, creado_en FROM outbox "
                "WHERE publicado_en IS NULL ORDER BY creado_en "
                "LIMIT %s FOR UPDATE SKIP LOCKED", (LOTE,))
            filas = cur.fetchall()
            for id_evento, agregado_id, tipo, version, carga, creado_en in filas:
                envoltura = {"id_evento": str(id_evento), "tipo": tipo, "version": version,
                             "fecha_ms": int(creado_en.timestamp() * 1000),
                             "origen": "pedidos", "datos": carga}
                productor.produce(TOPICO, key=agregado_id, value=json.dumps(envoltura).encode(),
                                  headers=[("tipo", tipo), ("id_evento", str(id_evento))])
            productor.flush()                       # espera la confirmación de Kafka
            if filas:
                cur.execute("UPDATE outbox SET publicado_en = now() WHERE id = ANY(%s)",
                            ([f[0] for f in filas],))
    return len(filas)


if __name__ == "__main__":
    with psycopg.connect(DSN) as conn:
        while True:
            n = publicar_pendientes(conn)
            if n == 0:
                time.sleep(0.2)                     # sin trabajo: espera corta
            else:
                print(f"[relay] publicados {n} eventos")

Analicemos el relay con la misma lupa que el consumidor. Si muere después de flush() y antes del UPDATE, la transacción se deshace, las filas siguen pendientes y en la siguiente pasada se vuelven a publicar: el relay es at-least-once, y produce duplicados con el mismo id_evento (el id de la fila de outbox). Esos duplicados los absorbe el consumidor idempotente del apartado 2. Los dos patrones se necesitan mutuamente: outbox garantiza que ningún evento se pierde; idempotencia garantiza que ninguno se aplica dos veces. Y ORDER BY creado_en con la clave de partición agregado_id mantiene el orden de los eventos de cada pedido.

CDC como alternativa al relay. En lugar de un proceso que hace polling a la tabla, herramientas de captura de cambios de datos (Change Data Capture, CDC) como Debezium leen el write-ahead log de PostgreSQL y publican cada inserción en outbox en Kafka con latencia de milisegundos y sin consultas repetidas. Es la implementación de outbox recomendada a escala; el relay por polling es más fácil de entender y suficiente para empezar. Solo lo nombramos: la mecánica del log de replicación pertenece a 03-04.

  1. Petición/respuesta sobre mensajería

A veces se quiere la desconexión temporal de la mensajería pero se necesita una respuesta: un servicio interno que pide a pagos que autorice una tarjeta y quiere el resultado, aunque tarde segundos. El patrón usa dos propiedades del mensaje AMQP:

  • reply_to: el nombre de la cola donde el solicitante espera la respuesta (a menudo una cola exclusiva y temporal creada por él).
  • correlation_id: un identificador que el solicitante pone en la petición y que el servidor copia en la respuesta, para que el solicitante empareje respuestas con peticiones cuando tiene varias en vuelo.
# km0/servicios/pedidos/cliente_pagos_rpc.py  (fragmento: solo el envío y la espera)
import json
import uuid

import pika

conexion = pika.BlockingConnection(pika.ConnectionParameters("localhost"))
canal = conexion.channel()
# Cola de respuestas exclusiva de este proceso; el broker la borra al desconectar
cola_respuestas = canal.queue_declare(queue="", exclusive=True).method.queue
respuestas = {}


def al_responder(ch, metodo, props, cuerpo):
    respuestas[props.correlation_id] = json.loads(cuerpo)


canal.basic_consume(queue=cola_respuestas, on_message_callback=al_responder, auto_ack=True)


def autorizar_tarjeta(pedido_id, importe_centimos, timeout_s=5.0):
    corr_id = str(uuid.uuid4())
    canal.basic_publish(
        exchange="", routing_key="pagos.autorizaciones",       # cola de trabajo de pagos
        properties=pika.BasicProperties(reply_to=cola_respuestas, correlation_id=corr_id,
                                        content_type="application/json"),
        body=json.dumps({"pedido_id": pedido_id, "importe_centimos": importe_centimos}))
    # Espera activa con límite: procesa eventos del canal hasta que llegue la respuesta
    import time
    fin = time.monotonic() + timeout_s
    while corr_id not in respuestas:
        if time.monotonic() > fin:
            raise TimeoutError(f"pagos no respondió a {corr_id} en {timeout_s} s")
        conexion.process_data_events(time_limit=0.1)
    return respuestas.pop(corr_id)

En el lado de pagos, el consumidor procesa la petición y publica la respuesta en props.reply_to con correlation_id=props.correlation_id. (Con BlockingConnection, process_data_events(time_limit=...) es la forma de atender las respuestas que llegan mientras se espera; en un cliente asíncrono el bucle sería distinto, pero la idea es la misma.) Este patrón conserva las ventajas de la mensajería (si pagos está saturado, las peticiones esperan en la cola en lugar de saturarlo más; si se reinicia, no se pierden) a cambio de más complejidad que una llamada gRPC. En Kilómetro Cero se usará poco: para consultas rápidas, gRPC (02-03) es más simple; para consecuencias, eventos puros. Su lugar son las operaciones lentas o con picos donde se quiere respuesta pero no bloquear al servidor.

  1. Consumidores competidores y orden por partición

Cuando inventario no da abasto durante la "Semana del Queso Artesano", la solución es arrancar más instancias que compartan la cola (RabbitMQ) o el grupo (Kafka): consumidores competidores. Cada mensaje lo procesa una sola instancia, y el rendimiento escala con el número de instancias... hasta un límite, que es distinto en cada broker:

  • En RabbitMQ, cualquier número de instancias compite por la misma cola, pero al repartir mensajes entre ellas se pierde el orden: pedido.creado de Ana puede ir a la instancia 1 y pedido.cancelado a la 2, que puede procesarlo antes. Si el orden importa, hay que serializar por otra vía (una cola por clave, o el plugin de consistent hash exchange).
  • En Kafka, el límite es el número de particiones (02-04) y el orden está garantizado dentro de cada partición: todos los eventos de P-2026-000123 los procesa la misma instancia en orden, mientras los de otros pedidos se reparten. Es la razón principal por la que los eventos de dominio de Kilómetro Cero van a Kafka.
flowchart LR
    K[(pedidos.eventos<br/>6 particiones)]
    K -- "p0, p1" --> I1[inventario 1]
    K -- "p2, p3" --> I2[inventario 2]
    K -- "p4, p5" --> I3[inventario 3]
    note["clave P-2026-000123 → siempre p2 → siempre inventario 2 → en orden"]

El precio del orden por partición es el bloqueo de cabeza de línea de toda la partición (el mismo fenómeno que en TCP, 02-01): si el mensaje del offset 41 de la partición 2 tarda diez segundos, los offsets 42 en adelante de esa partición esperan, aunque sean de otros pedidos. De aquí la importancia de que el consumidor no se quede atascado en un mensaje, que es el tema del siguiente apartado.

  1. Mensajes envenenados: reintentos con backoff y colas de mensajes muertos

Un mensaje envenenado (poison message) es uno que hace fallar al consumidor cada vez que lo intenta: JSON malformado, un producto que no existe en stock, un bug en el consumidor con ciertos datos. Con at-least-once y sin más, el broker lo reenvía indefinidamente y el consumidor se queda en bucle, bloqueando su partición o su cola. Hay que distinguir dos clases de fallo:

  • Transitorio: la base de datos no responde, un servicio remoto da timeout. Reintentar tras una espera tiene sentido, y la espera debe crecer (backoff exponencial: 1 s, 2 s, 4 s, 8 s...) para no martillear a un sistema enfermo.
  • Permanente: datos inválidos, violación de una regla de negocio (StockInsuficiente del apartado 2), bug. Reintentar no cambia nada; el mensaje debe apartarse para que el resto avance, y alguien debe examinarlo.

El destino de los mensajes apartados es la cola de mensajes muertos (dead-letter queue, DLQ). En RabbitMQ es nativa: se declara la cola con un dead-letter exchange y, al rechazar sin reencolar, el broker mueve el mensaje allí.

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

def al_recibir(ch, metodo, props, cuerpo):
    try:
        procesar(cuerpo)
        ch.basic_ack(delivery_tag=metodo.delivery_tag)
    except ErrorPermanente:
        ch.basic_nack(delivery_tag=metodo.delivery_tag, requeue=False)   # → DLQ
    except ErrorTransitorio:
        ch.basic_nack(delivery_tag=metodo.delivery_tag, requeue=True)    # vuelve a la cola

(El requeue=True sin límite puede degenerar en el bucle; en RabbitMQ la práctica es contar intentos en una cabecera, x-death, y usar una cola de espera con TTL para el backoff. La idea es la misma que vamos a implementar en Kafka.)

Kafka no tiene DLQ nativa: se construye con tópicos de reintento y un tópico de muertos, y el consumidor decide a cuál enviar cada fallo:

sequenceDiagram
    participant T as pedidos.eventos
    participant C as inventario
    participant R as pedidos.eventos.retry
    participant D as pedidos.eventos.dlq
    participant O as Operador
    T->>C: evento (intento 1)
    C--xC: ErrorTransitorio (BD no responde)
    C->>R: evento + cabeceras {intentos: 1, no_antes: t+1s}
    C->>T: commit (la partición avanza: sin bloqueo)
    R->>C: evento (espera hasta no_antes; intento 2)
    C--xC: ErrorTransitorio
    C->>R: evento {intentos: 2, no_antes: t+2s}
    R->>C: evento (intento 3)
    C--xC: ErrorTransitorio (máximo alcanzado)
    C->>D: evento {intentos: 3, motivo: "..."}
    D->>O: alerta: 1 mensaje en DLQ
    O->>T: tras corregir la causa, republica desde la DLQ
# km0/servicios/inventario/consumo_con_reintentos.py (fragmento)
import json
import time

MAX_INTENTOS = 3
ESPERA_BASE_S = 1.0
TOPICO, RETRY, DLQ = "pedidos.eventos", "pedidos.eventos.retry", "pedidos.eventos.dlq"


class ErrorTransitorio(Exception): ...
class ErrorPermanente(Exception): ...


def cabecera(msg, nombre, por_defecto):
    for k, v in (msg.headers() or []):
        if k == nombre:
            return v.decode()
    return por_defecto


def consumir_con_reintentos(msg, manejador, productor):
    intentos = int(cabecera(msg, "intentos", "0"))
    try:
        manejador(json.loads(msg.value()))
    except ErrorTransitorio as e:
        if intentos + 1 < MAX_INTENTOS:
            espera = ESPERA_BASE_S * (2 ** intentos)                # 1, 2, 4 segundos
            productor.produce(RETRY, key=msg.key(), value=msg.value(), headers=[
                ("intentos", str(intentos + 1)), ("no_antes", str(time.time() + espera)),
                ("ultimo_error", str(e)[:200])])
        else:
            productor.produce(DLQ, key=msg.key(), value=msg.value(), headers=[
                ("intentos", str(intentos + 1)), ("motivo", f"transitorio agotado: {e}"[:200])])
    except ErrorPermanente as e:
        productor.produce(DLQ, key=msg.key(), value=msg.value(), headers=[
            ("intentos", str(intentos + 1)), ("motivo", f"permanente: {e}"[:200])])
    productor.flush()
    # En todos los casos se confirma el offset del tópico de origen: el mensaje ya
    # está a salvo en retry o en dlq, y la partición no se bloquea.

Un consumidor del tópico .retry lee cada mensaje, duerme hasta no_antes si hace falta, y lo procesa con la misma función. Consideraciones que no se ven en el código: el mensaje que va a .retry pierde su posición en el orden de la partición original (se procesará después de otros más recientes), lo que para un descuento de stock es aceptable y para una secuencia creado → cancelado puede no serlo; la DLQ necesita alertas (07-01) y un procedimiento para examinar, corregir y republicar; y StockInsuficiente del apartado 2 es un buen ejemplo de error permanente cuya resolución no es técnica sino de negocio (avisar al cliente, cancelar el pedido), que es el territorio de las sagas de 03-05. Los reintentos, backoff y circuit breakers de las llamadas síncronas se tratan en 07-04; aquí solo hemos visto los que pertenecen al consumo de mensajes.

  1. Event-driven y event-sourcing: dos cosas distintas

Todo lo que hemos hecho en 02-04 y en esta lección es arquitectura dirigida por eventos (event-driven): los servicios comunican hechos ocurridos (pedido.creado) y otros reaccionan. El estado de cada servicio sigue viviendo en sus tablas (pedidos, stock); los eventos son notificaciones.

Event-sourcing es otra cosa, que a menudo se confunde con la anterior: consiste en que el estado no se guarda; se guardan los eventos, y el estado se reconstruye reproduciéndolos. El pedido de Ana no sería una fila con estado = 'pagado', sino la secuencia PedidoCreado, LineaAñadida, PedidoPagado, y la fila sería una proyección derivada. Da auditoría completa y permite reconstruir cualquier estado pasado, a cambio de complejidad considerable (proyecciones, instantáneas, evolución de eventos históricos). Se puede hacer event-driven sin event-sourcing (Kilómetro Cero lo hace) y viceversa. Lo nombramos para que el término no se confunda; su desarrollo excede este curso.

  1. Esquemas de eventos y versionado

Un evento es un contrato entre el productor y todos los consumidores presentes y futuros, incluidos los que leerán el histórico de Kafka dentro de seis días. Se le aplican las reglas de evolución de esquemas de 02-03 con más severidad, porque no se puede saber cuántos consumidores hay ni obligarlos a actualizarse. La envoltura estándar que hemos usado ya prevé el mecanismo:

{
  "id_evento": "6f1c9e2a-...",
  "tipo": "pedido.creado",
  "version": 2,
  "fecha_ms": 1789000000000,
  "origen": "pedidos",
  "datos": {
    "id": "P-2026-000123",
    "cliente": "ana",
    "lineas": [{"producto": "queso-curado", "cantidad": 2, "precio_centimos": 1450}],
    "direccion_entrega": {"ciudad": "Girona", "codigo_postal": "17001"}
  }
}

Las reglas, adaptadas de 02-03:

  • Cambios compatibles (no incrementan version): añadir campos opcionales a datos; añadir tipos de evento nuevos. Los consumidores deben ignorar campos desconocidos y tolerar la ausencia de los nuevos.
  • Cambios incompatibles (incrementan version): eliminar o renombrar un campo, cambiar su tipo o su significado. Durante la transición, el productor puede publicar ambas versiones o los consumidores aceptar ambas (if evento["version"] == 1: ...). Nunca se cambia el significado de un campo conservando su nombre.
  • Registro de esquemas (Schema Registry): con Avro o protobuf en Kafka, un servicio central guarda cada versión del esquema de cada tópico, valida en el productor que un cambio es compatible (hacia atrás, hacia delante o ambas, según la política) y permite a los consumidores deserializar mensajes antiguos con el esquema con el que se escribieron. Es la forma de convertir las reglas anteriores en una comprobación automática, y el paso natural cuando JSON se queda corto. Kilómetro Cero empieza con JSON y la envoltura; el registro de esquemas se introduce con el pipeline de datos del Módulo 5.

  1. Tabla de patrones: qué problema resuelve cada uno

Patrón Problema que resuelve Coste Dónde se usa en Kilómetro Cero
At-least-once + ack tras procesar Pérdida de mensajes al morir el consumidor Duplicados Todos los consumidores de eventos de dominio
Consumidor idempotente (tabla mensajes_procesados) Duplicados por reintento del productor, del relay o del broker Una fila por mensaje; purga periódica inventario (descuento de stock y ReservarStock gRPC), pagos (cobros)
Outbox transaccional + relay Escritura dual: estado guardado sin evento, o evento sin estado Un proceso más; latencia de milisegundos a segundos pedidos (todos sus eventos); después, cada servicio que publique
CDC (Debezium) Relay por polling a escala Infraestructura adicional Cuando el polling no baste
Petición/respuesta (reply_to, correlation_id) Necesitar respuesta con la amortiguación de una cola Complejidad frente a gRPC Operaciones lentas con picos (autorizaciones de pagos)
Consumidores competidores Un consumidor no da abasto Pérdida de orden (RabbitMQ) o límite por particiones (Kafka) Todos los servicios en campaña
Clave de partición Orden entre eventos de una misma entidad Bloqueo de cabeza de línea por partición; particiones desequilibradas pedidos.eventos por id de pedido; telemetría por repartidor
Reintentos con backoff en consumo Fallos transitorios del consumidor Pérdida de orden del mensaje reintentado Todos los consumidores
Dead-letter queue Mensajes envenenados que bloquean la cola o la partición Necesita alertas y procedimiento de reproceso Todos los consumidores
Envoltura + versionado de eventos Evolución del contrato sin coordinar despliegues Disciplina y revisión de cambios Todos los eventos

Errores Comunes y Consejos

  • Buscar exactly-once en la configuración del broker. No está ahí. Está en el consumidor idempotente, y no hay atajo.
  • Registrar la clave de idempotencia fuera de la transacción del efecto. Reabre la ventana que se quería cerrar. Misma transacción, siempre.
  • Usar como clave de idempotencia algo que cambia entre reintentos. Si el productor genera un id_evento nuevo en cada reintento, el consumidor no puede reconocer el duplicado. La clave se genera una vez y se reutiliza.
  • Publicar en Kafka "dentro" de la transacción de PostgreSQL. No la deshace el ROLLBACK. Outbox o nada.
  • Un relay que marca como publicado antes del flush. Si Kafka no confirma, el evento se pierde para siempre con la fila marcada. Confirmación primero, marca después.
  • Reintentar errores permanentes. Un JSON malformado no se arregla esperando 4 segundos; bloquea la partición tres veces más. Clasifica los errores y manda los permanentes a la DLQ a la primera.
  • Una DLQ sin alertas ni dueño. Es un agujero negro donde los pedidos desaparecen en silencio. Cada DLQ necesita una alerta, una persona y un procedimiento.
  • Cambiar el significado de un campo sin cambiar la versión. Un consumidor que releerá el histórico leerá datos con dos significados bajo el mismo nombre. Campo nuevo o versión nueva.
  • Consejo: prueba la idempotencia de forma explícita: en el entorno de pruebas, publica cada evento dos veces y comprueba que el estado final es el mismo. Es la prueba más barata y la que más incidentes evita.
  • Consejo: mide el lag del relay (filas de outbox con publicado_en IS NULL y su antigüedad) y de cada grupo de consumidores. Son los dos indicadores que dicen si la asincronía está funcionando o acumulando deuda invisible.

Ejercicios

Ejercicio 1: Auditar un consumidor

Este consumidor de pagos procesa pedido.creado y ejecuta el cobro contra la pasarela externa. Identifica todos los problemas de garantías de entrega que tiene, indica qué puede salir mal en cada uno (pérdida, duplicado, bloqueo) y reescríbelo aplicando los patrones de la lección. Asume que la pasarela acepta una cabecera Idempotency-Key y que pagos tiene su propia base de datos PostgreSQL.

consumidor = Consumer({"bootstrap.servers": "localhost:9092", "group.id": "pagos",
                       "enable.auto.commit": True})
consumidor.subscribe(["pedidos.eventos"])
while True:
    msg = consumidor.poll(1.0)
    if msg is None or msg.error():
        continue
    evento = json.loads(msg.value())
    if evento["tipo"] == "pedido.creado":
        importe = sum(l["cantidad"] * l["precio_centimos"] for l in evento["datos"]["lineas"])
        respuesta = pasarela.cobrar(tarjeta=evento["datos"]["tarjeta"], importe=importe)
        with psycopg.connect(DSN) as conn:
            conn.execute("INSERT INTO pagos (pedido_id, importe, ref_pasarela) VALUES (%s, %s, %s)",
                         (evento["datos"]["id"], importe, respuesta["ref"]))

Ejercicio 2: Outbox para inventario

inventario debe publicar stock.actualizado (con producto, antes, despues, motivo) cada vez que cambia el stock, tanto por la reserva gRPC de 02-03 como por el consumo de pedido.creado del apartado 2. Diseña la tabla outbox de inventario, modifica procesar_pedido_creado para escribir el evento en la misma transacción, y explica qué garantiza el conjunto (idempotencia + outbox) ante cada uno de estos fallos: (a) el consumidor muere entre el COMMIT y el commit de offset; (b) el relay muere tras publicar y antes de marcar; (c) Kafka no está disponible durante 10 minutos.

Ejercicio 3: Diseñar la política de reintentos y DLQ

Para el consumidor de reparto que asigna un repartidor a cada pedido.pagado, clasifica cada uno de estos fallos como transitorio o permanente, indica a dónde va el mensaje (reintento con backoff, DLQ) y qué hace el operador en cada caso: (1) el servicio de mapas externo devuelve 503; (2) el pedido tiene una dirección de entrega en una ciudad donde Kilómetro Cero no reparte; (3) KeyError: 'direccion_entrega' en un evento de version: 1; (4) la base de datos de reparto rechaza la conexión durante un reinicio; (5) no hay ningún repartidor libre en Tarragona ahora mismo.

Soluciones

Solución 1:

Problemas: (1) enable.auto.commit=True confirma offsets periódicamente con independencia de si se ha procesado: si el proceso muere tras el auto-commit y antes de cobrar, el cobro se pierde (at-most-once); si muere después de cobrar y antes del auto-commit, se cobra dos veces (at-least-once sin idempotencia). (2) El cobro a la pasarela y el INSERT son dos escrituras sin transacción común: si la pasarela cobra y el INSERT falla, hay un cobro sin registro; al reintentarse, otro cobro. (3) No hay clave de idempotencia ni hacia la pasarela ni en la base de datos: el mismo pedido.creado entregado dos veces cobra dos veces a Ana. (4) msg.error() se ignora en silencio. (5) Cualquier excepción (JSON malformado, tarjeta inválida) mata el bucle o, si se capturara, bloquearía la partición. (6) No publica pago.confirmado, así que nadie se entera del cobro (escritura dual pendiente). Reescritura:

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

def cobrar_pedido(conn, evento):
    pedido = evento["datos"]
    with conn.transaction():
        with conn.cursor() as cur:
            cur.execute("INSERT INTO mensajes_procesados (id_mensaje, consumidor) VALUES (%s, 'pagos.cobro') "
                        "ON CONFLICT DO NOTHING RETURNING id_mensaje", (evento["id_evento"],))
            if cur.fetchone() is None:
                return "duplicado"
            importe = sum(l["cantidad"] * l["precio_centimos"] for l in pedido["lineas"])
            # La pasarela deduplica por Idempotency-Key: usar el id del pedido hace que
            # un reintento (incluso con otro id_evento) no cobre dos veces.
            respuesta = pasarela.cobrar(tarjeta=pedido["tarjeta"], importe=importe,
                                        idempotency_key=f"cobro-{pedido['id']}")
            cur.execute("INSERT INTO pagos (pedido_id, importe, ref_pasarela) VALUES (%s, %s, %s)",
                        (pedido["id"], importe, respuesta["ref"]))
            cur.execute("INSERT INTO outbox (id, agregado, agregado_id, tipo, carga) "
                        "VALUES (%s, 'pago', %s, 'pago.confirmado', %s)",
                        (uuid.uuid4(), pedido["id"], json.dumps({"pedido_id": pedido["id"], "importe": importe})))
    return "procesado"

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:
            evento = json.loads(msg.value())
        except json.JSONDecodeError as e:
            a_dlq(msg, f"permanente: {e}"); consumidor.commit(message=msg); continue
        if evento["tipo"] == "pedido.creado":
            try:
                cobrar_pedido(conn, evento)
            except PasarelaNoDisponible as e:
                a_retry(msg, e)                  # transitorio: backoff
            except TarjetaRechazada as e:
                a_dlq(msg, f"permanente: {e}")   # de negocio: pedido a cancelar (03-05)
        consumidor.commit(message=msg)

Queda una ventana inevitable: si el proceso muere entre pasarela.cobrar y el COMMIT, la transacción se deshace (ni clave ni registro), el mensaje se reintenta, y el segundo cobrar con la misma Idempotency-Key hace que la pasarela devuelva el cobro ya hecho en lugar de repetirlo. La idempotencia del sistema externo cierra lo que la transacción local no puede cubrir. Sin esa cabecera, habría que registrar "cobro iniciado" antes de llamar y reconciliar después: es la situación en la que las sagas de 03-05 se hacen necesarias.

Solución 2:

CREATE TABLE outbox (
    id UUID PRIMARY KEY, agregado TEXT NOT NULL, agregado_id TEXT NOT NULL,
    tipo TEXT NOT NULL, version INTEGER NOT NULL DEFAULT 1, carga JSONB NOT NULL,
    creado_en TIMESTAMPTZ NOT NULL DEFAULT now(), publicado_en TIMESTAMPTZ);

Con agregado = 'producto' y agregado_id = producto, para que los cambios de un mismo producto vayan en orden a la misma partición de inventario.eventos. En procesar_pedido_creado, dentro del mismo with conn.transaction(), tras cada UPDATE con éxito:

cur.execute("SELECT unidades FROM stock WHERE producto = %s", (linea["producto"],))
despues = cur.fetchone()[0]
cur.execute("INSERT INTO outbox (id, agregado, agregado_id, tipo, carga) VALUES (%s, 'producto', %s, 'stock.actualizado', %s)",
            (uuid.uuid4(), linea["producto"],
             json.dumps({"producto": linea["producto"], "antes": despues + linea["cantidad"],
                         "despues": despues, "motivo": "pedido", "pedido_id": evento["datos"]["id"]})))

(Y lo mismo en ReservarStock del servidor gRPC, que dejaría de usar el diccionario en memoria para usar la misma base de datos.) Garantías: (a) el consumidor muere entre COMMIT y commit de offset: el stock está descontado y el evento está en outbox; el relay lo publicará; el mensaje reenviado se detecta como duplicado y no genera ni descuento ni segundo evento. (b) El relay muere tras publicar y antes de marcar: el evento se publica dos veces con el mismo id_evento; catalogo y analitica, si son idempotentes, lo ignoran la segunda vez; si catalogo solo hace SET indicador = (despues < 5), es idempotente por naturaleza y ni siquiera necesita la tabla. (c) Kafka caído 10 minutos: inventario sigue procesando (si lee de RabbitMQ) o se detiene (si lee de Kafka), pero en ambos casos no pierde nada: los eventos se acumulan en outbox con publicado_en IS NULL y el relay los publica en orden cuando Kafka vuelve. Es exactamente el comportamiento que se pedía en la solución 3.3 de la lección 01-06, ahora implementado.

Solución 3:

  1. Transitorio: 503 es "vuelve luego". Reintento con backoff (1, 2, 4 s...) hasta el máximo; después, DLQ. El operador comprueba el estado del proveedor de mapas; al recuperarse, republica la DLQ.
  2. Permanente, de negocio: no hay reintento que lo arregle. DLQ con motivo ciudad_no_cubierta, y probablemente un evento reparto.rechazado para que pedidos cancele y pagos reembolse (una saga, 03-05). El operador o, mejor, la validación en pedidos al crear el pedido, debería haberlo impedido antes: la DLQ revela un hueco en la validación.
  3. Permanente, de contrato: el consumidor asume el esquema version: 2 y recibe uno de version: 1 (probablemente del histórico o de un pedidos aún no actualizado). DLQ con motivo esquema, alerta al equipo, y la corrección es en el consumidor (aceptar ambas versiones, apartado 8); tras desplegarla, se republican los mensajes de la DLQ. No es un fallo de datos sino de un consumidor que no siguió las reglas de compatibilidad.
  4. Transitorio: reinicio de la base de datos. Reintento con backoff. Aquí es especialmente importante no confirmar el offset sin haber apartado el mensaje a .retry; con la base de datos caída, ni siquiera se puede registrar la idempotencia, así que el mensaje debe quedar íntegro para reintentarse.
  5. Transitorio, pero de negocio y de larga duración: puede tardar 20 minutos en haber un repartidor libre. Un backoff de segundos no encaja; es mejor que el consumidor procese el mensaje registrando el pedido como "pendiente de asignación" en su base de datos (idempotentemente) y confirme, y que un proceso periódico intente asignar los pendientes. Convertir una espera larga en estado persistido es preferible a mantener mensajes rebotando entre tópicos de reintento, que además perderían el orden respecto a un posible pedido.cancelado posterior.

Conclusión

Esta lección ha cerrado el módulo resolviendo los problemas que las anteriores fueron dejando abiertos. Las garantías de entrega dependen de dónde se confirma respecto a dónde se procesa: confirmar antes pierde mensajes (at-most-once), confirmar después los duplica (at-least-once), y exactly-once no existe de extremo a extremo, así que la estrategia es at-least-once en el transporte e idempotencia en el consumidor. El consumidor idempotente de inventario, con su tabla mensajes_procesados en la misma transacción que el efecto, ha cerrado el problema de los duplicados que arrastrábamos desde la simulación de 01-04 y las semánticas de 02-02, y la misma técnica hace seguros los reintentos de la llamada gRPC de 02-03. El patrón Outbox, con su relay FOR UPDATE SKIP LOCKED, ha eliminado la escritura dual en pedidos: estado y evento se guardan juntos o no se guardan. Hemos visto además petición/respuesta sobre colas con reply_to y correlation_id, consumidores competidores con orden por partición, reintentos con backoff, colas de mensajes muertos con sus alertas y su procedimiento, la diferencia entre event-driven y event-sourcing, y el versionado de eventos con la envoltura estándar y las reglas heredadas de 02-03. La tabla del apartado 9 resume qué patrón resuelve qué problema y dónde lo usa Kilómetro Cero.

Con esto, el Módulo 2 ha cumplido su promesa: pedidos e inventario son dos procesos que hablan de forma fiable, síncronamente por gRPC cuando hace falta respuesta y asíncronamente por Kafka para las consecuencias, y analitica escucha sin que nadie tenga que esperarla. Pero fíjate en lo que ha ocurrido por el camino: inventario tiene ahora su propia tabla stock y pedidos su propia tabla pedidos, y la verdad sobre el último queso curado de la Quesería Montblanc ya no está en un único sitio. Cuando Ana lo reserva y catalogo recibe el evento stock.actualizado medio segundo después, durante ese medio segundo Marc ve en el catálogo un queso que ya no existe. Cuando inventario tenga las réplicas inv-bcn e inv-vlc, ¿cuál de las dos tiene razón si difieren? ¿Qué significa exactamente "consistente" cuando los datos viven en varios lugares, y qué se puede prometer a un cliente ante un fallo de red entre ellos? Estas son las preguntas del Módulo 3, Consistencia y Replicación, que empieza con los modelos de consistencia: el vocabulario preciso para decir qué garantiza un sistema distribuido sobre sus datos y qué no.

Curso de Arquitecturas Distribuidas

Módulo 1: Introducción a los Sistemas Distribuidos

Módulo 2: Comunicación en Sistemas Distribuidos

Módulo 3: Consistencia y Replicación

Módulo 4: Almacenamiento Distribuido

Módulo 5: Computación Distribuida

Módulo 6: Seguridad en Sistemas Distribuidos

Módulo 7: Monitoreo y Mantenimiento

Módulo 8: Casos de Estudio y Aplicaciones

© Copyright 2026. Todos los derechos reservados