Eventos y mensajería
Arquitectura orientada a eventos: comunicación síncrona y asíncrona, colas y publicar/suscribir, productores y consumidores, garantías de entrega, idempotencia, orden y reintentos, colas de mensajes fallidos, patrón outbox, y RabbitMQ frente a Kafka.
Definición
Un evento (event) es el registro inmutable de un hecho que ya ocurrió, nombrado en pasado: «pedido pagado», «cliente registrado», «pieza enviada». No pide nada a nadie; solo informa. Un mensaje (message) es el paquete de datos que viaja de un componente a otro: lleva una carga útil (payload, normalmente JSON) y metadatos como un identificador único y una marca de tiempo. Un evento se transmite como un mensaje.
Un sistema de mensajería es la infraestructura que recibe esos mensajes, los conserva y los entrega. Su pieza central es el broker (intermediario de mensajes): un servicio independiente de las aplicaciones que sirve de punto de encuentro. Los actores son:
- Productor (producer o publisher): el componente que publica el mensaje.
- Consumidor (consumer o subscriber): el componente que lo recibe y actúa.
- Broker: guarda los mensajes y decide a quién se entregan.
Una arquitectura orientada a eventos (event-driven architecture, EDA) es un estilo en el que los componentes se comunican publicando y reaccionando a eventos, en lugar de invocarse directamente. El productor no sabe quién escucha ni cuántos son. Martin Fowler advierte que «orientado a eventos» agrupa patrones distintos (notificación de eventos, transferencia de estado, event sourcing, CQRS), por lo que conviene precisar cuál se usa.
Comunicación síncrona frente a asíncrona
En una comunicación síncrona (una llamada HTTP, por ejemplo) quien llama se queda esperando la respuesta. En una asíncrona deja el mensaje y continúa; la respuesta, si existe, llega después por otro camino.
La diferencia de fondo es el acoplamiento temporal (temporal coupling): en el modo síncrono, emisor y receptor deben estar disponibles al mismo tiempo. Si la API de facturación está caída o lenta, el pago también falla o se retrasa.
| Aspecto | Síncrona (petición/respuesta) | Asíncrona (mensajes) |
|---|---|---|
| Disponibilidad | Ambos extremos deben estar activos | El receptor puede estar caído; el mensaje espera |
| Latencia para el usuario | Suma de todas las llamadas encadenadas | Solo la del primer paso; el resto ocurre después |
| Acoplamiento | El emisor conoce a cada receptor | El emisor solo conoce el tema o la cola |
| Picos de carga | Se propagan al receptor | La cola los amortigua |
| Resultado | Inmediato y directo | Eventual; requiere seguimiento |
| Depuración | Una traza lineal | Flujos distribuidos y repartidos en el tiempo |
Ninguna es superior: la consulta «¿cuánto cuesta este producto?» necesita respuesta inmediata y es síncrona por naturaleza; «enviar el correo de confirmación» no necesita que el cliente espere.
Modelos: cola y publicar/suscribir
Los brokers ofrecen dos modelos básicos.
Cola punto a punto (queue): los mensajes se acumulan en una cola y cada uno lo recibe un solo consumidor, aunque varios compitan por leer de ella (competing consumers). Sirve para repartir trabajo: cien miniaturas por generar entre cuatro procesos.
Publicar/suscribir (publish/subscribe, pub/sub): el productor publica en un tema (topic) y cada suscriptor interesado recibe su propia copia. Sirve para difundir un hecho a varias partes: el tema pedido_pagado llega a correo, facturación y envío. En la práctica se combinan: cada suscriptor tiene su propia cola (los avisos esperan ahí si el consumidor no está disponible) y, dentro de un suscriptor, pueden competir varias instancias para repartir la carga.
┌── cola "correo" ──► consumidor correo
tienda ──► tema pedido_pagado ┼── cola "factura" ──► consumidor factura
└── cola "envio" ──► consumidor envío
Tipos de mensaje
No todo lo que viaja por un broker es un evento. La distinción importa porque cambia quién tiene la responsabilidad:
| Tipo | Significado | Ejemplo | Destinatarios |
|---|---|---|---|
| Evento (event) | Algo ocurrió; en pasado | PedidoPagado | Quien le interese; el emisor no espera nada |
| Comando (command) | Se ordena hacer algo; imperativo | EnviarCorreo | Un destinatario concreto, que debe ejecutarlo |
| Notificación de evento (event notification) | Evento mínimo: identificador y poco más | {pedido_id: "P-100"} | El consumidor consulta el detalle al emisor si lo necesita |
| Evento con estado transferido (event-carried state transfer) | Evento que lleva los datos completos | {pedido_id, cliente, total, partidas} | El consumidor ya no necesita llamar al emisor |
La notificación mínima mantiene los mensajes pequeños, pero reintroduce una llamada síncrona de vuelta. El evento con estado evita esa dependencia a costa de mensajes mayores y del riesgo de datos desactualizados.
Garantías de entrega
Una red puede perder, duplicar o retrasar mensajes, y cualquier proceso puede caer a mitad del trabajo. Los brokers declaran qué prometen:
- A lo sumo una vez (at-most-once): se entrega 0 o 1 veces. Puede perderse, nunca se duplica. Apto para métricas o lecturas de sensores donde perder un dato es tolerable.
- Al menos una vez (at-least-once): se entrega 1 o más veces. No se pierde, pero puede duplicarse. Funciona con confirmación (acknowledgement, ack): el broker retiene el mensaje hasta que el consumidor confirma que terminó; si el ack no llega, lo reenvía. Es el modo habitual.
- Exactamente una vez (exactly-once): cada mensaje produce su efecto una sola vez. Como garantía de red es inalcanzable en el caso general (un ack puede perderse sin que el emisor sepa si el trabajo se hizo). Lo que se logra es un exactly-once efectivo: entrega al menos una vez más un consumidor idempotente, o transacciones dentro de un mismo sistema (Kafka las ofrece entre sus propios temas).
La consecuencia práctica: diseña asumiendo duplicados.
Idempotencia y deduplicación
Una operación es idempotente cuando aplicarla varias veces produce el mismo resultado que aplicarla una. «Establecer el estado a pagado» es idempotente; «sumar 100 al saldo» no lo es.
Para volver idempotente a un consumidor se usa una clave de idempotencia: el identificador único del mensaje (o del hecho de negocio, como el id del pedido) y una verificación de que no se procesó antes. Dos técnicas:
- Restricción de unicidad en la base de datos: la factura tiene
pedido_idcomo clave primaria; un segundo intento no inserta nada. Es la más robusta porque la verificación y el efecto son atómicos. - Tabla de mensajes procesados (inbox): se registra el id al terminar y se descarta todo mensaje ya visto. Debe guardarse en la misma transacción que el efecto.
Evita el esquema «consultar si existe y luego insertar» sin restricción: dos entregas simultáneas pasan ambas la consulta.
Orden, particiones y reintentos
Orden. Una cola única con un solo consumidor conserva el orden de llegada. Al añadir consumidores en paralelo, o al reintentar un mensaje fallido mientras otros avanzan, el orden global se pierde. Las plataformas distribuidas dividen un tema en particiones (partitions): el orden solo se garantiza dentro de una partición. Se asigna una clave de partición (por ejemplo el id del pedido) para que todos los eventos de una misma entidad caigan en la misma partición y se lean en secuencia.
Reintentos con retroceso. Un fallo transitorio (un servicio saturado) suele resolverse esperando. Reintentar de inmediato y sin límite agrava el problema, por lo que se usa retroceso exponencial (exponential backoff): esperas crecientes (1 s, 2 s, 4 s…), normalmente con algo de aleatoriedad (jitter) para que muchos consumidores no reintenten a la vez.
Cola de mensajes fallidos (dead-letter queue, DLQ). Un mensaje venenoso (poison message) es el que siempre falla, por ejemplo por una dirección inválida. Tras un número máximo de intentos se aparta a una cola aparte, para que no bloquee a los demás ni se reintente eternamente. La DLQ no es un basurero: requiere alertas, inspección y un procedimiento para corregir y reinyectar los mensajes.
Patrón outbox transaccional
Un servicio que guarda un pedido en su base de datos y publica el evento en un broker enfrenta un problema de escritura dual (dual write): son dos sistemas sin transacción común. Si guarda y se cae antes de publicar, el pedido existe pero nadie se entera. Si publica y la base de datos falla, se anuncia un pedido que no existe.
El patrón outbox transaccional (transactional outbox, «bandeja de salida») lo resuelve en dos pasos:
- En la misma transacción local que escribe el pedido, se inserta el evento en una tabla
outboxde la misma base. Por la atomicidad (las transacciones y las propiedades ACID que vimos en el día 3), ocurren ambos o ninguno. - Un proceso aparte, el relé (message relay), lee los eventos pendientes de la tabla, los publica en el broker y los marca como publicados. Puede ser un sondeo periódico o la captura de cambios del registro de la base (change data capture).
Si el relé cae entre publicar y marcar, al reiniciar publicará de nuevo: el resultado es entrega al menos una vez, y de ahí que los consumidores deban ser idempotentes. Ambas ideas se complementan.
RabbitMQ frente a Kafka
Los dos dominan el panorama, con modelos distintos. RabbitMQ es un broker de mensajes clásico (protocolo AMQP): enruta mensajes mediante exchanges hacia colas y los elimina cuando el consumidor confirma. Apache Kafka es un registro distribuido de eventos (distributed log): los eventos se añaden a particiones y se conservan durante un período de retención; cada consumidor lleva su propia posición (offset) y puede releer el pasado.
| Aspecto | RabbitMQ | Apache Kafka |
|---|---|---|
| Modelo | Broker con colas y enrutamiento flexible | Registro particionado y persistente |
| Tras consumirse | El mensaje se elimina | Permanece hasta vencer la retención |
| Releer historial | No es el uso habitual | Sí, moviendo el offset |
| Orden | Por cola | Por partición |
| Rendimiento | Alto; destaca en enrutamiento y baja latencia | Muy alto; pensado para grandes volúmenes |
| Uso típico | Colas de trabajo, tareas, integración entre servicios | Flujos de datos, analítica, event sourcing |
| Operación | Más sencillo de iniciar | Más complejo de operar |
Los servicios gestionados en la nube evitan operar el broker: Amazon SQS (colas) junto con SNS (publicar/suscribir) para difundir a una cola por suscriptor; Google Cloud Pub/Sub (temas y suscripciones); y variantes administradas de Kafka. Cambias control por menos mantenimiento.
Cuándo usarlo y cuándo no
Conviene cuando:
- Hay tareas lentas que el usuario no debe esperar (correos, facturas, reportes).
- Varios sistemas deben enterarse del mismo suceso y se quiere añadir nuevos sin tocar al emisor.
- Los picos de carga necesitan amortiguarse con una cola.
- Un componente puede fallar sin detener al resto.
No conviene cuando:
- La respuesta se necesita de inmediato para continuar (una consulta de precio).
- El sistema es pequeño y una llamada directa basta: un broker es infraestructura que alguien debe operar y vigilar.
- Se requiere consistencia inmediata. La mensajería produce consistencia eventual (eventual consistency): durante un intervalo la factura aún no existe aunque el pago ya esté registrado, y la interfaz debe contemplarlo.
Los costos reales: depuración más difícil (un flujo atraviesa varios servicios; hacen falta identificadores de correlación y trazas distribuidas), manejo explícito de duplicados, orden y mensajes fallidos, y una operación más compleja. Adopta el patrón por una necesidad concreta, no por moda.
Práctica: un broker mínimo, tres consumidores y un outbox
Todo con la librería estándar de Python 3.10 o superior (más pytest para las pruebas). El broker vive en memoria; la factura y el outbox usan SQLite (sqlite3). Crea mensajeria.py:
import json
import sqlite3
from collections import deque
class Broker:
"""Broker mínimo en memoria: publicar/suscribir con cola por suscriptor."""
def __init__(self, max_intentos=3, base_espera=1.0):
self.max_intentos = max_intentos
self.base_espera = base_espera
self.suscriptores = {} # topic -> {nombre: manejador}
self.colas = {} # (topic, nombre) -> deque de mensajes
self.fallidos = [] # cola de mensajes fallidos (DLQ)
self.esperas = [] # esperas de backoff calculadas (no se duerme)
def suscribir(self, topic, nombre, manejador):
self.suscriptores.setdefault(topic, {})[nombre] = manejador
self.colas[(topic, nombre)] = deque()
def publicar(self, topic, mensaje):
for nombre in self.suscriptores.get(topic, {}):
self.colas[(topic, nombre)].append(
{"mensaje": mensaje, "intentos": 0})
def procesar(self):
"""Entrega hasta vaciar las colas. Sin ack no hay baja: se reintenta."""
hay_trabajo = True
while hay_trabajo:
hay_trabajo = False
for (topic, nombre), cola in self.colas.items():
if not cola:
continue
hay_trabajo = True
entrega = cola.popleft()
try:
self.suscriptores[topic][nombre](entrega["mensaje"])
except Exception as error: # ack negativo
entrega["intentos"] += 1
if entrega["intentos"] >= self.max_intentos:
self.fallidos.append(
(nombre, entrega["mensaje"], str(error)))
else:
self.esperas.append(
self.base_espera * 2 ** (entrega["intentos"] - 1))
cola.append(entrega) # reintento
class Factura:
"""Consumidor idempotente: la clave primaria es el id del pedido."""
def __init__(self, db):
self.db = db
db.execute("CREATE TABLE IF NOT EXISTS facturas ("
"pedido_id TEXT PRIMARY KEY, total REAL NOT NULL)")
def __call__(self, mensaje):
with self.db:
self.db.execute(
"INSERT OR IGNORE INTO facturas VALUES (?, ?)",
(mensaje["pedido_id"], mensaje["total"]))
def cuenta(self):
return self.db.execute("SELECT COUNT(*) FROM facturas").fetchone()[0]
def crear_pedido_pagado(db, pedido_id, total):
"""Pedido y evento en la MISMA transacción (outbox transaccional)."""
with db:
db.execute("CREATE TABLE IF NOT EXISTS pedidos ("
"id TEXT PRIMARY KEY, total REAL NOT NULL)")
db.execute("CREATE TABLE IF NOT EXISTS outbox ("
"id INTEGER PRIMARY KEY AUTOINCREMENT, topic TEXT NOT NULL,"
"cuerpo TEXT NOT NULL, publicado INTEGER NOT NULL DEFAULT 0)")
db.execute("INSERT INTO pedidos VALUES (?, ?)", (pedido_id, total))
db.execute("INSERT INTO outbox (topic, cuerpo) VALUES (?, ?)",
("pedido_pagado",
json.dumps({"pedido_id": pedido_id, "total": total})))
def retransmitir(db, broker, caer_antes_de_marcar=False):
"""Relé: publica lo pendiente y lo marca. Si cae antes de marcar, repite."""
filas = db.execute(
"SELECT id, topic, cuerpo FROM outbox WHERE publicado = 0 ORDER BY id"
).fetchall()
for id_, topic, cuerpo in filas:
broker.publicar(topic, json.loads(cuerpo))
if caer_antes_de_marcar:
raise RuntimeError("el relé se cayó tras publicar")
with db:
db.execute("UPDATE outbox SET publicado = 1 WHERE id = ?", (id_,))
Puntos clave del código:
Broker.publicarcopia el mensaje a la cola de cada suscriptor (pub/sub);procesarlo entrega y, si el manejador lanza una excepción, cuenta como ack negativo y el mensaje vuelve a la cola.- Tras
max_intentosel mensaje pasa afallidos(la DLQ). Las esperas del retroceso exponencial (1 s, 2 s…) se calculan y se registran, pero no se duermen, para que las pruebas sean instantáneas; en un sistema real el consumidor aguardaría ese tiempo antes de reintentar. - Un reintento se coloca al final de su cola: ilustra por qué reintentar altera el orden.
Facturaes idempotente por la clave primariapedido_idconINSERT OR IGNORE.crear_pedido_pagadoescribe el pedido y el evento en una sola transacción;retransmitires el relé, con un parámetro para simular su caída justo antes de marcar el evento.
Ahora las pruebas, en test_mensajeria.py:
import sqlite3
import pytest
from mensajeria import Broker, Factura, crear_pedido_pagado, retransmitir
PEDIDO = {"pedido_id": "P-100", "total": 1500.0}
@pytest.fixture
def db():
return sqlite3.connect(":memory:")
def test_tres_consumidores_reciben_el_evento(db):
broker, vistos = Broker(), []
factura = Factura(db)
broker.suscribir("pedido_pagado", "correo", lambda m: vistos.append(("correo", m["pedido_id"])))
broker.suscribir("pedido_pagado", "factura", factura)
broker.suscribir("pedido_pagado", "envio", lambda m: vistos.append(("envio", m["pedido_id"])))
broker.publicar("pedido_pagado", PEDIDO)
broker.procesar()
assert sorted(vistos) == [("correo", "P-100"), ("envio", "P-100")]
assert factura.cuenta() == 1
def test_al_menos_una_vez_reintenta_hasta_tener_exito():
broker, intentos = Broker(max_intentos=3), []
def envio_inestable(m):
intentos.append(m["pedido_id"])
if len(intentos) < 3:
raise ConnectionError("transportista caído")
broker.suscribir("pedido_pagado", "envio", envio_inestable)
broker.publicar("pedido_pagado", PEDIDO)
broker.procesar()
assert len(intentos) == 3
assert broker.esperas == [1.0, 2.0] # backoff exponencial
assert broker.fallidos == []
def test_mensaje_venenoso_va_a_la_cola_de_fallidos():
broker = Broker(max_intentos=3)
def siempre_falla(m):
raise ValueError("dirección inválida")
broker.suscribir("pedido_pagado", "envio", siempre_falla)
broker.publicar("pedido_pagado", PEDIDO)
broker.procesar()
assert len(broker.fallidos) == 1
assert broker.fallidos[0][0] == "envio"
def test_un_consumidor_roto_no_bloquea_a_los_demas(db):
broker, factura = Broker(), Factura(db)
broker.suscribir("pedido_pagado", "correo", lambda m: 1 / 0)
broker.suscribir("pedido_pagado", "factura", factura)
broker.publicar("pedido_pagado", PEDIDO)
broker.procesar()
assert factura.cuenta() == 1
assert len(broker.fallidos) == 1
def test_idempotencia_mismo_id_dos_veces_una_factura(db):
broker, factura = Broker(), Factura(db)
broker.suscribir("pedido_pagado", "factura", factura)
broker.publicar("pedido_pagado", PEDIDO)
broker.publicar("pedido_pagado", PEDIDO) # duplicado
broker.procesar()
assert factura.cuenta() == 1
def test_outbox_atomico_si_falla_no_queda_pedido_ni_evento(db):
crear_pedido_pagado(db, "P-1", 10.0)
with pytest.raises(sqlite3.IntegrityError):
crear_pedido_pagado(db, "P-1", 10.0) # id repetido: revierte todo
assert db.execute("SELECT COUNT(*) FROM pedidos").fetchone()[0] == 1
assert db.execute("SELECT COUNT(*) FROM outbox").fetchone()[0] == 1
def test_outbox_con_caida_del_rele_duplica_y_la_idempotencia_lo_absorbe(db):
broker, factura = Broker(), Factura(db)
broker.suscribir("pedido_pagado", "factura", factura)
crear_pedido_pagado(db, "P-200", 800.0)
with pytest.raises(RuntimeError):
retransmitir(db, broker, caer_antes_de_marcar=True)
retransmitir(db, broker) # reinicio: vuelve a publicar
broker.procesar()
assert db.execute("SELECT COUNT(*) FROM outbox WHERE publicado = 0").fetchone()[0] == 0
assert factura.cuenta() == 1 # 2 entregas, 1 factura
Ejecútalas con python -m pytest -v. Verificadas con Python 3.12 y pytest 9: las 7 pruebas pasan. Cubren, en orden: los tres consumidores reciben el evento; un consumidor que falla dos veces tiene éxito al tercer intento (con esperas de 1 s y 2 s); un mensaje venenoso termina en la DLQ; un consumidor roto no bloquea a los demás; dos entregas del mismo pedido generan una sola factura; el outbox es atómico; y una caída del relé duplica la entrega pero la idempotencia lo absorbe.
Ejercicios sugeridos: añade un cuarto consumidor (inventario) sin modificar la tienda; sustituye la lista fallidos por una tabla SQLite; y cambia Factura por un consumidor que sume el total a un saldo para comprobar cómo el duplicado sí corrompe el resultado.
Preguntas frecuentes
¿Qué es la mensajería entre servicios?
Es una forma de comunicación asíncrona en la que un servicio deja un mensaje en un intermediario (broker) y otros lo procesan cuando pueden, sin tener que esperar la respuesta en ese momento.
¿Qué diferencia hay entre una cola y publicar/suscribir?
En una cola, cada mensaje lo procesa un solo consumidor. En publicar/suscribir, cada evento llega a todos los suscriptores interesados, y cada uno hace su parte.
¿Qué diferencia hay entre RabbitMQ y Kafka?
RabbitMQ es un broker de mensajes orientado a repartir tareas y enrutarlas con flexibilidad. Kafka es un registro distribuido de eventos que guarda el historial y permite volver a leerlo, pensado para grandes volúmenes.
¿Qué es la idempotencia y por qué importa?
Es la propiedad de que procesar el mismo mensaje dos veces dé el mismo resultado que procesarlo una. Importa porque muchos sistemas entregan cada mensaje «al menos una vez» y pueden repetirlo.
En el reel lo explicamos así
Imagina el tablón de avisos de una tienda. Cuando alguien paga, la tienda pega una sola nota: «pedido pagado». No llama a nadie ni espera respuesta. Correo, facturación y envío pasan por el tablón y leen la nota a su ritmo; cada uno tiene su propia copia pendiente, y si uno está ocupado, su aviso espera en la cola hasta que pueda atenderlo. Si alguien lee la nota dos veces, la factura no debe duplicarse: de eso se encarga la idempotencia.
Referencias
- Martin Fowler, «What do you mean by "Event-Driven"?».
- Gregor Hohpe y Bobby Woolf, Enterprise Integration Patterns: catálogo de patrones de mensajería.
- Documentación oficial de RabbitMQ: conceptos, confirmaciones y colas de mensajes fallidos.
- Documentación de Apache Kafka: particiones, offsets y semántica de entrega.
- Chris Richardson, Transactional outbox (microservices.io).
- Amazon SQS y Google Cloud Pub/Sub: servicios gestionados.
Un tablón de avisos: alguien pega una nota y cada interesado la lee a su ritmo.
Glosario: del inglés al español
| Event | evento (aviso de que algo pasó) |
|---|---|
| Queue | cola (fila de mensajes) |
| Producer / Consumer | quien publica / quien lee |