Hasta ahora, toda la comunicación de Kilómetro Cero ha sido síncrona: pedidos llama a inventario por gRPC, espera y recibe. Es lo correcto cuando la respuesta hace falta en ese instante, y es lo incorrecto para todo lo demás. Recordemos el segundo síntoma del monolito (lección 01-06): la pasarela de pagos externa se quedó lenta y, como el cobro estaba dentro de la misma transacción que la reserva de stock y la creación del pedido, cada pedido bloqueado retenía conexiones y hilos hasta tumbar la plataforma entera. Convertir esa cadena en llamadas gRPC no la arregla: seguiría siendo una cadena de disponibilidades multiplicadas y latencias sumadas. Lo que hace falta es que pedidos pueda decir "ha ocurrido esto" y seguir adelante, y que quien deba reaccionar lo haga cuando pueda.

Esa es la función de la mensajería: un intermediario (el broker) que guarda los mensajes y los entrega a sus destinatarios, desacoplando a quien envía de quien recibe en el tiempo, en el espacio y en el ritmo. En esta lección fijaremos los conceptos (productor, consumidor, cola, tópico, ack, persistencia), los dos modelos (punto a punto y publicación/suscripción), y las dos tecnologías que dominan el sector, RabbitMQ y Apache Kafka, con sus diferencias de fondo. Construiremos el primer flujo asíncrono de Kilómetro Cero: pedidos publica el evento pedido.creado, e inventario y analitica lo consumen cada uno a su ritmo. Las garantías de entrega y los patrones que las hacen seguras (idempotencia, outbox, colas de mensajes muertos) se nombran aquí y se desarrollan en la lección 02-05.

Contenido

  1. Por qué desacoplar en el tiempo
  2. Vocabulario de la mensajería
  3. Dos modelos: punto a punto y publicación/suscripción
  4. RabbitMQ y AMQP: exchanges, colas, bindings y routing keys
  5. pedido.creado con RabbitMQ y pika
  6. Apache Kafka: el log distribuido
  7. pedido.creado con Kafka
  8. RabbitMQ frente a Kafka: criterios de elección
  9. Ampliar el docker-compose.yml de km0/
  10. Errores comunes y consejos
  11. Ejercicios
  12. Conclusión

  1. Por qué desacoplar en el tiempo

Una llamada síncrona acopla a los dos participantes de tres maneras:

  • En el tiempo: ambos deben estar vivos y disponibles en el mismo instante. Si analitica está reiniciándose cuando pedidos intenta notificarle una venta, la notificación falla o pedidos espera.
  • En el espacio: el emisor necesita saber quién es el receptor y dónde está (dirección, puerto). Si mañana reparto también quiere saber de los pedidos, hay que cambiar pedidos.
  • En el ritmo: el emisor no puede ir más rápido que el receptor más lento. En la "Semana de la Vendimia", pedidos produce 1.200 pedidos por segundo; si analitica solo procesa 300, pedidos se ralentiza hasta 300.

Un broker rompe los tres acoplamientos: pedidos entrega el mensaje al broker (que sí está disponible, porque es infraestructura replicada) y sigue; el broker lo guarda; los consumidores lo recogen cuando quieren y al ritmo que pueden; y añadir un consumidor nuevo no toca al productor.

flowchart LR
    subgraph Antes["Síncrono: cadena de dependencias"]
        P1[pedidos] --> I1[inventario]
        P1 --> PA1[pagos]
        P1 --> A1[analitica]
        P1 --> R1[reparto]
    end
    subgraph Despues["Asíncrono: el broker en medio"]
        P2[pedidos] -- pedido.creado --> B[(broker)]
        B --> I2[inventario]
        B --> A2[analitica]
        B --> R2[reparto]
    end

El coste es igual de claro y hay que aceptarlo con los ojos abiertos: pedidos ya no sabe si el mensaje se ha procesado, ni cuándo. Las decisiones que requieren respuesta inmediata (¿hay stock?) siguen siendo síncronas; lo asíncrono es para las consecuencias de una decisión ya tomada. Y aparece una pieza de infraestructura crítica más, con su propia disponibilidad y su propio modelo de fallos.

  1. Vocabulario de la mensajería

