⏱️ Lectura: 14 min
Guardar un pedido en la base de datos y publicar el evento correspondiente en Kafka parecen dos pasos triviales, pero nunca son atómicos: si el proceso muere entre uno y otro, el pedido existe sin que el resto del sistema se entere. El patrón outbox transaccional resuelve ese problema sin transacciones distribuidas ni bloqueos entre sistemas distintos.
📑 En este artículo
El problema se conoce como dual write: dos escrituras que deberían ocurrir juntas, en dos sistemas que no comparten una transacción. Aparece entre una base de datos y una cola de mensajes, entre dos bases distintas, o entre una base y una llamada a una API externa.
TL;DR
- El patrón outbox transaccional guarda el evento en la misma transacción que el dato de negocio, en una tabla outbox.
- SELECT … FOR UPDATE SKIP LOCKED evita que dos relays procesen el mismo evento outbox dos veces.
- Debezium lee el WAL de PostgreSQL y publica cada fila nueva de la tabla outbox en Kafka sin tocarla con SELECTs.
- Los consumidores deben ser idempotentes porque el outbox garantiza at-least-once, nunca exactly-once.
- Una tabla outbox sin purgar crece sin límite: hace falta un job que borre o archive los eventos ya publicados.
¿Qué es el patrón outbox transaccional?
El patrón outbox transaccional es un patrón de diseño que garantiza que una escritura en la base de datos y la publicación de un evento asociado ocurran de forma atómica: guarda el evento como una fila más, dentro de la misma transacción, y delega su envío real a un proceso independiente llamado relay.
Nació como respuesta a un problema concreto en arquitecturas de microservicios: cuando un servicio necesita persistir un cambio y además notificar a otros servicios de ese cambio, hacerlo en dos pasos separados deja una ventana donde el sistema puede quedar en un estado a medias. Chris Richardson documentó el patrón en su catálogo de patrones de microservicios, y Gunnar Morling, del equipo de Debezium, popularizó su combinación con captura de datos modificados en un artículo de referencia de 2019.
Por qué importa
Un servicio que actualiza su base de datos y después publica un mensaje tiene dos escrituras que deberían pasar o fallar juntas, pero corren en dos sistemas que no comparten transacción. Si el proceso muere justo después del commit en la base y antes de publicar, el evento se pierde para siempre: el pedido existe, pero nadie se entera. Si el orden se invierte y se publica antes de confirmar en la base, puede llegar un evento sobre un pedido que después falla y nunca se guarda.
La solución ingenua es envolver ambas operaciones en una transacción distribuida con dos fases (2PC). Funciona, pero exige que el broker de mensajes soporte XA, agrega bloqueos entre sistemas y no escala bien. Kafka, por ejemplo, no ofrece soporte nativo para transacciones XA con una base de datos externa. La clave del patrón outbox transaccional es evitar el problema por completo. En vez de escribir en dos sistemas, se escribe en uno solo.
flowchart TD
A["Servicio procesa un pedido"] --> B["Guarda el pedido en la base"]
A --> C["Publica el evento en Kafka"]
B --> D["Confirmado: el pedido queda guardado"]
C --> E["Falla de red: el evento se pierde"]
Cómo funciona el outbox por dentro
La tabla outbox vive en la misma base de datos que las tablas de negocio. Cuando el servicio necesita persistir un cambio y emitir un evento, hace ambas escrituras en la misma transacción: un INSERT en la tabla de negocio (por ejemplo, orders) y un INSERT en la tabla outbox con el payload del evento. Si la transacción falla, ninguna de las dos escrituras queda registrada; si confirma, las dos quedan juntas, siempre.
Un proceso separado, llamado relay o publisher, lee periódicamente las filas nuevas de la tabla outbox y las envía a la cola de mensajes real (Kafka, RabbitMQ, SQS). Una vez que confirma el envío, marca la fila como publicada. Este relay puede implementarse de dos formas: por polling, con un SELECT cada cierto intervalo, o por CDC, leyendo el log de transacciones de la base de datos sin ejecutar consultas contra la tabla.
sequenceDiagram
participant S as Servicio
participant DB as Base de datos
participant R as Relay
participant K as Kafka
S->>DB: BEGIN transacción
S->>DB: INSERT en orders
S->>DB: INSERT en outbox
S->>DB: COMMIT
R->>DB: SELECT eventos pendientes
DB-->>R: filas de la tabla outbox
R->>K: publica el evento
R->>DB: UPDATE published_at
Ejemplos prácticos y cómo empezar
El siguiente ejemplo implementa un outbox con Node.js y PostgreSQL, sin depender de Kafka ni de infraestructura adicional: primero un relay por polling simple, después la mejora con SKIP LOCKED para que puedan correr varias instancias del relay al mismo tiempo.
Requisitos
Node.js 18 o superior, una instancia de PostgreSQL corriendo localmente o en Docker, y la librería pg. No hace falta Kafka para este ejemplo: el relay puede simular la publicación con un log antes de conectarlo a un broker real.
npm init -y
npm install pg
El esquema tiene dos tablas: la tabla de negocio orders y la tabla outbox, que guarda cada evento pendiente de publicar.
CREATE TABLE orders (
id SERIAL PRIMARY KEY,
customer TEXT NOT NULL,
total NUMERIC NOT NULL
);
CREATE TABLE outbox (
id SERIAL PRIMARY KEY,
aggregate_type TEXT NOT NULL,
aggregate_id INTEGER NOT NULL,
event_type TEXT NOT NULL,
payload JSONB NOT NULL,
created_at TIMESTAMPTZ NOT NULL DEFAULT now(),
published_at TIMESTAMPTZ
);
El servicio que crea un pedido inserta en las dos tablas dentro de la misma transacción. Si el COMMIT falla, ninguna de las dos filas queda guardada.
const { Pool } = require('pg');
const pool = new Pool();
async function crearPedido(customer, total) {
const client = await pool.connect();
try {
await client.query('BEGIN');
const { rows } = await client.query(
'INSERT INTO orders (customer, total) VALUES ($1, $2) RETURNING id',
[customer, total]
);
const orderId = rows[0].id;
await client.query(
`INSERT INTO outbox (aggregate_type, aggregate_id, event_type, payload)
VALUES ($1, $2, $3, $4)`,
['order', orderId, 'OrderCreated', JSON.stringify({ orderId, customer, total })]
);
await client.query('COMMIT');
return orderId;
} catch (err) {
await client.query('ROLLBACK');
throw err;
} finally {
client.release();
}
}
Al correr crearPedido('Ana', 49.90), la tabla orders gana una fila y la tabla outbox gana otra, atómicamente. Una consulta de verificación confirma que ambas existen:
SELECT o.id, o.customer, x.event_type, x.published_at
FROM orders o JOIN outbox x ON x.aggregate_id = o.id
WHERE o.id = 1;
id | customer | event_type | published_at
----+----------+---------------+--------------
1 | Ana | OrderCreated | NULL
published_at en NULL significa que el evento todavía no salió. Ahí entra el relay: un proceso que hace polling sobre la tabla outbox, publica lo que encuentra y marca la fila.
async function relayPolling() {
const { rows } = await pool.query(
`SELECT * FROM outbox WHERE published_at IS NULL ORDER BY id LIMIT 50
FOR UPDATE SKIP LOCKED`
);
for (const evento of rows) {
await publicarEnKafka(evento.event_type, evento.payload);
await pool.query('UPDATE outbox SET published_at = now() WHERE id = $1', [evento.id]);
}
}
setInterval(relayPolling, 2000);
FOR UPDATE SKIP LOCKED es la pieza clave si corrés más de una instancia del relay: cada instancia se queda solo con las filas que nadie más está procesando en ese momento, en vez de bloquearse esperando a la otra. Sin esa cláusula, dos relays compitiendo por la misma fila terminan publicando el mismo evento dos veces o se traban entre sí.
💡 Tip: correr varias instancias del relay tolera caídas, no multiplica el rendimiento. Con SKIP LOCKED, cada evento lo procesa una sola instancia, así que sumar relays agrega disponibilidad, no velocidad de publicación.
Del otro lado, el consumidor necesita su propia tabla de deduplicación, a veces llamada inbox, para ignorar un evento que ya procesó antes.
CREATE TABLE processed_events (
event_id INTEGER PRIMARY KEY,
processed_at TIMESTAMPTZ NOT NULL DEFAULT now()
);
async function manejarEventoOrderCreated(evento) {
const client = await pool.connect();
try {
await client.query('BEGIN');
const { rowCount } = await client.query(
'INSERT INTO processed_events (event_id) VALUES ($1) ON CONFLICT DO NOTHING',
[evento.id]
);
if (rowCount === 0) {
await client.query('ROLLBACK');
return; // ya se procesó antes, se ignora
}
await client.query(
'UPDATE inventory SET reserved = reserved + 1 WHERE product_id = $1',
[evento.payload.productId]
);
await client.query('COMMIT');
} catch (err) {
await client.query('ROLLBACK');
throw err;
} finally {
client.release();
}
}
Casos de uso reales
El caso más común es sincronizar un cambio de estado con una notificación a otros servicios: un pedido que pasa a pagado y dispara el envío de un correo, la actualización de inventario y el cálculo de comisiones, todo en servicios distintos que no comparten base de datos.
Otro uso frecuente es alimentar un sistema de búsqueda o un caché derivado: cada escritura en la tabla principal genera un evento outbox que un consumidor usa para reindexar en Elasticsearch o invalidar una entrada en Redis, sin acoplar la transacción principal a la disponibilidad de esos sistemas.
También aparece como base para event sourcing parcial: en vez de reconstruir todo el estado desde eventos, un servicio guarda su estado normal en tablas relacionales y usa el outbox solo para exponer los cambios como flujo de eventos hacia el resto de la organización.
Errores comunes y buenas prácticas
- Publicar antes de confirmar la transacción. Si el relay llegara a leer un evento cuya transacción todavía no hizo commit, un consumidor podría procesar un pedido que después se cancela con un rollback. Siempre hay que leer solo filas ya confirmadas.
- Tabla outbox sin purgar. Si nunca se borran ni archivan las filas publicadas, la tabla crece sin límite y termina degradando los índices de toda la base. Un job periódico que borre filas con
published_atanterior a unos días es obligatorio en producción. - Consumidores no idempotentes. El outbox da garantías at-least-once: un mismo evento puede llegar duplicado si el relay falla justo después de publicar y antes de marcar la fila. El consumidor necesita una clave de deduplicación, y el propio id de la fila outbox sirve.
- Orden de eventos entre agregados distintos. Publicar sin clave de partición no preserva el orden dentro de un mismo pedido. Hay que usar
aggregate_idcomo clave al enviar a Kafka:
await kafkaProducer.send({
topic: 'orders',
messages: [{ key: String(evento.aggregate_id), value: JSON.stringify(evento.payload) }]
});
Comparativa con alternativas
| Opción | Cuándo usarla | Ventaja | Limitación |
|---|---|---|---|
| Outbox con polling | Volumen bajo o medio, sin infraestructura de CDC disponible | Fácil de implementar con lo que ya tenés | Latencia igual al intervalo de polling y carga extra de SELECTs |
| Outbox con CDC (Debezium) | Alto volumen, baja latencia, ya operás Kafka | Latencia mínima, lee el log de la base sin tocarla | Requiere desplegar y mantener Kafka Connect |
| Transacciones distribuidas (2PC/XA) | Sistemas legacy con soporte XA en ambos extremos | Atomicidad real entre los dos recursos | Bloqueante, no escala y la mayoría de los brokers modernos no lo soportan |
Profundizando
La implementación con CDC reemplaza el polling por lectura directa del log de transacciones, el WAL en PostgreSQL. Debezium corre como conector de Kafka Connect, se suscribe a la replicación lógica de PostgreSQL y emite un evento de Kafka por cada fila nueva en la tabla outbox, sin ejecutar un solo SELECT contra ella. La ventaja no es solo de latencia. Un relay por polling compite por locks con las transacciones normales de la aplicación, mientras que leer el WAL no interfiere con nada.
Ninguna de las dos variantes da exactly-once real. Tanto el polling como Debezium pueden publicar el mismo evento más de una vez si el proceso falla entre publicar y confirmar el offset. La garantía que sí sostienen es at-least-once con orden preservado por partición, y eso empuja la responsabilidad de la deduplicación al consumidor.
El ciclo de vida de un evento outbox pasa por estados bien definidos, y modelarlos así ayuda a decidir qué reintentar:
stateDiagram-v2
[*] --> Pendiente
Pendiente --> Publicado: el relay lo envía a Kafka
Publicado --> Confirmado: el consumidor lo procesa
Publicado --> Pendiente: fallo de red, se reintenta
Confirmado --> [*]
Monitorear el lag del outbox
La señal de salud más simple es medir cuánto tiempo llevan sin publicarse los eventos pendientes. Si ese número crece de forma sostenida, el relay dejó de mantenerse al día.
SELECT extract(epoch FROM now() - min(created_at)) AS lag_seconds
FROM outbox
WHERE published_at IS NULL;
lag_seconds
-------------
3.42
⚠️ Ojo: si varios servicios escriben en la misma tabla outbox con relojes de servidor distintos, ordenar solo por created_at puede procesar eventos casi simultáneos en un orden distinto al real. El id autoincremental de la tabla es una referencia de orden más confiable que el timestamp.
Un detalle que rara vez se documenta: el outbox no necesita una tabla nueva por cada tipo de agregado. Un esquema genérico con aggregate_type, aggregate_id, event_type y payload en JSONB sirve para toda la aplicación, y es el mismo esquema que espera el conector de outbox de Debezium sin transformar nada más.
📖 Resumen en Telegram: Ver resumen
Tu próximo paso: montá las dos tablas de este artículo en una base PostgreSQL local, corré crearPedido tres veces seguidas y confirmá con el SELECT de verificación que cada pedido tiene su fila outbox antes de escribir el relay.
Preguntas frecuentes
¿Qué problema resuelve el patrón outbox transaccional?
Resuelve el dual write: evita que una escritura en la base de datos y la publicación de un evento en una cola de mensajes queden inconsistentes cuando el proceso falla entre una y otra.
¿El outbox transaccional garantiza exactly-once?
No. Garantiza at-least-once: un evento puede llegar duplicado al consumidor, que necesita ser idempotente para procesarlo una sola vez pese a los reintentos.
¿Hace falta Kafka para usar un outbox?
No necesariamente. El outbox funciona con cualquier cola de mensajes, como RabbitMQ o SQS, o incluso con un relay que llama directamente a una API externa. Kafka es solo la opción más común en arquitecturas de microservicios.
¿Cuál es la diferencia entre el outbox por polling y con Debezium?
El polling ejecuta un SELECT periódico sobre la tabla outbox. Debezium lee el log de transacciones y no toca la tabla con consultas, lo que reduce la latencia y la carga sobre la base.
¿Cómo se evita que la tabla outbox crezca sin límite?
Con un job periódico que borra o archiva las filas con published_at distinto de NULL después de un tiempo razonable, separando el histórico si hace falta conservarlo para auditoría.
Referencias
- microservices.io: catálogo de patrones de Chris Richardson con la definición original del transactional outbox.
- Debezium Blog: artículo de Gunnar Morling sobre cómo combinar el outbox con CDC.
- PostgreSQL Docs: referencia oficial de la cláusula
SELECT ... FOR UPDATE SKIP LOCKED. - Wikipedia: definición general de change data capture (CDC).
📱 ¿Te gusta este contenido? Únete a nuestro canal de Telegram @programacion donde publicamos a diario lo más relevante de tecnología, IA y desarrollo. Resúmenes rápidos, contenido fresco todos los días.
Imagen destacada: Foto de Markus Spiske en Unsplash
¿Te sirvió? ¿Te dio otro error? Contalo abajo: las preguntas se responden y le sirven al siguiente que llegue.
Dejar un comentario
0 Comentarios