Patrones de gestión de estado en sistemas distribuidos: Saga, Outbox y Event Sourcing con PostgreSQL
Introducción: el problema de las transacciones distribuidas
Cuando un monolito se divide en microservicios, lo primero que se sacrifica es la transaccionalidad. En el mundo ACID clásico, una sola base de datos garantiza la atomicidad: todo o nada. En un sistema distribuido, la misma operación —por ejemplo, registrar un pedido— afecta al servicio de pedidos, al de pagos, al de inventario y al de notificaciones. Cada uno tiene su propia base de datos y no existe una transacción única que los abarque a todos.
El protocolo de confirmación en dos fases (2PC) resuelve este problema en teoría, pero en la práctica crea un cuello de botella: el coordinador puede caerse entre una fase y otra, dejando el sistema en un estado indefinido. Además, el 2PC bloquea los recursos de los participantes durante toda la transacción, lo que destruye el rendimiento bajo alta carga.
La comunidad de arquitectura respondió a este desafío con tres patrones que se han convertido en el estándar de facto entre 2024 y 2026: Saga, Transactional Outbox y Event Sourcing. En este artículo analizaremos cada uno en detalle y mostraremos una implementación práctica en Go con PostgreSQL y Redis.
Patrón Saga: gestión de transacciones de larga duración
Una Saga es una secuencia de transacciones locales, cada una de las cuales publica un evento o mensaje que desencadena el siguiente paso. Si un paso falla, la Saga ejecuta transacciones compensatorias para todos los pasos anteriores, devolviendo el sistema a un estado consistente.
Coreografía vs. Orquestación
Existen dos variantes de implementación de Saga con arquitecturas radicalmente distintas:
- Coreografía (Choreography): cada servicio reacciona a eventos y publica los suyos propios. No hay un coordinador central. El servicio de pedidos publica
OrderCreated, el servicio de pagos se suscribe a ese evento, reserva los fondos y publicaPaymentReserved, el servicio de inventario se suscribe aPaymentReserved, y así sucesivamente. - Orquestación (Orchestration): un orquestador central (Saga Orchestrator) coordina explícitamente cada servicio: «bloquea el pago», «reserva el artículo», «envía la notificación». El orquestador hace seguimiento del estado de todo el proceso.
La coreografía es más sencilla al principio —no hay un punto único de fallo—, pero a medida que crece el número de servicios, el grafo de dependencias se vuelve difícil de depurar. La orquestación ofrece control centralizado y observabilidad, pero el orquestador se convierte en el punto donde se concentra la lógica de negocio y en un posible cuello de botella.
Transacciones compensatorias
El requisito clave de Saga: cada transacción local debe tener su correspondiente transacción compensatoria. Si el paso «cobrar el dinero» tuvo éxito pero el siguiente paso «reservar el artículo» falló, el sistema debe ejecutar «devolver el dinero». Las compensaciones deben ser idempotentes: una llamada repetida no debe generar un doble reembolso.
Ventajas de Saga: sin bloqueos distribuidos, alta disponibilidad, cada servicio es independiente. Desventajas: inconsistencia temporal (eventual consistency), complejidad en el diseño de las compensaciones, dificultades de depuración en la coreografía.
Transactional Outbox Pattern: entrega garantizada de eventos
Saga y Event Sourcing presuponen una publicación fiable de eventos. Pero ¿cómo garantizar que un evento sea publicado incluso si el servicio cae justo después de hacer commit de la transacción? Aquí es donde entra en escena el Transactional Outbox.
La idea es simple y elegante: en lugar de publicar el evento directamente en el broker de mensajes, el servicio escribe el evento en una tabla especial outbox dentro de la misma transacción que los datos de negocio. Un proceso de fondo separado (Message Relay o Polling Worker) lee los registros no procesados de la tabla outbox y los publica en el broker. La atomicidad está garantizada por la base de datos: o se crean tanto el registro de negocio como el de outbox, o no se crea ninguno.
Esquema SQL para Outbox en PostgreSQL
CREATE TABLE outbox_events (\n id UUID PRIMARY KEY DEFAULT gen_random_uuid(),\n aggregate_type VARCHAR(100) NOT NULL,\n aggregate_id VARCHAR(100) NOT NULL,\n event_type VARCHAR(100) NOT NULL,\n payload JSONB NOT NULL,\n created_at TIMESTAMPTZ NOT NULL DEFAULT NOW(),\n published_at TIMESTAMPTZ,\n retry_count INT NOT NULL DEFAULT 0,\n status VARCHAR(20) NOT NULL DEFAULT 'PENDING'\n CHECK (status IN ('PENDING', 'PUBLISHED', 'FAILED'))\n);\n\nCREATE INDEX idx_outbox_status_created\n ON outbox_events (status, created_at)\n WHERE status = 'PENDING';\nEl índice con condición parcial (WHERE status = 'PENDING') es crítico para el rendimiento: el polling worker accederá a este índice en cada iteración, y sin él, un escaneo completo de la tabla se convertirá en un problema a partir de varios millones de registros.
Event Sourcing: el estado como secuencia de eventos
Event Sourcing invierte el modelo tradicional de almacenamiento de datos. En lugar de guardar el estado actual de una entidad (actualizando una fila en la tabla), almacenamos la secuencia de eventos que llevaron a ese estado. El estado actual se reconstruye reproduciendo (replay) todos los eventos.
Por ejemplo, para una cuenta bancaria no se almacena el campo balance = 1500. En su lugar se almacenan los eventos: AccountOpened(0), MoneyDeposited(2000), MoneyWithdrawn(500). El saldo de 1500 se calcula en el momento de la lectura.
Almacenamiento de eventos en PostgreSQL
CREATE TABLE event_store (\n id BIGSERIAL PRIMARY KEY,\n stream_id VARCHAR(200) NOT NULL,\n stream_version BIGINT NOT NULL,\n event_type VARCHAR(100) NOT NULL,\n event_data JSONB NOT NULL,\n metadata JSONB NOT NULL DEFAULT '{}',\n created_at TIMESTAMPTZ NOT NULL DEFAULT NOW(),\n UNIQUE (stream_id, stream_version)\n);\n\nCREATE INDEX idx_event_store_stream\n ON event_store (stream_id, stream_version);\nLa restricción UNIQUE (stream_id, stream_version) garantiza el bloqueo optimista: si dos procesos intentan escribir un evento con la misma versión para el mismo stream, PostgreSQL rechazará una de las operaciones, previniendo el conflicto.
Snapshots: con streams muy grandes, el replay completo se vuelve costoso. La solución es guardar periódicamente un snapshot del estado actual y reproducir solo los eventos ocurridos después del snapshot.
Implementación práctica del Outbox en Go con PostgreSQL
Veamos una implementación completa del Polling Worker para el Transactional Outbox en Go. Usamos pgx para trabajar con PostgreSQL y patrones estándar de concurrencia.
Estructuras de datos
package outbox\n\nimport (\n "context"\n "time"\n "github.com/google/uuid"\n)\n\ntype Event struct {\n ID uuid.UUID\n AggregateType string\n AggregateID string\n EventType string\n Payload []byte\n CreatedAt time.Time\n RetryCount int\n}\n\ntype Publisher interface {\n Publish(ctx context.Context, event Event) error\n}\nPolling Worker
package outbox\n\nimport (\n "context"\n "fmt"\n "log/slog"\n "time"\n\n "github.com/jackc/pgx/v5/pgxpool"\n)\n\ntype Worker struct {\n db *pgxpool.Pool\n publisher Publisher\n batchSize int\n interval time.Duration\n}\n\nfunc NewWorker(db *pgxpool.Pool, pub Publisher) *Worker {\n return &Worker{\n db: db,\n publisher: pub,\n batchSize: 100,\n interval: 500 * time.Millisecond,\n }\n}\n\nfunc (w *Worker) Run(ctx context.Context) error {\n ticker := time.NewTicker(w.interval)\n defer ticker.Stop()\n\n for {\n select {\n case <-ctx.Done():\n return ctx.Err()\n case <-ticker.C:\n if err := w.processBatch(ctx); err != nil {\n slog.Error("outbox: batch processing failed", "error", err)\n }\n }\n }\n}\n\nfunc (w *Worker) processBatch(ctx context.Context) error {\n tx, err := w.db.Begin(ctx)\n if err != nil {\n return fmt.Errorf("begin tx: %w", err)\n }\n defer tx.Rollback(ctx)\n\n // SELECT FOR UPDATE SKIP LOCKED — clave para el escalado horizontal\n rows, err := tx.Query(ctx, `\n SELECT id, aggregate_type, aggregate_id, event_type, payload, created_at, retry_count\n FROM outbox_events\n WHERE status = 'PENDING'\n ORDER BY created_at\n LIMIT $1\n FOR UPDATE SKIP LOCKED\n `, w.batchSize)\n if err != nil {\n return fmt.Errorf("query events: %w", err)\n }\n\n var events []Event\n for rows.Next() {\n var e Event\n if err := rows.Scan(\n &e.ID, &e.AggregateType, &e.AggregateID,\n &e.EventType, &e.Payload, &e.CreatedAt, &e.RetryCount,\n ); err != nil {\n return fmt.Errorf("scan event: %w", err)\n }\n events = append(events, e)\n }\n rows.Close()\n\n for _, event := range events {\n if err := w.publisher.Publish(ctx, event); err != nil {\n slog.Warn("outbox: publish failed", "event_id", event.ID, "error", err)\n _, _ = tx.Exec(ctx,\n `UPDATE outbox_events SET retry_count = retry_count + 1,\n status = CASE WHEN retry_count >= 4 THEN 'FAILED' ELSE 'PENDING' END\n WHERE id = $1`, event.ID)\n continue\n }\n _, _ = tx.Exec(ctx,\n `UPDATE outbox_events SET status = 'PUBLISHED', published_at = NOW() WHERE id = $1`,\n event.ID)\n }\n\n return tx.Commit(ctx)\n}\nLa directiva FOR UPDATE SKIP LOCKED en PostgreSQL permite ejecutar varias instancias del polling worker en paralelo: cada una capturará su propio conjunto de filas sin competir con las demás. Esto garantiza el escalado horizontal del procesamiento del Outbox.
Escritura en el Outbox dentro de una transacción de negocio
func (s *OrderService) CreateOrder(ctx context.Context, req CreateOrderRequest) error {\n tx, err := s.db.Begin(ctx)\n if err != nil {\n return err\n }\n defer tx.Rollback(ctx)\n\n // 1. Creamos el pedido\n orderID := uuid.New()\n _, err = tx.Exec(ctx,\n `INSERT INTO orders (id, user_id, total) VALUES ($1, $2, $3)`,\n orderID, req.UserID, req.Total)\n if err != nil {\n return fmt.Errorf("insert order: %w", err)\n }\n\n // 2. En la misma transacción, escribimos el evento en el Outbox\n payload, _ := json.Marshal(map[string]any{\n "order_id": orderID,\n "user_id": req.UserID,\n "total": req.Total,\n })\n _, err = tx.Exec(ctx, `\n INSERT INTO outbox_events (aggregate_type, aggregate_id, event_type, payload)\n VALUES ('Order', $1, 'OrderCreated', $2)\n `, orderID.String(), payload)\n if err != nil {\n return fmt.Errorf("insert outbox: %w", err)\n }\n\n return tx.Commit(ctx)\n}\nIntegración con Redis para buffer y deduplicación
Redis complementa de forma natural el Outbox Pattern en dos escenarios: el buffer de eventos de alta frecuencia y la deduplicación en el lado del consumidor.
Deduplicación con Redis
Incluso usando Outbox, los eventos pueden llegar al consumidor más de una vez (entrega at-least-once). En el lado del consumidor, Redis proporciona una deduplicación eficiente mediante SET NX EX:
func (c *Consumer) isProcessed(ctx context.Context, eventID string) (bool, error) {\n key := "processed_event:" + eventID\n // SET key 1 NX EX 86400 — establecer si no existe, TTL 24 horas\n set, err := c.redis.SetNX(ctx, key, 1, 24*time.Hour).Result()\n if err != nil {\n return false, err\n }\n // set=true significa que la clave fue creada — el evento se procesa por primera vez\n return !set, nil\n}\n\nfunc (c *Consumer) Handle(ctx context.Context, event Event) error {\n already, err := c.isProcessed(ctx, event.ID.String())\n if err != nil || already {\n return err // ignoramos el duplicado\n }\n return c.processEvent(ctx, event)\n}\nRedis Streams como buffer intermedio
Durante picos de carga, el polling worker puede no ser capaz de leer de PostgreSQL con suficiente rapidez. Redis Streams (XADD/XREADGROUP) actúan como buffer entre el polling worker y el broker final (Kafka, RabbitMQ). El worker publica en el Redis Stream de forma atómica, y un consumer group separado reenvía los eventos a Kafka. Esto reduce la latencia y protege a Kafka de los picos de carga.
Monitorización y depuración de transacciones distribuidas
El tracing distribuido es una herramienta imprescindible para sistemas que utilizan Saga y Outbox. El Trace ID debe propagarse a través de todos los eventos: almacenarse en el campo metadata de la tabla outbox_events y ser leído por los consumidores para restaurar el contexto. Utiliza OpenTelemetry para instrumentar los servicios en Go.
Métricas clave para monitorizar el Outbox:
- outbox_pending_count — número de eventos no publicados (indicador crítico de SLO);
- outbox_publish_latency_seconds — latencia entre la creación del evento y su publicación;
- outbox_failed_count — eventos en estado FAILED que requieren intervención manual;
- saga_compensation_total — número de transacciones compensatorias ejecutadas.
Para diagnosticar Sagas bloqueadas, añade una tabla saga_state con el campo last_updated_at y configura una alerta cuando no haya actualizaciones durante más de N minutos. Esto permitirá detectar orquestaciones «atascadas» antes de que comiencen a afectar a los usuarios.
PostgreSQL ofrece herramientas de depuración potentes: pg_stat_activity mostrará los bloqueos activos, y EXPLAIN ANALYZE mostrará el plan de consulta del polling worker. Vigila la métrica de autovacuum en la tabla outbox_events: con una alta frecuencia de actualizaciones de filas (transición PENDING → PUBLISHED), sin un vacuum oportuno la tabla crecerá sin control.
Comparación de patrones: cuándo usar cada uno
La elección entre Saga, Outbox y Event Sourcing no es excluyente —en la práctica suelen usarse juntos—. Sin embargo, sus objetivos y compromisos difieren considerablemente.
- Saga resuelve el problema de coordinación de procesos de negocio que abarcan varios servicios. Aplica Saga en cualquier caso donde haya transacciones de negocio de múltiples pasos con posibilidad de reversión. Prefiere la orquestación cuando la lógica es compleja y se necesita monitorización centralizada; la coreografía es adecuada para procesos lineales simples con pocos participantes.
- Transactional Outbox es un patrón universal para la publicación fiable de eventos. Úsalo siempre que un servicio modifique el estado en la base de datos y necesite notificar a otros servicios. Outbox elimina el problema de la «doble escritura» y es el bloque de construcción de cualquier arquitectura orientada a eventos.
- Event Sourcing está justificado cuando el historial de cambios es un requisito de negocio de primer nivel: sistemas financieros, auditoría, sistemas que necesitan consultas temporales. Event Sourcing incrementa considerablemente la complejidad del sistema: consultas más difíciles, problema de evolución del esquema, necesidad de snapshots. No lo apliques «por defecto» solo porque esté de moda.
La configuración típica de producción para un sistema de microservicios en 2026: Saga Orchestration para procesos de negocio + Transactional Outbox en cada servicio + Event Sourcing para agregados con un rico historial de cambios + Redis para deduplicación y buffer. PostgreSQL con las extensiones pgcrypto y pg_partman (particionamiento de la tabla event_store por tiempo) es una base sólida para los tres patrones sin necesidad de introducir almacenes de eventos especializados en las etapas tempranas.
Conclusión
La gestión del estado en sistemas distribuidos es uno de los problemas fundamentales de la arquitectura backend. Saga, Outbox y Event Sourcing no son balas de plata, pero combinados correctamente garantizan la consistencia de los datos, la entrega fiable de eventos y un historial completo de cambios, manteniendo la independencia de los servicios. PostgreSQL como base sólida, Go como implementación de alto rendimiento y Redis como buffer rápido forman un stack capaz de dar servicio a sistemas de alto tráfico con un comportamiento predecible bajo carga y una depuración transparente.
Tecnologías
Etiquetas
Ruslan Ismailov
Desarrollador Senior Web / Backend. Desarrollador senior web/backend con 9 años de experiencia. Stack: PHP, Laravel, PostgreSQL, Redis, Docker, Kubernetes, REST, microservicios, CI/CD. Más sobre mí →