Patrón Outbox

7 min read

Tu servicio de órdenes hace dos cosas cuando se crea una orden: guarda la orden en la base de datos, y publica un evento OrderCreated a Kafka para que el servicio de inventario y el de facturación se enteren. Código simple:


orderRepository.save(order);       // 1. Se guarda en la base de datos
kafkaProducer.send("OrderCreated", order); // 2. Se publica el evento

¿Qué pasa si el servidor se cae justo entre la línea 1 y la línea 2? La orden queda guardada, pero el evento nunca se publicó. Inventario nunca se enteró. Facturación nunca se enteró. Tu base de datos y tu sistema de mensajería quedaron inconsistentes, y nadie se dio cuenta hasta que el cliente reclamó porque nunca le llegó la factura.

Este es el problema clásico de la **doble escritura** (dual write) en sistemas distribuidos, y el Outbox Pattern (Patrón Buzón de Salida) es la forma estándar de resolverlo sin necesidad de transacciones distribuidas.


¿Qué es el Outbox Pattern?

La idea central: en lugar de escribir en dos sistemas distintos (base de datos + broker de mensajes), escribís en un solo sistema — tu base de datos — dentro de una única transacción local. Guardás el cambio de negocio (la orden) y el evento a publicar (en una tabla llamada outbox) en la misma transacción. O se guardan ambos, o no se guarda ninguno.

Después, un proceso separado — llamado Message Relay — lee periódicamente la tabla outbox y publica los eventos pendientes al broker real (Kafka, RabbitMQ). Una vez publicados con éxito, los marca como enviados.


El diagrama


Guía paso a paso

Paso 1 — Crear la tabla outbox
Necesita al menos: id, tipo de evento, payload (el contenido en JSON), timestamp de creación, y un flag o timestamp de “enviado”.

Paso 2 — Escribir el cambio de negocio y el evento en la misma transacción
Ambos INSERT (a la tabla de negocio y a la tabla outbox) van dentro del mismo @Transactional. Si algo falla, ambos se revierten juntos.

Paso 3 — Implementar el Message Relay
Un proceso — puede ser un scheduled job o un mecanismo de Change Data Capture (Captura de Cambios de Datos) como Debezium — que consulta periódicamente los eventos no enviados en la tabla outbox.

Paso 4 — Publicar y marcar como enviado
Por cada evento pendiente, publicalo al broker. Si la publicación tiene éxito, marcá el evento como enviado. Si falla, se reintenta en el siguiente ciclo.

Paso 5 — Garantizar idempotencia en los consumidores
El relay puede publicar el mismo evento más de una vez si falla justo después de publicar y antes de marcar como enviado (at-least-once delivery). Los consumidores del evento deben poder procesar duplicados sin efectos secundarios.


El código

Implementamos el patrón completo en Java: el servicio de órdenes que escribe en una sola transacción, y el relay que publica los eventos pendientes.


import lombok.AllArgsConstructor;
import lombok.Getter;
import java.time.Instant;
import java.util.ArrayList;
import java.util.List;
import java.util.stream.Collectors;

// --- Entidad de negocio ---

@Getter
@AllArgsConstructor
class Order {
    private final String orderId;
    private final String customerId;
    private final double amount;
}

// --- Registro en la tabla outbox ---

@Getter
class OutboxEvent {
    private final String eventId;
    private final String eventType;
    private final String payload;
    private final Instant createdAt;
    private boolean sent;

    public OutboxEvent(String eventId, String eventType, String payload) {
        this.eventId = eventId;
        this.eventType = eventType;
        this.payload = payload;
        this.createdAt = Instant.now();
        this.sent = false;
    }

    public void markAsSent() {
        this.sent = true;
    }
}

// --- Simulación de la base de datos: una sola transacción escribe ambas tablas ---

class Database {
    private final List<Order> orders = new ArrayList<>();
    private final List<OutboxEvent> outbox = new ArrayList<>();

    // Simula una transacción: ambos INSERT se confirman juntos
    public synchronized void saveOrderWithEvent(Order order, OutboxEvent event) {
        orders.add(order);
        outbox.add(event);
        System.out.println("[DB] Transacción confirmada: orden + evento outbox guardados juntos");
    }

    public List<OutboxEvent> getPendingEvents() {
        return outbox.stream().filter(e -> !e.isSent()).toList();
    }
}

// --- Simulación del broker de mensajes ---

class MessageBroker {
    public void publish(String eventType, String payload) {
        System.out.println("[Kafka] Publicado → " + eventType + ": " + payload);
    }
}

// --- Servicio de órdenes: escribe orden + evento en una sola transacción ---

