Patrón Outbox transaccional en Go con PostgreSQL

Escribe el evento con los datos. Nunca los dividas.

Índice

Dos escrituras que deberían tener éxito juntas eventualmente fallan por separado.

Tu servicio de pedidos guarda el pedido en la base de datos y luego publica un evento order.created en un corredor de mensajes.

Estas dos operaciones se ejecutan una después de la otra.

Entre ellas, algo sale mal: el corredor está caído, la red se agota, el proceso se reinicia o el contenedor es eliminado. La escritura en la base de datos tuvo éxito. La publicación no lo hizo. El servicio descendente que necesita saber sobre el nuevo pedido nunca se entera. Nadie se dio cuenta hasta que un cliente llamó.

Este es el problema de la doble escritura, y es una de las fuentes más comunes de pérdida silenciosa de datos en sistemas distribuidos. El patrón del buzón transaccional es la solución estándar.

Patrón del buzón transaccional – evento y datos escritos juntos

El problema de la doble escritura

El modo de fallo es fácil de razonar una vez que lo ves:

BEGIN;
  INSERT INTO orders ...   -- tiene éxito
COMMIT;

PUBLISH order.created ...  -- falla, se bloquea o nunca se alcanza

La base de datos y el corredor de mensajes no comparten un límite de transacción. No hay un rollback que cubra ambos. Cada servicio que realiza save -> publish en secuencia tiene esta brecha. El patrón aparece en muchas formas:

  • db.Save(order) seguido de events.Publish(OrderCreated{...})
  • Manejador HTTP que confirma una transacción y luego llama a un webhook externo
  • Trabajador que procesa un registro de una cola y escribe los resultados en otra

El resultado en todos los casos es el mismo: un lado tiene éxito mientras el otro falla, y el sistema termina en un estado invisible para la monitorización porque ambas operaciones individuales devolvieron éxito en algún momento.

Un bucle de reintento no soluciona esto. Reintentar la publicación después del commit de la base de datos solo funciona si el reintento en sí es confiable, lo cual requiere la garantía de durabilidad que no tienes.

Qué hace el patrón del buzón transaccional

El patrón del buzón elimina la brecha eliminando por completo la publicación directa. En lugar de llamar al corredor desde dentro de tu lógica comercial, escribes un registro de evento en una tabla outbox en la misma transacción de base de datos que los datos comerciales. Un proceso en segundo plano separado, el relé, lee desde la tabla del buzón y publica en el corredor.

BEGIN;
  INSERT INTO orders ...         -- datos comerciales
  INSERT INTO outbox_events ...  -- registro de evento
COMMIT;

-- Proceso de relé (separado):
SELECT ... FROM outbox_events FOR UPDATE SKIP LOCKED;
PUBLISH order.created ...
UPDATE outbox_events SET processed_at = NOW() WHERE id = $1;

Ambas escrituras tienen éxito o ambas fallan. La garantía de transacción que ya tienes de PostgreSQL ahora cubre también el registro de evento. El relé puede reintentar la publicación tantas veces como sea necesario porque el evento reside en almacenamiento durable. Si el relé se bloquea a mitad de ejecución, se reinicia y reintenta. El peor resultado es que el evento se publique más de una vez, lo cual se maneja haciendo que los consumidores sean idempotentes (ver Idempotencia en Sistemas Distribuidos).

Esquema de PostgreSQL para la tabla del buzón

El esquema es intencionalmente simple:

CREATE TABLE outbox_events (
    id             UUID         PRIMARY KEY DEFAULT gen_random_uuid(),
    aggregate_type VARCHAR(100) NOT NULL,
    aggregate_id   VARCHAR(100) NOT NULL,
    event_type     VARCHAR(100) NOT NULL,
    payload        JSONB        NOT NULL,
    attempts       INT          NOT NULL DEFAULT 0,
    created_at     TIMESTAMPTZ  NOT NULL DEFAULT NOW(),
    processed_at   TIMESTAMPTZ
);

-- Índice parcial: solo indexa filas no procesadas, permanece pequeño a medida que las filas se marcan como hechas
CREATE INDEX idx_outbox_unprocessed
    ON outbox_events (created_at)
    WHERE processed_at IS NULL;