Término Qué es En Kilómetro Cero
Mensaje Unidad de datos que viaja: cabeceras (metadatos) + cuerpo (bytes serializados, 02-03) Un evento pedido.creado con el pedido de Ana en el cuerpo
Productor (producer, publisher) Quien envía mensajes al broker pedidos
Consumidor (consumer, subscriber) Quien recibe mensajes del broker inventario, analitica, reparto
Broker El servidor intermediario que recibe, almacena y entrega RabbitMQ o Kafka
Cola (queue) Almacén ordenado del que los consumidores extraen mensajes; cada mensaje se entrega normalmente a un consumidor inventario.pedidos en RabbitMQ
Tópico (topic) Canal con nombre al que se publican mensajes y al que se suscriben interesados; cada mensaje puede llegar a varios pedidos.eventos en Kafka
Ack (acknowledgement) Confirmación del consumidor al broker de que ha procesado el mensaje; hasta entonces el broker lo conserva basic_ack en RabbitMQ, commit de offset en Kafka
Persistencia El broker escribe el mensaje en disco, no solo en memoria Sobrevive a un reinicio del broker
Durabilidad La cola o el tópico sobrevive al reinicio del broker (además de sus mensajes, si son persistentes) Las colas de inventario deben ser durables
Retención Cuánto tiempo (o cuántos bytes) conserva el broker los mensajes RabbitMQ: hasta el ack; Kafka: días, aunque ya se hayan leído
Prefetch Cuántos mensajes puede tener un consumidor pendientes de ack a la vez 1 para procesamiento seguro, más para rendimiento

  1. Dos modelos: punto a punto y publicación/suscripción

Punto a punto (cola de trabajo): los productores dejan mensajes en una cola y uno o varios consumidores los extraen; cada mensaje lo procesa exactamente un consumidor. Es el modelo para repartir trabajo: generar las facturas en PDF de los pedidos, enviar correos, calcular rutas. Añadir consumidores aumenta el ritmo de procesamiento (consumidores competidores, que trataremos en 02-05).

Publicación/suscripción (pub/sub): los productores publican en un tópico sin saber quién escucha; cada suscriptor recibe su propia copia de cada mensaje. Es el modelo para eventos: "ha ocurrido un pedido" interesa a inventario, a analitica y a reparto, y cada uno hace algo distinto con él.

flowchart LR
    subgraph PP["Punto a punto"]
        Pr1[productor] --> Q[(cola facturas)]
        Q --> C1[consumidor A]
        Q --> C2[consumidor B]
        note1[cada mensaje va a UNO de los dos]
    end
    subgraph PS["Publicación/suscripción"]
        Pr2[pedidos] --> T[(tópico pedido.creado)]
        T --> S1[inventario]
        T --> S2[analitica]
        T --> S3[reparto]
        note2[cada mensaje va a TODOS]
    end

En la práctica los dos modelos se combinan: cada suscriptor de un tópico suele ser un grupo de instancias que compiten entre sí por los mensajes de ese suscriptor (tres réplicas de analitica se reparten los eventos, pero analitica como conjunto recibe todos). Tanto RabbitMQ como Kafka soportan esta combinación, con mecanismos distintos.

  1. RabbitMQ y AMQP: exchanges, colas, bindings y routing keys

RabbitMQ es un broker que implementa AMQP 0-9-1 (Advanced Message Queuing Protocol), un protocolo binario sobre TCP con un modelo de enrutamiento muy flexible. La clave para entenderlo es que los productores nunca publican en colas: publican en un exchange, y el exchange decide a qué colas copiar el mensaje según unas reglas (bindings) y una etiqueta del mensaje (routing key).

flowchart LR
    P[pedidos] -- "routing key: pedido.creado" --> X{{exchange km0.pedidos<br/>tipo topic}}
    X -- "binding: pedido.creado" --> Q1[(cola inventario.pedidos)]
    X -- "binding: pedido.#" --> Q2[(cola analitica.pedidos)]
    X -- "binding: pedido.pagado" --> Q3[(cola reparto.pedidos)]
    Q1 --> I[inventario]
    Q2 --> A[analitica]
    Q3 --> R[reparto]

Los cuatro tipos de exchange:

Tipo Regla de enrutamiento Uso
direct La routing key del mensaje debe ser igual a la del binding Colas de trabajo con destino explícito
fanout Ignora la routing key: copia a todas las colas enlazadas Difusión pura
topic Routing key con puntos (pedido.creado); los bindings usan comodines: * (una palabra), # (cero o más) Eventos clasificados: pedido.*, reparto.furgoneta-3.#
headers Enruta por cabeceras del mensaje en lugar de por routing key Casos especiales

Otros conceptos de RabbitMQ: cada consumidor mantiene una conexión TCP y dentro de ella uno o más canales (multiplexación ligera); los mensajes se entregan al consumidor mediante push y quedan "sin confirmar" hasta el ack; si el consumidor muere sin confirmar, el broker reencola el mensaje para otro consumidor; y una vez confirmado, el mensaje desaparece de la cola. Esa última propiedad es la diferencia de fondo con Kafka.

  1. pedido.creado con RabbitMQ y pika

pika es el cliente Python oficial de RabbitMQ (pip install pika). Primero, el productor en pedidos:

# km0/servicios/pedidos/publicador_rabbit.py
import json
import time
import uuid

import pika

EXCHANGE = "km0.pedidos"


def conectar():
    conexion = pika.BlockingConnection(pika.ConnectionParameters("localhost"))
    canal = conexion.channel()
    # Declarar es idempotente: si el exchange existe con los mismos parámetros, no pasa nada.
    # durable=True: el exchange sobrevive a reinicios del broker.
    canal.exchange_declare(exchange=EXCHANGE, exchange_type="topic", durable=True)
    return conexion, canal


