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
- Por qué desacoplar en el tiempo
- Vocabulario de la mensajería
- Dos modelos: punto a punto y publicación/suscripción
- RabbitMQ y AMQP: exchanges, colas, bindings y routing keys
pedido.creadocon RabbitMQ ypika- Apache Kafka: el log distribuido
pedido.creadocon Kafka- RabbitMQ frente a Kafka: criterios de elección
- Ampliar el
docker-compose.ymldekm0/ - Errores comunes y consejos
- Ejercicios
- Conclusión
- 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
analiticaestá reiniciándose cuandopedidosintenta notificarle una venta, la notificación falla opedidosespera. - En el espacio: el emisor necesita saber quién es el receptor y dónde está (dirección, puerto). Si mañana
repartotambién quiere saber de los pedidos, hay que cambiarpedidos. - En el ritmo: el emisor no puede ir más rápido que el receptor más lento. En la "Semana de la Vendimia",
pedidosproduce 1.200 pedidos por segundo; sianaliticasolo procesa 300,pedidosse 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.
- 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 |
- 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.
- 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.
pedido.creado con RabbitMQ y pika
pedido.creado con RabbitMQ y pikapika 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.
inventariosolo quierepedido.creado;analiticaquierepedido.#(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:
inventarioyanaliticareciben cada uno su ejemplar (pub/sub). Si arrancas dos instancias deconsumidor_rabbit.pydeinventario, compartirán la colainventario.pedidosy 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=2son las dos mitades de la supervivencia a un reinicio del broker: la cola y el mensaje. Una sin la otra no sirve.prefetch_count=1eninventario: 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.
- 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 Npara 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.pagadoypedido.canceladodeP-2026-000123llegan 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.idse 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 grupoinventarioha 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; yanaliticapuede 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.
pedido.creado con Kafka
pedido.creado con KafkaUsaremos 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 1El 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 salirproduce 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.
- 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:
- ¿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:
analiticaquerrá reprocesar el mes; un servicio nuevo querrá el histórico. - ¿Son tareas que hay que ejecutar una vez y olvidar, con enrutamiento fino, prioridades o respuesta? RabbitMQ.
- ¿Importa el orden por entidad (todos los eventos de un pedido en orden)? Kafka con clave de partición lo da de forma natural.
- ¿Volumen? Por debajo de unos pocos miles de mensajes por segundo, cualquiera de los dos; muy por encima, Kafka.
- ¿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.
- Ampliar el
docker-compose.yml de km0/
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.producesolo encola; un proceso que termina sinflushpierde 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.pagadopuede procesarse antes quepedido.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
inventarionunca 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:
- Generar el PDF de la factura de cada pedido pagado (tarea pesada, una vez, sin orden).
- Los eventos de stock (
inventario.eventos) quecatalogoconsume para el indicador "quedan pocas unidades" y queanaliticaquiere reprocesar cada mes para estudiar roturas de stock. - Posiciones de 400 repartidores, una por segundo cada uno, para el mapa en vivo y para el análisis posterior de rutas.
- Un servicio interno que necesita preguntar a
pagossi 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:
- 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í.
- 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ñó.
- 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. - RabbitMQ, si se decide hacerlo por mensajería: el patrón petición/respuesta con
reply_toycorrelation_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 apagos(02-03) es más simple y más rápida. La mensajería solo se justificaría sipagosfuera 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
- Conceptos Básicos de Sistemas Distribuidos
- Modelos de Sistemas Distribuidos
- Ventajas y Desafíos de los Sistemas Distribuidos
- Las Falacias de la Computación Distribuida
- Tiempo, Relojes y Ordenación de Eventos
- Del Monolito a la Plataforma Distribuida: el Caso Kilómetro Cero
Módulo 2: Comunicación en Sistemas Distribuidos
- Protocolos de Comunicación
- RPC y RMI
- gRPC y Serialización de Datos
- Mensajería y Colas de Mensajes
- Patrones de Comunicación Asíncrona
Módulo 3: Consistencia y Replicación
- Modelos de Consistencia
- El Teorema CAP y PACELC
- Algoritmos de Consenso
- Replicación de Datos
- Transacciones Distribuidas y Sagas
Módulo 4: Almacenamiento Distribuido
- Particionado de Datos y Hashing Consistente
- Sistemas de Archivos Distribuidos
- Almacenamiento de Objetos
- Bases de Datos Distribuidas
- Cachés Distribuidos
Módulo 5: Computación Distribuida
- Modelos de Computación Distribuida
- MapReduce y Hadoop
- Spark y Computación en Memoria
- Procesamiento de Flujos de Datos
- Planificación de Trabajos y Pipelines de Datos
Módulo 6: Seguridad en Sistemas Distribuidos
- Autenticación y Autorización
- Cifrado y Protección de Datos
- Gestión de Identidades
- Seguridad entre Servicios: mTLS y Gestión de Secretos
- Puertas de Enlace, Limitación de Tasa y Auditoría
Módulo 7: Monitoreo y Mantenimiento
- Monitoreo de Sistemas Distribuidos
- Logs Centralizados y Trazabilidad Distribuida
- Gestión de Fallos y Recuperación
- Patrones de Resiliencia: Timeouts, Reintentos y Circuit Breaker
- Automatización y Orquestación
- Pruebas en Sistemas Distribuidos e Ingeniería del Caos