El índice parcial en created_at WHERE processed_at IS NULL es importante. Sin él, el índice crece con cada evento escrito nunca y la consulta de sondeo del relé se vuelve más lenta con el tiempo. Con él, el índice cubre solo las filas pendientes, que en estado estable son un conjunto pequeño y acotado independientemente de cuántos eventos se hayan publicado.

Elecciones clave de campos:

  • aggregate_type y aggregate_id describen a qué entidad pertenece el evento. Útil para garantías de ordenación y enrutamiento.
  • event_type es el nombre del evento que tus consumidores esperan.
  • payload JSONB almacena el cuerpo del evento. Usa JSONB en lugar de TEXT para que puedas consultarlo si es necesario.
  • attempts rastrea cuántas veces el relé ha intentado publicar esta fila. Se usa para límites de reintento y manejo de mensajes envenenados.
  • processed_at es NULL para filas pendientes y se establece cuando el relé publica con éxito.

Escribir datos comerciales y evento del buzón en una transacción

La lógica comercial escribe ambos registros dentro de una única llamada BeginTx / Commit. No hay llamada de publicación aquí, solo escrituras de base de datos.

type OrderService struct {
    db *sql.DB
}

func (s *OrderService) CreateOrder(ctx context.Context, order Order) error {
    tx, err := s.db.BeginTx(ctx, nil)
    if err != nil {
        return fmt.Errorf("begin tx: %w", err)
    }
    defer tx.Rollback()

    if _, err := tx.ExecContext(ctx, `
        INSERT INTO orders (id, customer_id, total, created_at)
        VALUES ($1, $2, $3, NOW())
    `, order.ID, order.CustomerID, order.Total); err != nil {
        return fmt.Errorf("insert order: %w", err)
    }

    payload, err := json.Marshal(map[string]any{
        "order_id":    order.ID,
        "customer_id": order.CustomerID,
        "total":       order.Total,
    })
    if err != nil {
        return fmt.Errorf("marshal payload: %w", err)
    }

    if _, err := tx.ExecContext(ctx, `
        INSERT INTO outbox_events
            (aggregate_type, aggregate_id, event_type, payload)
        VALUES ($1, $2, $3, $4)
    `, "order", order.ID, "order.created", payload); err != nil {
        return fmt.Errorf("insert outbox event: %w", err)
    }

    return tx.Commit()
}

Si tx.Commit() falla, ni la fila del pedido ni la fila del buzón se persisten. Si tiene éxito, ambas están garantizadas para estar en la base de datos. El relé puede publicar el evento en cualquier momento después de eso, inmediatamente, en un segundo o después de que el relé se reinicie tras un bloqueo.

Este es el único cambio de código requerido en tu capa comercial. El resto del patrón vive en el relé.

Implementación del relé en Go

El relé es un trabajador en segundo plano que realiza sondeos en la tabla del buzón en un temporizador. Busca un lote de filas no procesadas, publica cada una y la marca como terminada. Mantén el mismo binario que tu aplicación o ejecútalo como un proceso separado, ambos funcionan, pero el mismo binario es más simple de operar.

type OutboxRelay struct {
    db          *sql.DB
    publisher   Publisher
    logger      *slog.Logger
    batchSize   int
    pollInterval time.Duration
    maxAttempts  int
}

func (r *OutboxRelay) Run(ctx context.Context) error {
    ticker := time.NewTicker(r.pollInterval)
    defer ticker.Stop()

    for {
        select {
        case <-ctx.Done():
            return ctx.Err()
        case <-ticker.C:
            if err := r.processBatch(ctx); err != nil {
                r.logger.Error("outbox relay batch failed", "err", err)
            }
        }
    }
}

El relé respeta la cancelación del contexto, lo que facilita su integración con el apagado graceful. Para un tratamiento detallado de la duración del contexto y los patrones de cancelación, ver Go context.Context Done Right).

FOR UPDATE SKIP LOCKED: el patrón de trabajador concurrente

La función processBatch usa FOR UPDATE SKIP LOCKED para manejar de forma segura los trabajadores del relé concurrentes:

func (r *OutboxRelay) processBatch(ctx context.Context) error {
    tx, err := r.db.BeginTx(ctx, nil)
    if err != nil {
        return fmt.Errorf("begin tx: %w", err)
    }
    defer tx.Rollback()

    rows, err := tx.QueryContext(ctx, `
        SELECT id, aggregate_type, aggregate_id, event_type, payload
        FROM outbox_events
        WHERE processed_at IS NULL
          AND attempts < $1
        ORDER BY created_at
        LIMIT $2
        FOR UPDATE SKIP LOCKED
    `, r.maxAttempts, r.batchSize)
    if err != nil {
        return fmt.Errorf("query outbox: %w", err)
    }
    defer rows.Close()

    type row struct {
        id            string
        aggregateType string
        aggregateID   string
        eventType     string
        payload       json.RawMessage
    }

    var batch []row
    for rows.Next() {
        var e row
        if err := rows.Scan(
            &e.id, &e.aggregateType, &e.aggregateID, &e.eventType, &e.payload,
        ); err != nil {
            return fmt.Errorf("scan row: %w", err)
        }
        batch = append(batch, e)
    }
    if err := rows.Err(); err != nil {
        return err
    }

    for _, e := range batch {
        if err := r.publisher.Publish(ctx, e.eventType, e.aggregateID, e.payload); err != nil {
            r.logger.Error("publish failed", "event_id", e.id, "err", err)
            if _, err := tx.ExecContext(ctx,
                `UPDATE outbox_events SET attempts = attempts + 1 WHERE id = $1`, e.id,
            ); err != nil {
                r.logger.Error("increment attempts failed", "event_id", e.id, "err", err)
            }
            continue
        }

        if _, err := tx.ExecContext(ctx,
            `UPDATE outbox_events SET processed_at = NOW() WHERE id = $1`, e.id,
        ); err != nil {
            return fmt.Errorf("mark processed: %w", err)
        }
    }

    return tx.Commit()
}

FOR UPDATE SKIP LOCKED hace dos cosas. Primero, FOR UPDATE bloquea las filas seleccionadas durante la duración de la transacción, impidiendo que cualquier otra transacción las seleccione. Segundo, SKIP LOCKED significa que si una fila ya está bloqueada por otra transacción, la consulta la omite en lugar de esperar. El resultado es que múltiples trabajadores del relé pueden ejecutarse en paralelo y cada uno recogerá un subconjunto no superpuesto de filas.

Sin SKIP LOCKED, un segundo trabajador bloquearía hasta que la primera transacción confirme antes de ver las mismas filas, momento en el cual ya estarían marcadas como terminadas. Con SKIP LOCKED, el segundo trabajador recoge inmediatamente filas diferentes en lugar de esperar, lo que te da una escalabilidad horizontal segura.

Nota la separación de escaneo-entonces-publicación en el código anterior: todas las filas se escanean en un slice antes de que comience el bucle de publicación. Esto evita mantener un cursor *sql.Rows abierto durante llamadas de red al corredor, lo que mantendría la transacción abierta más tiempo del necesario.

Idempotencia y deduplicación

El relé publica al menos una vez. Si publica un evento y luego se bloquea antes de confirmar la actualización de processed_at, publicará el mismo evento nuevamente al reiniciar. Esto es inevitable, la entrega exactamente una vez a través de una base de datos y un corredor de mensajes sin un coordinador de transacciones distribuidas requiere este compromiso.

Los consumidores deben ser idempotentes. El enfoque más simple es rastrear los IDs de eventos procesados en una tabla processed_events:

CREATE TABLE processed_events (
    event_id   UUID PRIMARY KEY,
    processed_at TIMESTAMPTZ NOT NULL DEFAULT NOW()
);
func (h *OrderHandler) HandleOrderCreated(ctx context.Context, eventID string, payload []byte) error {
    // Deduplicar usando el ID del evento como clave natural
    _, err := h.db.ExecContext(ctx, `
        INSERT INTO processed_events (event_id) VALUES ($1)
        ON CONFLICT (event_id) DO NOTHING
    `, eventID)
    if err != nil {
        return fmt.Errorf("dedup check: %w", err)
    }

    // Verificar si la inserción realmente ocurrió (1 fila) o fue un no-op (0 filas)
    // Un enfoque más simple: usar RETURNING o verificar filas afectadas
    // Si 0 filas afectadas, esto es un duplicado, saltarlo
    ...
}