def publicar_pedido_creado(canal, pedido):
    evento = {
        "id_evento": str(uuid.uuid4()),     # identidad del mensaje (clave en 02-05)
        "tipo": "pedido.creado",
        "version": 1,                        # versión del esquema del evento
        "fecha_ms": int(time.time() * 1000),
        "origen": "pedidos",
        "datos": pedido,
    }
    canal.basic_publish(
        exchange=EXCHANGE,
        routing_key="pedido.creado",
        body=json.dumps(evento).encode("utf-8"),
        properties=pika.BasicProperties(
            content_type="application/json",
            message_id=evento["id_evento"],
            delivery_mode=2,                 # 2 = persistente: se escribe en disco
        ),
    )
    print(f"[pedidos] publicado pedido.creado {pedido['id']}")


if __name__ == "__main__":
    conexion, canal = conectar()
    pedido_ana = {
        "id": "P-2026-000123", "cliente": "ana", "mercado": "girona",
        "lineas": [{"producto": "queso-curado", "cantidad": 2, "precio_centimos": 1450},
                   {"producto": "tomate-rosa", "cantidad": 3, "precio_centimos": 320}],
    }
    publicar_pedido_creado(canal, pedido_ana)
    conexion.close()

Observa la envoltura del evento: además de los datos del pedido lleva un identificador único, un tipo, una versión, una fecha y un origen. Es un convenio que mantendremos en todos los eventos de Kilómetro Cero, y cada campo tiene un uso que aparecerá en 02-05 (el id_evento para detectar duplicados, la version para evolucionar el esquema).

Ahora el consumidor de inventario. Su trabajo, en esta primera versión, es descontar el stock reservado (más adelante decidiremos si la reserva síncrona por gRPC y el descuento por evento conviven o se sustituyen; por ahora, lo importante es el mecanismo):

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

import pika

EXCHANGE = "km0.pedidos"
COLA = "inventario.pedidos"

stock = {"tomate-rosa": 120, "calabacin": 80, "queso-curado": 5,
         "queso-fresco": 30, "vino-crianza": 200}


def al_recibir(canal, metodo, propiedades, cuerpo):
    evento = json.loads(cuerpo)
    pedido = evento["datos"]
    for linea in pedido["lineas"]:
        stock[linea["producto"]] -= linea["cantidad"]
    print(f"[inventario] procesado {evento['tipo']} {pedido['id']} "
          f"(evento {evento['id_evento'][:8]}); queso-curado={stock['queso-curado']}")
    # Ack DESPUÉS de procesar: si morimos antes, el broker reencola el mensaje
    canal.basic_ack(delivery_tag=metodo.delivery_tag)


conexion = pika.BlockingConnection(pika.ConnectionParameters("localhost"))
canal = conexion.channel()
canal.exchange_declare(exchange=EXCHANGE, exchange_type="topic", durable=True)
canal.queue_declare(queue=COLA, durable=True)                    # la cola sobrevive al reinicio
canal.queue_bind(queue=COLA, exchange=EXCHANGE, routing_key="pedido.creado")
canal.basic_qos(prefetch_count=1)     # un mensaje sin confirmar a la vez por consumidor
canal.basic_consume(queue=COLA, on_message_callback=al_recibir)
print("[inventario] esperando pedido.creado en", COLA)
canal.start_consuming()

Y el de analitica, que quiere todos los eventos de pedido, no solo los de creación, para acumular estadísticas:

# km0/servicios/analitica/consumidor_rabbit.py
import json
from collections import Counter

import pika

ventas_por_producto = Counter()


def al_recibir(canal, metodo, propiedades, cuerpo):
    evento = json.loads(cuerpo)
    if evento["tipo"] == "pedido.creado":
        for linea in evento["datos"]["lineas"]:
            ventas_por_producto[linea["producto"]] += linea["cantidad"]
    print(f"[analitica] {evento['tipo']} -> {dict(ventas_por_producto)}")
    canal.basic_ack(delivery_tag=metodo.delivery_tag)


conexion = pika.BlockingConnection(pika.ConnectionParameters("localhost"))
canal = conexion.channel()
canal.exchange_declare(exchange="km0.pedidos", exchange_type="topic", durable=True)
canal.queue_declare(queue="analitica.pedidos", durable=True)
canal.queue_bind(queue="analitica.pedidos", exchange="km0.pedidos", routing_key="pedido.#")
canal.basic_qos(prefetch_count=50)    # analitica tolera lotes: más rendimiento
canal.basic_consume(queue="analitica.pedidos", on_message_callback=al_recibir)
canal.start_consuming()

