Saltar al contenido

El patrón Transactional Outbox: resolver la doble escritura

Cómo el patrón Transactional Outbox elimina la doble escritura y garantiza la publicación fiable de eventos sin acoplar la base de datos al broker.

6 min de lectura
Diagrama del patrón Transactional Outbox que muestra una tabla outbox de base de datos alimentando un proceso de relay de mensajes

Casi todos los sistemas distribuidos esconden este bug en algún rincón del código: guardar el pedido en la base de datos y después publicar un evento order.created en Kafka. Cuando la publicación falla —y tarde o temprano fallará—, el pedido ya existe, pero todos los servicios que dependen de él quedan ciegos para siempre. Se puede añadir lógica de reintentos, pero eso solo desplaza el problema. La solución real es arquitectónica: dejar de tratar la escritura en la base de datos y la publicación en el broker como dos operaciones separadas.

El patrón Transactional Outbox resuelve esto a nivel de infraestructura. Cambia la ilusión de atomicidad entre dos sistemas por una escritura atómica real en uno solo, y luego propaga el cambio de forma asíncrona mediante un proceso de relay independiente. El resultado es consistente, fácil de depurar y escalable horizontalmente.

El problema de la doble escritura

Cada vez que escribes en dos sistemas distintos de forma secuencial, estás apostando a que no haya un fallo parcial. El caso más común es escribir en una base de datos relacional y en un broker de mensajes dentro del mismo request handler:

tstypescript
// ❌ Dual-write — the publish can fail after the DB commit
async function createOrder(data: CreateOrderInput): Promise<Order> {
  const order = await db.orders.create(data);
  await kafka.publish("order.created", order); // what if this throws?
  return order;
}
 
// ✅ Single transaction — the event is part of the same write
async function createOrder(data: CreateOrderInput): Promise<Order> {
  return db.transaction(async (tx) => {
    const order = await tx.orders.create(data);
    await tx.outbox.insert({
      aggregateType: "order",
      aggregateId: order.id,
      eventType: "order.created",
      payload: order,
    });
    return order;
  });
}

La fila outbox y la fila del pedido se confirman de forma atómica. Si la transacción hace rollback, ni el pedido ni la entrada outbox llegan a existir. Si se confirma, existen ambas. El broker de mensajes nunca está en el camino crítico de la transacción: se entera del evento más tarde, a través de un proceso de relay que lee las filas outbox ya confirmadas.

El esquema de la tabla outbox

La tabla outbox es el corazón del patrón. Diséñala desde el primer día para soportar un procesamiento ordenado y fiable:

sqlsql
CREATE TABLE outbox_events (
  id             UUID        DEFAULT gen_random_uuid() PRIMARY KEY,
  aggregate_type TEXT        NOT NULL,
  aggregate_id   TEXT        NOT NULL,
  event_type     TEXT        NOT NULL,
  payload        JSONB       NOT NULL,
  created_at     TIMESTAMPTZ DEFAULT now() NOT NULL,
  processed_at   TIMESTAMPTZ,
  attempts       INT         DEFAULT 0 NOT NULL,
  last_error     TEXT
);
 
CREATE INDEX idx_outbox_unprocessed
  ON outbox_events (created_at)
  WHERE processed_at IS NULL;

El índice parcial sobre processed_at IS NULL es fundamental. Cuando la tabla acumula millones de filas ya procesadas, el proceso de relay solo recorre el pequeño subconjunto sin procesar. Sin ese índice, cada poll se convierte en un escaneo secuencial completo que se degrada con el tiempo.

Las columnas attempts y last_error no son un lujo opcional: son la forma de detectar eventos atascados en producción sin tener que rastrear logs a mano.

Cómo construir el proceso de relay

El relay hace polling de eventos sin procesar y los publica en el broker. Mantenlo simple y prioriza la corrección antes que nada:

tstypescript
async function processOutboxBatch(batchSize = 100): Promise<number> {
  const events = await db.$transaction(async (tx) => {
    // Lock rows to prevent concurrent relays from double-publishing
    const rows = await tx.$queryRaw<OutboxEvent[]>`
      SELECT * FROM outbox_events
      WHERE processed_at IS NULL
      ORDER BY created_at ASC
      LIMIT ${batchSize}
      FOR UPDATE SKIP LOCKED
    `;
 
    if (rows.length === 0) return [];
 
    await Promise.all(
      rows.map((event) =>
        kafka.publish(event.event_type, event.payload, {
          key: event.aggregate_id, // same aggregate → same partition → ordered delivery
        })
      )
    );
 
    const ids = rows.map((r) => r.id);
    await tx.$executeRaw`
      UPDATE outbox_events
      SET processed_at = now()
      WHERE id = ANY(${ids}::uuid[])
    `;
 
    return rows;
  });
 
  return events.length;
}
 
// Relay loop — run one per service replica
async function startRelay(): Promise<void> {
  while (true) {
    const processed = await processOutboxBatch();
    // Back off when idle to avoid hammering the DB unnecessarily
    await sleep(processed === 0 ? 500 : 50);
  }
}