En la práctica, muchos equipos confían en los propios encabezados de deduplicación del corredor (como el campo key de Kafka para temas con compactación de logs, o el encabezado message-id de RabbitMQ) y tratan la deduplicación a nivel de base de datos como un respaldo. Ambas son capas válidas para aplicar.

Incluye el id del evento del buzón (un UUID) en el mensaje publicado como clave de deduplicación. Los consumidores pueden luego usarlo independientemente del mecanismo de deduplicación que prefieran.

Política de reintento y mensajes envenenados

La columna attempts impulsa la política de reintento. El relé omite filas donde attempts >= maxAttempts y trata esas filas como cartas muertas. Un proceso separado o alerta de operador las maneja.

Una vista simple de cartas muertas:

CREATE VIEW outbox_dead_letters AS
SELECT *
FROM outbox_events
WHERE attempts >= 5
  AND processed_at IS NULL
ORDER BY created_at;

Una buena política de reintento en producción:

  • Establece maxAttempts entre 5 y 10 dependiendo de lo costosos que sean los reintentos.
  • Considera backoff exponencial: incluye una columna retry_after y omite filas donde retry_after > NOW().
  • Alerta cuando COUNT(*) FROM outbox_dead_letters exceda un umbral.
  • Proporciona un camino de reintento manual: un endpoint de administrador o script que restablezca attempts = 0 y retry_after = NULL para filas específicas.

Los mensajes envenenados, filas que fallan consistentemente debido a un error en el consumidor o una discrepancia de esquema, no deben bloquear mensajes saludables. Dado que el relé procesa un lote por tick y marca las fallas con un incremento de intento en lugar de eliminarlas de la cola, las filas saludables proceden normalmente mientras las envenenadas acumulan intentos hasta alcanzar el umbral de cartas muertas. Esta vista outbox_dead_letters es una versión a nivel de base de datos del mismo patrón que las colas de cartas muertas nativas del corredor) implementan — cuarentena después de un umbral, alerta por volumen y requiere una decisión deliberada antes del reprocesamiento.

Ordenación de eventos y particionamiento

La consulta de sondeo ordena por created_at, lo que da ordenación primero en entrar primero en salir dentro de un lote. Para la mayoría de los casos eso es suficiente. Cuando la ordenación por entidad estricta importa, por ejemplo, asegurando que order.updated nunca se publique antes que order.created para el mismo pedido, necesitas ordenación por agregado.

Añade aggregate_id a la cláusula ORDER BY y úsalo como clave de mensaje al publicar en un tema particionado como Apache Kafka). Kafka enruta todos los mensajes con la misma clave a la misma partición, y las particiones se consumen en orden. Esto te da garantías de ordenación por agregado sin ordenación global, lo cual requeriría una única instancia del relé.

ORDER BY aggregate_id, created_at

Para corredores que no soportan ordenación particionada (como colas AMQP básicas), el relé de instancia única o las comprobaciones de ordenación a nivel de aplicación en el consumidor son las alternativas prácticas.

Reduce la latencia de sondeo con LISTEN/NOTIFY

Un intervalo de sondeo de un segundo significa una latencia promedio de evento de 500 milisegundos. Para la mayoría de las cargas de trabajo eso está bien. Para casos donde necesitas latencia cercana a cero, el mecanismo LISTEN/NOTIFY de PostgreSQL permite que el relé se despierte inmediatamente cuando se inserta una nueva fila del buzón.

Añade un disparador a la tabla del buzón:

CREATE OR REPLACE FUNCTION notify_outbox_insert() RETURNS trigger AS $$
BEGIN
    PERFORM pg_notify('outbox_event', NEW.id::text);
    RETURN NEW;
END;
$$ LANGUAGE plpgsql;

CREATE TRIGGER outbox_insert_notify
AFTER INSERT ON outbox_events
FOR EACH ROW EXECUTE FUNCTION notify_outbox_insert();

En el relé, escucha en el canal y despiértate por notificaciones mientras aún caes en el sondeo periódico:

func (r *OutboxRelay) Run(ctx context.Context) error {
    listener := pq.NewListener(r.dsn, 10*time.Second, time.Minute, nil)
    defer listener.Close()

    if err := listener.Listen("outbox_event"); err != nil {
        return fmt.Errorf("listen: %w", err)
    }

    ticker := time.NewTicker(5 * time.Second) // sondeo de respaldo
    defer ticker.Stop()

    for {
        select {
        case <-ctx.Done():
            return ctx.Err()
        case <-listener.Notify:
            if err := r.processBatch(ctx); err != nil {
                r.logger.Error("outbox batch failed (notify)", "err", err)
            }
        case <-ticker.C:
            if err := r.processBatch(ctx); err != nil {
                r.logger.Error("outbox batch failed (poll)", "err", err)
            }
        }
    }
}

El temporizador de respaldo maneja cualquier notificación perdida durante un reinicio del relé o un tropiezo de red. Mantén el intervalo de respaldo en unos pocos segundos en lugar de milisegundos, su trabajo es la recuperación, no la baja latencia.

Observabilidad: métricas, registros y alertas

El buzón es infraestructura. Trátalo como infraestructura e instrumenta según corresponda.

Métricas clave:

var (
    outboxPublished = prometheus.NewCounter(prometheus.CounterOpts{
        Name: "outbox_events_published_total",
        Help: "Total de eventos del buzón publicados con éxito.",
    })
    outboxFailed = prometheus.NewCounterVec(prometheus.CounterOpts{
        Name: "outbox_events_failed_total",
        Help: "Total de fallos de publicación del buzón por tipo de evento.",
    }, []string{"event_type"})
    outboxPending = prometheus.NewGauge(prometheus.GaugeOpts{
        Name: "outbox_events_pending",
        Help: "Número actual de eventos del buzón no procesados.",
    })
    outboxBatchDuration = prometheus.NewHistogram(prometheus.HistogramOpts{
        Name:    "outbox_batch_duration_seconds",
        Help:    "Duración de cada lote de procesamiento del buzón.",
        Buckets: prometheus.DefBuckets,
    })
)

Actualización de la Gauge: ejecuta una consulta periódica para mantener outbox_events_pending precisa:

SELECT COUNT(*) FROM outbox_events WHERE processed_at IS NULL;

Umbral de alertas a considerar:

  • outbox_events_pending > 1000 durante más de dos minutos: el relé se está quedando atrás o está atascado.
  • outbox_events_pending creciendo monotonamente: el corredor está caído o el relé se ha bloqueado.
  • Conteo de cartas muertas no cero: error de esquema o consumidor que necesita investigación.
  • outbox_batch_duration_seconds p95 > 5s: la base de datos es lenta o el tamaño del lote es demasiado grande.

Campos de registro estructurados: incluye event_id, event_type, aggregate_id y attempt en cada línea de registro del relé. Estos campos te permiten correlacionar una publicación fallida con la fila específica del buzón y el rastro del consumidor descendente.

Buzón vs. cola directa vs. saga

El patrón del buzón no es la herramienta adecuada para cada problema de coordinación. Aquí está la comparación:

Enfoque Atomicidad Complejidad Cuándo usar
Publicación directa Ninguna Baja Aceptable perder eventos ocasionalmente
Buzón transaccional Fuerte Media Entrega de eventos confiable desde un único servicio
Patrón Saga Final Alta Transacciones multi-servicio que abarcan múltiples bases de datos
Commit de dos fases Fuerte Muy alta Raramente práctico; evitado en la mayoría de sistemas distribuidos

El patrón del buzón garantiza que un único servicio emite de forma confiable eventos que reflejan sus propios cambios de estado. No coordina cambios de estado a través de múltiples servicios, eso es para lo que está el Patrón Saga). La elección del corredor, ya sea RabbitMQ, SQS, o Kafka, es independiente del patrón del buzón en sí; el relé publica a cualquier corredor que use tu sistema.

Si estás construyendo una saga, el patrón del buzón sigue siendo útil: cada participante en la saga escribe su cambio de estado local y su evento de saga en una transacción usando el buzón, luego el orquestador o coreografía de la saga lee esos eventos de forma confiable.

CDC basado en WAL como relé alternativo

En lugar de sondear, puedes seguir el Write-Ahead Log (WAL) de PostgreSQL y leer las inserciones del buzón directamente desde el stream de replicación. Herramientas como Debezium hacen esto. Las ventajas son menor latencia y ninguna presión de bloqueo en la tabla del buzón. Las desventajas son la complejidad operativa, un slot de replicación dedicado de PostgreSQL y un servicio externo para ejecutar y monitorizar.