Puntos importantes de los tres programas:

  • Cada consumidor declara su propia cola y la enlaza al exchange con el patrón que le interesa. inventario solo quiere pedido.creado; analitica quiere pedido.# (creado, pagado, cancelado...). El productor no sabe nada de ninguna de las dos colas: eso es el desacoplamiento en el espacio.
  • El mensaje se copia a cada cola enlazada: inventario y analitica reciben cada uno su ejemplar (pub/sub). Si arrancas dos instancias de consumidor_rabbit.py de inventario, compartirán la cola inventario.pedidos y cada mensaje irá a una de ellas (punto a punto dentro del suscriptor).
  • Arranca los consumidores después de publicar y verás que reciben el mensaje igualmente: estaba guardado en la cola. Pero si la cola no existía cuando se publicó (porque el consumidor nunca había arrancado), el mensaje se descartó: el exchange no tenía dónde copiarlo. Por eso las colas de los consumidores críticos deben declararse en el despliegue, no en el arranque del primer consumidor.
  • durable=True + delivery_mode=2 son las dos mitades de la supervivencia a un reinicio del broker: la cola y el mensaje. Una sin la otra no sirve.
  • prefetch_count=1 en inventario: el broker no le envía el siguiente mensaje hasta que confirme el actual, lo que evita que una instancia acapare mensajes que no puede procesar si muere. Cuesta rendimiento; analitica, que puede permitirse reprocesar lotes, usa 50.