@AllArgsConstructor
class OrderService {
    private final Database database;

    public void createOrder(String orderId, String customerId, double amount) {
        Order order = new Order(orderId, customerId, amount);

        String payload = String.format(
                "{\"orderId\":\"%s\",\"customerId\":\"%s\",\"amount\":%.2f}",
                orderId, customerId, amount);

        OutboxEvent event = new OutboxEvent(
                "EVT-" + orderId, "OrderCreated", payload);

        // Ambos se guardan en la MISMA transacción local
        database.saveOrderWithEvent(order, event);
    }
}

// --- Message Relay: publica los eventos pendientes ---

@AllArgsConstructor
class OutboxRelay {
    private final Database database;
    private final MessageBroker broker;

    public void relayPendingEvents() {
        List<OutboxEvent> pending = database.getPendingEvents();
        System.out.println("[Relay] Eventos pendientes encontrados: " + pending.size());

        for (OutboxEvent event : pending) {
            try {
                broker.publish(event.getEventType(), event.getPayload());
                event.markAsSent();
                System.out.println("[Relay] Evento " + event.getEventId() + " marcado como enviado");
            } catch (Exception e) {
                System.out.println("[Relay] Falló publicación de " + event.getEventId() + ". Se reintentará.");
            }
        }
    }
}

// --- Main ---

public class OutboxPatternDemo {

    public static void main(String[] args) throws InterruptedException {
        Database database = new Database();
        MessageBroker broker = new MessageBroker();
        OrderService orderService = new OrderService(database);
        OutboxRelay relay = new OutboxRelay(database, broker);

        System.out.println("=== Se crea una orden ===");
        orderService.createOrder("ORD-500", "CUST-77", 249.90);

        System.out.println("\n=== El Message Relay corre cada pocos segundos (scheduled job) ===");
        Thread.sleep(200); // simula el intervalo del scheduler
        relay.relayPendingEvents();

        System.out.println("\n=== Segunda corrida del relay: no hay nada pendiente ===");
        relay.relayPendingEvents();
    }
}

Output esperado:


=== Se crea una orden ===
[DB] Transacción confirmada: orden + evento outbox guardados juntos

=== El Message Relay corre cada pocos segundos (scheduled job) ===
[Relay] Eventos pendientes encontrados: 1
[Kafka] Publicado → OrderCreated: {"orderId":"ORD-500","customerId":"CUST-77","amount":249.90}
[Relay] Evento EVT-ORD-500 marcado como enviado

=== Segunda corrida del relay: no hay nada pendiente ===
[Relay] Eventos pendientes encontrados: 0

Fijate que en ningún momento hicimos dos escrituras a sistemas distintos de forma independiente. La orden y el evento nacieron juntos, en la misma transacción de base de datos. El relay es quien se encarga de la entrega al broker, con reintentos si algo falla.


Para tener en cuenta en producción

  • Change Data Capture es la implementación más robusta: en lugar de un scheduled job haciendo polling (consultando repetidamente), herramientas como Debezium leen directamente el log de transacciones de la base de datos y publican los eventos casi en tiempo real, sin sobrecargar la tabla outbox con SELECTs constantes.
  • Limpiá la tabla outbox periódicamente: los eventos ya enviados no necesitan quedarse ahí para siempre. Un job de limpieza que borra eventos enviados con más de X días evita que la tabla crezca indefinidamente.
  • Conecta directo con el Patrón Saga: si ya leíste el post de Saga, notarás que el Outbox Pattern resuelve exactamente el problema de “¿cómo notifico de forma confiable el resultado de una transacción local a los demás pasos de la Saga?”. Se usan juntos con frecuencia.
  • At-least-once, nunca exactly-once garantizado por el patrón en sí: el relay puede publicar el mismo evento dos veces si falla entre publicar y marcar como enviado. La idempotencia del consumidor no es opcional, es parte del contrato.

Conclusión

El Outbox Pattern no elimina la complejidad de los sistemas distribuidos, la reubica: en lugar de coordinar dos sistemas distintos en el momento crítico, aprovechás la garantía transaccional que ya tiene tu base de datos y delegás la entrega al broker a un proceso separado y reintentable.

Es uno de esos patrones que, una vez que lo entendés, te hace mirar con sospecha cualquier código que haga un save() seguido de un publish() como si fueran una sola operación atómica. No lo son, y ahora sabés cómo arreglarlo.

Recordá suscribirte aquí para recibir los próximos posts directamente en tu correo.


Referencias:
Microservices Patterns — Chris Richardson
Debezium Documentation — https://debezium.io/documentation/

Leave a Reply

Your email address will not be published. Required fields are marked *