Para la mayoría de los equipos, el relé de sondeo descrito anteriormente es el punto de partida adecuado. Seguir el WAL tiene sentido cuando tienes altas tasas de inserción del buzón (decenas de miles por segundo), necesitas latencia de evento sub-100ms o ya estás ejecutando Debezium para otras necesidades de captura de cambios.

Integración con sqlc

Si usas sqlc para código de base de datos Go type-safe, las consultas del buzón encajan naturalmente:

-- name: InsertOutboxEvent :exec
INSERT INTO outbox_events (aggregate_type, aggregate_id, event_type, payload)
VALUES (@aggregate_type, @aggregate_id, @event_type, @payload);

-- name: FetchOutboxBatch :many
SELECT id, aggregate_type, aggregate_id, event_type, payload
FROM outbox_events
WHERE processed_at IS NULL
  AND attempts < @max_attempts
ORDER BY created_at
LIMIT @batch_size
FOR UPDATE SKIP LOCKED;

-- name: MarkOutboxProcessed :exec
UPDATE outbox_events SET processed_at = NOW() WHERE id = @id;

-- name: IncrementOutboxAttempts :exec
UPDATE outbox_events SET attempts = attempts + 1 WHERE id = @id;

-- name: OutboxPendingCount :one
SELECT COUNT(*) FROM outbox_events WHERE processed_at IS NULL;

sqlc genera funciones type-safe para cada consulta, lo que evita errores de interpolación de cadenas y mantiene la lógica de consulta del buzón co-ubicada con el resto de tu capa de acceso a datos.

Lista de verificación de producción

Usa esto antes de enviar una implementación del buzón:

Base de datos

  • La tabla del buzón tiene el índice parcial en created_at WHERE processed_at IS NULL
  • Columna attempts presente con un valor predeterminado de 0
  • Vista o consulta de cartas muertas definida
  • Las filas procesadas antiguas se archivan o eliminan periódicamente (un trabajo de limpieza nocturno es suficiente)

Relé

  • FOR UPDATE SKIP LOCKED usado en la consulta de sondeo
  • El relé se ejecuta dentro de una transacción (comienza antes de la consulta, confirma después de todas las actualizaciones)
  • El tamaño del lote está acotado (50-200 filas es típico)
  • El relé respeta la cancelación del contexto para un apagado graceful
  • Las publicaciones fallidas incrementan attempts en lugar de causar que el lote se aborte

Idempotencia

  • El mensaje publicado incluye el id del buzón como clave de deduplicación
  • Los consumidores son idempotentes o el corredor proporciona deduplicación
  • Ver Idempotencia en Sistemas Distribuidos) para patrones de deduplicación

Observabilidad

  • La gauge outbox_events_pending es monitorizada y alertada
  • El conteo de cartas muertas es alertado
  • La duración del lote del relé es rastreada
  • Los registros estructurados incluyen event_id, event_type y aggregate_id

Operaciones

  • Existe un camino de reintento manual para filas de cartas muertas
  • El comportamiento de reinicio del relé es probado (¿vuelve a publicar correctamente?)
  • El comportamiento ante la caída del corredor es probado (¿el buzón crece y se drena correctamente?)

Reflexiones finales

El problema de la doble escritura es fácil de descartar como un caso edge hasta que cause un incidente. El patrón del buzón transaccional lo soluciona con herramientas que ya tienes: una transacción de PostgreSQL, una goroutine en segundo plano y una tabla extra. El relé es simple de construir, simple de operar y simple de razonar.

El costo es que los consumidores deben estar diseñados para entrega al menos una vez. Ese es un compromiso razonable. La entrega exactamente una vez a través de una base de datos y un corredor sin transacciones distribuidas no es alcanzable en la práctica, y fingir lo contrario lleva a sistemas que pierden silenciosamente o reprocesan eventos en condiciones de fallo.

Escribe el evento con los datos. Retransmítelo de forma confiable. Haz que los consumidores sean idempotentes. Ese es todo el patrón.

Este artículo es parte del clúster App Architecture in Production).

Fuentes

Suscribirse

Recibe nuevas publicaciones sobre sistemas, infraestructura e ingeniería de IA.