Con la interfaz web de administración (http://localhost:15672, usuario y contraseña km0) se ven exchanges, colas, mensajes pendientes y consumidores conectados: es la primera herramienta de diagnóstico.

  1. Apache Kafka: el log distribuido

Kafka nació en LinkedIn (2011) con una idea distinta: en lugar de una cola de la que los mensajes desaparecen al consumirse, un log, es decir, un fichero de solo añadir (append-only) en el que cada mensaje ocupa una posición fija (offset) y permanece durante un tiempo de retención (por defecto 7 días) independientemente de que alguien lo haya leído. Los consumidores no extraen mensajes: leen el log desde una posición y recuerdan por dónde van.

flowchart LR
    subgraph T["tópico pedidos.eventos (3 particiones)"]
        P0["partición 0: [0][1][2][3][4] →"]
        P1["partición 1: [0][1][2] →"]
        P2["partición 2: [0][1][2][3] →"]
    end
    Pr[pedidos<br/>clave = id del pedido] --> T
    subgraph G1["grupo inventario"]
        C1[instancia 1] -.- P0
        C1 -.- P1
        C2[instancia 2] -.- P2
    end
    subgraph G2["grupo analitica"]
        C3[instancia única] -.- P0
        C3 -.- P1
        C3 -.- P2
    end

Los conceptos que hay que dominar:

  • Tópico: el nombre lógico del flujo (pedidos.eventos).
  • Partición: cada tópico se divide en N logs independientes. Es la unidad de paralelismo (cada partición la lee una sola instancia de cada grupo) y de orden (los mensajes están ordenados dentro de una partición, no entre particiones).
  • Clave de partición: el productor puede asignar una clave a cada mensaje; Kafka calcula hash(clave) mod N para elegir partición. Todos los mensajes con la misma clave van a la misma partición y, por tanto, se leen en orden. Usando el identificador del pedido como clave, pedido.creado, pedido.pagado y pedido.cancelado de P-2026-000123 llegan siempre en ese orden al mismo consumidor. Sin clave, se reparten y el orden se pierde. (El hashing para repartir claves entre particiones es el mismo problema del particionado de datos, lección 04-01.)
  • Offset: posición de un mensaje dentro de su partición. Es un entero creciente; los consumidores lo usan como marcador.
  • Grupo de consumidores: instancias que comparten un group.id se reparten las particiones del tópico (punto a punto dentro del grupo); grupos distintos leen el tópico completo cada uno (pub/sub entre grupos). El commit del offset es el equivalente del ack: "el grupo inventario ha procesado hasta el offset 4 de la partición 0".
  • Retención: los mensajes se borran por edad o por tamaño, no por consumo. Un consumidor nuevo puede leer el histórico completo (auto.offset.reset=earliest); uno que estuvo caído tres días recupera lo perdido; y analitica puede releer todo el mes si cambia su lógica. Es una capacidad que RabbitMQ no tiene.
  • Réplicas: cada partición se replica en varios brokers (factor de replicación 3 en producción); uno es el líder y atiende lecturas y escrituras. La replicación y sus garantías se estudian en 03-04.

  1. pedido.creado con Kafka

Usaremos confluent-kafka (pip install confluent-kafka), el cliente Python más eficiente (envuelve la biblioteca C librdkafka); kafka-python es una alternativa en Python puro con una API parecida. Primero creamos el tópico con 6 particiones (con un solo broker de desarrollo, factor de replicación 1):

docker compose exec kafka /opt/kafka/bin/kafka-topics.sh --bootstrap-server localhost:9092 \
    --create --topic pedidos.eventos --partitions 6 --replication-factor 1

El productor:

# km0/servicios/pedidos/publicador_kafka.py
import json
import time
import uuid

from confluent_kafka import Producer

TOPICO = "pedidos.eventos"

productor = Producer({
    "bootstrap.servers": "localhost:9092",
    "acks": "all",             # el líder y las réplicas sincronizadas confirman antes de responder
    "client.id": "pedidos",
})


def al_confirmar(err, msg):
    """Callback asíncrono: Kafka acumula mensajes en lotes y confirma después."""
    if err is not None:
        print(f"[pedidos] ERROR al publicar: {err}")
    else:
        print(f"[pedidos] confirmado en {msg.topic()}[{msg.partition()}] offset {msg.offset()}")


def publicar_pedido_creado(pedido):
    evento = {"id_evento": str(uuid.uuid4()), "tipo": "pedido.creado", "version": 1,
              "fecha_ms": int(time.time() * 1000), "origen": "pedidos", "datos": pedido}
    productor.produce(
        TOPICO,
        key=pedido["id"],                        # misma clave -> misma partición -> orden
        value=json.dumps(evento).encode("utf-8"),
        headers=[("tipo", "pedido.creado"), ("id_evento", evento["id_evento"])],
        callback=al_confirmar,
    )
    productor.poll(0)          # atiende callbacks pendientes sin bloquear


if __name__ == "__main__":
    for i, (cliente, producto, cantidad) in enumerate(
            [("ana", "queso-curado", 2), ("marc", "vino-crianza", 6), ("lucia", "tomate-rosa", 3)]):
        publicar_pedido_creado({"id": f"P-2026-00012{3 + i}", "cliente": cliente, "mercado": "girona",
                                "lineas": [{"producto": producto, "cantidad": cantidad}]})
    productor.flush()          # espera a que todo esté confirmado antes de salir

produce no envía: encola el mensaje en un buffer interno que se envía en lotes por un hilo en segundo plano (por eso Kafka consigue cientos de miles de mensajes por segundo). flush() bloquea hasta que todos están confirmados; acks=all hace que "confirmado" signifique escrito en el líder y en las réplicas sincronizadas. Sin flush al terminar el programa, los mensajes del buffer se perderían.

El consumidor de inventario:

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

from confluent_kafka import Consumer, KafkaError

consumidor = Consumer({
    "bootstrap.servers": "localhost:9092",
    "group.id": "inventario",             # todas las instancias de inventario comparten grupo
    "auto.offset.reset": "earliest",      # un grupo nuevo empieza por el principio del log
    "enable.auto.commit": False,          # confirmaremos nosotros, después de procesar
})
consumidor.subscribe(["pedidos.eventos"])

stock = {"tomate-rosa": 120, "calabacin": 80, "queso-curado": 5, "queso-fresco": 30, "vino-crianza": 200}

try:
    while True:
        msg = consumidor.poll(timeout=1.0)        # pull: el consumidor pide, el broker no empuja
        if msg is None:
            continue
        if msg.error():
            if msg.error().code() != KafkaError._PARTITION_EOF:
                print("[inventario] error:", msg.error())
            continue
        evento = json.loads(msg.value())
        if evento["tipo"] == "pedido.creado":
            for linea in evento["datos"]["lineas"]:
                stock[linea["producto"]] -= linea["cantidad"]
            print(f"[inventario] partición {msg.partition()} offset {msg.offset()} "
                  f"clave {msg.key().decode()}: queso-curado={stock['queso-curado']}")
        consumidor.commit(message=msg)            # equivalente al ack: "hasta aquí procesado"
finally:
    consumidor.close()

Para analitica, el mismo código con "group.id": "analitica" y su propia lógica: como es otro grupo, recibe todos los mensajes de nuevo, con independencia de lo que haya confirmado inventario. Prueba a arrancar dos instancias del consumidor de inventario: verás en los logs cómo Kafka reasigna las 6 particiones (3 y 3) entre ellas, y cómo cada pedido va siempre a la misma instancia por su clave. Para a una y verás las particiones volver a la otra: es el rebalanceo de grupos, que da tolerancia a fallos a los consumidores sin que el productor se entere.

Un detalle sobre el commit: commit(message=msg) marca hasta ese offset en esa partición, no "este mensaje". Si se procesan los mensajes 3, 4 y 5 y se confirma el 5, al reiniciar se sigue desde el 6. Si se confirma el 5 sin haber procesado el 4 (por ejemplo, procesando en paralelo), el 4 se pierde. El commit es un marcador de posición, y las consecuencias de confirmar antes o después de procesar (perder o duplicar) son exactamente las garantías de entrega que analiza 02-05.

  1. RabbitMQ frente a Kafka: criterios de elección

Criterio RabbitMQ Apache Kafka
Modelo Cola inteligente: el broker enruta y hace seguimiento de cada mensaje Log tonto: el broker guarda; el consumidor recuerda por dónde va
El mensaje tras consumirse Se borra Permanece hasta la retención
Enrutamiento Muy flexible (exchanges, comodines, cabeceras) Por tópico y partición; el filtrado lo hace el consumidor
Orden Por cola, con un consumidor Por partición (por clave)
Entrega Push al consumidor Pull del consumidor
Rendimiento típico Decenas de miles de mensajes/s Cientos de miles a millones de mensajes/s
Releer el histórico No Sí (por diseño)
Consumidores lentos Acumulan mensajes en la cola; puede degradar al broker No afectan al broker ni a otros grupos
Prioridades, TTL, colas de muertos Nativas No nativas (se construyen con tópicos, 02-05)
Petición/respuesta Cómodo (reply_to, 02-05) Incómodo
Complejidad operativa Baja-media Media-alta (particiones, rebalanceos, retención), aunque KRaft la ha reducido
Protocolo AMQP (también MQTT y STOMP con plugins) Propio, binario sobre TCP
Ecosistema de datos Poco Enorme: Kafka Connect, Streams, integración con Spark y Flink (Módulo 5)
Papel en Kilómetro Cero Colas de trabajo (facturas, correos), petición/respuesta interna Eventos de dominio (pedidos.eventos, inventario.eventos, telemetría de reparto)

Criterios de decisión:

  1. ¿Los mensajes son eventos que varios sistemas querrán, quizá en el futuro, quizá releídos? Kafka. Esta es la razón por la que la arquitectura objetivo de 01-06 elige Kafka para los eventos de dominio: analitica querrá reprocesar el mes; un servicio nuevo querrá el histórico.
  2. ¿Son tareas que hay que ejecutar una vez y olvidar, con enrutamiento fino, prioridades o respuesta? RabbitMQ.
  3. ¿Importa el orden por entidad (todos los eventos de un pedido en orden)? Kafka con clave de partición lo da de forma natural.
  4. ¿Volumen? Por debajo de unos pocos miles de mensajes por segundo, cualquiera de los dos; muy por encima, Kafka.
  5. ¿Quién lo va a operar? Un broker mal operado es peor que ninguno. Si el equipo es pequeño, empieza con uno solo y con un servicio gestionado.

Kilómetro Cero, siguiendo la arquitectura objetivo, usará Kafka para los eventos de dominio y se reserva RabbitMQ para colas de trabajo y para el patrón petición/respuesta que veremos en 02-05. Tener ambos no es obligatorio; muchos sistemas viven bien con uno solo.

  1. Ampliar el docker-compose.yml de km0/

Añadimos los dos brokers al entorno de 01-06. Kafka en modo KRaft (sin ZooKeeper, el modo estándar desde la versión 3.x), con dos listeners: uno interno para los servicios de la red de Compose y otro para los scripts que ejecutamos desde el host.

  rabbitmq:
    image: rabbitmq:3.13-management
    environment:
      RABBITMQ_DEFAULT_USER: km0
      RABBITMQ_DEFAULT_PASS: km0_dev
    ports:
      - "5672:5672"      # AMQP
      - "15672:15672"    # interfaz web de administración
    volumes:
      - rabbitmq_data:/var/lib/rabbitmq
    healthcheck:
      test: ["CMD", "rabbitmq-diagnostics", "-q", "ping"]
      interval: 10s
      timeout: 5s
      retries: 5

  kafka:
    image: apache/kafka:3.8.0
    environment:
      KAFKA_NODE_ID: 1
      KAFKA_PROCESS_ROLES: broker,controller
      KAFKA_LISTENERS: INTERNO://:19092,EXTERNO://:9092,CONTROLLER://:9093
      KAFKA_ADVERTISED_LISTENERS: INTERNO://kafka:19092,EXTERNO://localhost:9092
      KAFKA_LISTENER_SECURITY_PROTOCOL_MAP: CONTROLLER:PLAINTEXT,INTERNO:PLAINTEXT,EXTERNO:PLAINTEXT
      KAFKA_INTER_BROKER_LISTENER_NAME: INTERNO
      KAFKA_CONTROLLER_LISTENER_NAMES: CONTROLLER
      KAFKA_CONTROLLER_QUORUM_VOTERS: 1@kafka:9093
      KAFKA_OFFSETS_TOPIC_REPLICATION_FACTOR: 1          # un solo broker en desarrollo
      KAFKA_TRANSACTION_STATE_LOG_REPLICATION_FACTOR: 1
      KAFKA_TRANSACTION_STATE_LOG_MIN_ISR: 1
      KAFKA_LOG_RETENTION_HOURS: 168                     # 7 días
    ports:
      - "9092:9092"
    volumes:
      - kafka_data:/var/lib/kafka/data

volumes:
  postgres_data:
  rabbitmq_data:
  kafka_data:

Los servicios que corran dentro de Compose (como inventario a partir de 02-03) usarán kafka:19092 y rabbitmq:5672; los scripts lanzados desde el host, localhost:9092 y localhost:5672. Las contraseñas en claro son solo para desarrollo: la gestión de secretos es tema de 06-04. Con docker compose up -d rabbitmq kafka y los tres scripts de este capítulo, tienes el primer flujo asíncrono de Kilómetro Cero funcionando.

Errores Comunes y Consejos

  • Usar mensajería para lo que necesita respuesta inmediata. "¿Hay stock?" no puede esperar a que un consumidor procese un evento. Síncrono para decisiones, asíncrono para consecuencias.
  • Colas no durables o mensajes no persistentes en flujos críticos. Un reinicio del broker se lleva los pedidos. Durable + persistente, siempre, salvo para telemetría desechable.
  • Confirmar (ack/commit) antes de procesar. Si el consumidor muere después de confirmar y antes de terminar, el mensaje se pierde. Confirmar después de procesar duplica en el caso opuesto, y 02-05 explica cómo convivir con ello.
  • Olvidar flush() en el productor de Kafka. produce solo encola; un proceso que termina sin flush pierde el buffer, sin error.
  • Publicar sin clave de partición y esperar orden. Sin clave, los eventos de un mismo pedido se reparten entre particiones y pedido.pagado puede procesarse antes que pedido.creado.
  • Un tópico por cada tipo de evento en Kafka. Rompe el orden entre eventos de la misma entidad (están en tópicos distintos) y multiplica las particiones. Un tópico por entidad (pedidos.eventos) con el tipo en una cabecera suele ser mejor.
  • Consumidores que declaran colas que "deberían existir". Si inventario nunca arrancó, su cola de RabbitMQ no existe y los eventos publicados entretanto se pierden. Declara la topología (exchanges, colas, bindings, tópicos) en el despliegue.
  • Mensajes enormes. Ni RabbitMQ ni Kafka están hechos para transportar fotos de productos de 5 MB. En el mensaje va una referencia al almacén de objetos (04-03).
  • Consejo: pon siempre una envoltura estándar en los eventos (id_evento, tipo, version, fecha_ms, origen, datos). Cuesta cinco líneas y lo agradecerás en cada patrón de 02-05 y en cada investigación de 07-02.
  • Consejo: mira la interfaz de RabbitMQ y los lag de los grupos de Kafka (kafka-consumer-groups.sh --describe) desde el primer día. Una cola que crece o un lag que no baja son el primer síntoma de un consumidor enfermo, mucho antes de que nadie se queje.

Ejercicios

Ejercicio 1: Enrutamiento con topic exchange

reparto quiere recibir solo los pedidos ya pagados (pedido.pagado) y atencion-cliente quiere todos los eventos de cancelación de cualquier entidad (pedido.cancelado, reparto.cancelado, ...). Escribe las declaraciones de cola y binding de ambos sobre el exchange km0.pedidos (asume que los eventos de reparto también se publican en él con routing keys reparto.<tipo>), e indica qué recibiría cada uno si se publicaran, en orden: pedido.creado, pedido.pagado, reparto.asignado, reparto.cancelado, pedido.cancelado.

Ejercicio 2: Particiones, claves y orden

El tópico pedidos.eventos tiene 6 particiones y inventario tiene 3 instancias en el mismo grupo. (a) ¿Cuántas particiones lee cada instancia? (b) Si se añaden 4 instancias más (7 en total), ¿qué ocurre? (c) pedidos publica pedido.creado y pedido.cancelado para P-2026-000123 con clave P-2026-000123, y pedido.creado para P-2026-000124. ¿Está garantizado que inventario procesa el creado del 123 antes que su cancelación? ¿Y antes que el creado del 124? (d) Un desarrollador propone usar cliente como clave en lugar de id del pedido, "para que todos los pedidos de Ana vayan al mismo consumidor". ¿Qué gana y qué arriesga?

Ejercicio 3: Elegir broker

Para cada flujo de Kilómetro Cero, indica RabbitMQ o Kafka y justifica con la tabla del apartado 8:

  1. Generar el PDF de la factura de cada pedido pagado (tarea pesada, una vez, sin orden).
  2. Los eventos de stock (inventario.eventos) que catalogo consume para el indicador "quedan pocas unidades" y que analitica quiere reprocesar cada mes para estudiar roturas de stock.
  3. Posiciones de 400 repartidores, una por segundo cada uno, para el mapa en vivo y para el análisis posterior de rutas.
  4. Un servicio interno que necesita preguntar a pagos si una tarjeta está bloqueada, y esperar la respuesta.

Soluciones

Solución 1:

# reparto: solo pedidos pagados
canal.queue_declare(queue="reparto.pedidos", durable=True)
canal.queue_bind(queue="reparto.pedidos", exchange="km0.pedidos", routing_key="pedido.pagado")

# atención al cliente: cualquier cancelación de cualquier entidad
canal.queue_declare(queue="atencion.cancelaciones", durable=True)
canal.queue_bind(queue="atencion.cancelaciones", exchange="km0.pedidos", routing_key="*.cancelado")

Con los cinco eventos publicados: reparto.pedidos recibe solo pedido.pagado. atencion.cancelaciones recibe reparto.cancelado y pedido.cancelado (el comodín * casa exactamente una palabra: pedido o reparto). Ninguna de las dos recibe pedido.creado ni reparto.asignado. Y las colas ya existentes siguen recibiendo lo suyo: inventario.pedidos solo pedido.creado; analitica.pedidos (pedido.#) los tres eventos de pedido, pero no los de reparto. Cada consumidor decide qué quiere; el productor no cambia.

Solución 2:

(a) Kafka reparte las 6 particiones entre las 3 instancias: 2 cada una. (b) Con 7 instancias y 6 particiones, 6 instancias leen una partición cada una y la séptima queda ociosa: el número de particiones es el límite superior del paralelismo de un grupo. Para aprovechar más instancias habría que crear el tópico con más particiones (y no se pueden añadir sin cambiar la asignación clave→partición de los mensajes futuros, lo que rompe temporalmente el orden por clave). (c) Sí para el 123: ambos mensajes tienen la misma clave, van a la misma partición, y una partición la lee una sola instancia en orden. No respecto al 124: puede estar en otra partición leída por otra instancia; su orden relativo no está definido, y no importa, porque son pedidos independientes. (d) Gana que todos los eventos de un mismo cliente se procesan en orden y en la misma instancia (útil si hubiera reglas por cliente, como un límite de pedidos por hora). Arriesga dos cosas: particiones desequilibradas (un cliente empresarial que hace el 30 % de los pedidos concentra el 30 % de la carga en una partición y una instancia) y, en cierto sentido, orden innecesario: los pedidos de Ana son independientes entre sí, así que serializarlos no aporta nada y reduce el paralelismo. La clave debe ser la entidad cuyo orden importa, ni más fina ni más gruesa.

Solución 3:

  1. RabbitMQ: cola de trabajo clásica; cada factura la genera una instancia y desaparece; sin orden ni relectura; con prioridades si hiciera falta (facturas de campaña primero). Kafka también podría, pero no aporta nada aquí.
  2. Kafka: eventos de dominio con dos consumidores con necesidades distintas (uno en tiempo real, otro releyendo un mes), orden por producto (clave = slug del producto) y retención larga. Es exactamente el caso de uso para el que se diseñó.
  3. Kafka (con un puente MQTT delante para el tramo móvil, 08-02): volumen alto (400 mensajes/s sostenidos, mucho más en campaña), dos consumidores (mapa en vivo y análisis de rutas), orden por repartidor (clave = furgoneta-3), y retención para el análisis posterior, que es un trabajo por lotes del Módulo 5. RabbitMQ acumularía las posiciones y no permitiría releerlas.
  4. RabbitMQ, si se decide hacerlo por mensajería: el patrón petición/respuesta con reply_to y correlation_id (02-05) es natural en AMQP e incómodo en Kafka. Pero la pregunta previa es si debe ser mensajería en absoluto: es una consulta síncrona que necesita respuesta inmediata, y una llamada gRPC a pagos (02-03) es más simple y más rápida. La mensajería solo se justificaría si pagos fuera lento o intermitente y se quisiera absorber picos, o si la respuesta pudiera tardar y el solicitante pudiera esperar sin bloquear.

Conclusión

Esta lección ha introducido la segunda mitad de la comunicación entre servicios. Una llamada síncrona acopla en el tiempo, en el espacio y en el ritmo; un broker rompe los tres acoplamientos a cambio de que el emisor deje de saber cuándo, y si, se procesará su mensaje. Con ese vocabulario (productor, consumidor, cola, tópico, ack, persistencia, durabilidad, retención) hemos distinguido el modelo punto a punto, para repartir trabajo, del de publicación/suscripción, para difundir eventos, y hemos visto cómo cada broker los combina. RabbitMQ es una cola inteligente con enrutamiento flexible (exchanges, bindings, routing keys con comodines) que borra los mensajes al confirmarse; Kafka es un log distribuido en el que los mensajes permanecen, los consumidores recuerdan su offset, las particiones dan paralelismo y las claves dan orden por entidad, y cualquier grupo puede releer el histórico. Kilómetro Cero tiene ahora su primer flujo asíncrono: pedidos publica pedido.creado con una envoltura estándar, inventario y analitica lo consumen a su ritmo, y ambos brokers están en el docker-compose.yml.

Pero hemos dejado varias preguntas abiertas a propósito. ¿Qué pasa si inventario muere después de descontar el stock y antes de confirmar el mensaje? El broker lo reenviará y el stock se descontará dos veces, el mismo problema que dejamos pendiente en 01-04 y en 02-02. ¿Qué pasa si pedidos guarda el pedido en PostgreSQL y se cae antes de publicar el evento, o publica el evento y luego falla el COMMIT? ¿Y con un mensaje malformado que hace fallar al consumidor una y otra vez, bloqueando a todos los que vienen detrás? Estas preguntas son las garantías de entrega y los patrones que las hacen manejables (consumidores idempotentes, outbox transaccional, colas de mensajes muertos), y son el tema de la última lección del módulo: Patrones de Comunicación Asíncrona.

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