FOR UPDATE SKIP LOCKED es el detalle que hace posible el escalado horizontal. Varias instancias del relay pueden ejecutarse en paralelo: cada una toma las filas que ningún otro relay tiene bloqueadas. Así se obtiene concurrencia segura sin necesidad de un lock distribuido ni un cuello de botella de réplica única.

Polling frente a Change Data Capture

El polling es el punto de partida correcto para la mayoría de los equipos. Con el índice parcial, un intervalo de poll de 200–500ms produce una latencia inferior al segundo y una carga insignificante sobre la base de datos.

EnfoqueLatenciaCarga en BDComplejidad operativa
Polling (intervalo de 500ms)~500msBajaMínima — un loop por réplica
Polling (intervalo de 100ms)~100msModeradaMínima
CDC (Debezium + Kafka Connect)Casi en tiempo realMuy bajaAlta — clúster de Connect, schema registry, gestión de conectores

Change Data Capture lee directamente del stream WAL de PostgreSQL, eliminando el polling por completo. Es la opción correcta cuando necesitas latencia inferior a 100ms a gran escala, o cuando la tabla outbox está bajo una presión de escritura extrema. Para la inmensa mayoría de los servicios, el polling es el punto de partida correcto y debería seguir siendo la implementación hasta contar con datos de producción que justifiquen la sobrecarga operativa de Debezium.

~

Instrumenta tu relay con una métrica: la antigüedad del evento outbox sin procesar más antiguo. Si ese número se mantiene por debajo de un segundo en estado estable, el polling está cumpliendo su función. Recurre a CDC solo cuando ese SLO sea, de verdad, inalcanzable con polling.

La idempotencia del consumidor no es opcional

El patrón outbox garantiza la entrega at-least-once (al menos una vez), nunca exactly-once. Si el relay publica un evento correctamente pero el commit de processed_at falla antes de que la transacción se cierre (un corte de red, un proceso que muere a mitad de camino), ese mismo evento se volverá a publicar en el siguiente ciclo de poll. Es un comportamiento esperado, no un bug.

Todo consumidor debe manejar duplicados:

tstypescript
async function handleOrderCreated(event: OrderCreatedEvent): Promise<void> {
  const alreadyProcessed = await db.processedEvents.findUnique({
    where: { eventId: event.id },
  });
 
  if (alreadyProcessed) {
    logger.debug({ eventId: event.id }, "skipping duplicate event");
    return;
  }
 
  // Transactionally apply the effect and record the deduplication key
  await db.$transaction(async (tx) => {
    await tx.inventory.reserveItems(event.payload.lineItems);
    await tx.processedEvents.create({ data: { eventId: event.id } });
  });
}

La tabla processedEvents funciona como un registro de deduplicación. La verificación seguida de la escritura ocurre dentro de una transacción con una restricción unique sobre eventId: cuando dos consumidores compiten por el mismo evento, uno tendrá éxito y el otro recibirá una violación de restricción, que es seguro ignorar tratándola como un duplicado.

Retención y limpieza

Las filas outbox ya procesadas son peso muerto. Un job de limpieza programado evita que la tabla crezca sin límite:

sqlsql
-- Run daily via pg_cron or your job scheduler
DELETE FROM outbox_events
WHERE processed_at IS NOT NULL
  AND processed_at < now() - INTERVAL '7 days';

Siete días de retención te dan una ventana de replay para depurar incidentes, sin que el audit trail crezca para siempre. Si tu turno de guardia suele resolver los problemas en menos de 24 horas, con tres días es suficiente. El índice parcial mantiene estable el rendimiento de las consultas sin importar cuántos datos históricos conserves, pero las tablas más pequeñas significan vacuums más rápidos y menores costos de almacenamiento.

Conclusiones clave

  1. Nunca hagas dual-write — escribir en una base de datos y publicar después en un broker es un bug de consistencia esperando a manifestarse en el peor momento posible
  2. La fila outbox es tu evento — confírmala de forma atómica junto con tu escritura de dominio; nada más debe cruzar el límite de la transacción
  3. FOR UPDATE SKIP LOCKED permite escalar el relay horizontalmente de forma segura — varias instancias se coordinan a través de la propia base de datos, sin necesidad de locks externos
  4. Empieza con polling, mide, y luego evalúa CDC — la latencia casi en tiempo real rara vez justifica el costo operativo de Debezium al principio del ciclo de vida de un sistema
  5. La entrega at-least-once es un contrato, no un bug — diseña cada consumidor para que sea idempotente desde el primer día, no como un parche posterior
  6. Registra la antigüedad del evento sin procesar más viejo — es la métrica más útil para medir la salud del relay, más informativa que el throughput de eventos procesados por segundo
Wilfredo Rujel

Wilfredo Rujel

Ingeniero de Software Full Stack

Compartir esta publicaciónX