Arquitectura

Implementación del patrón Saga en microservicios con Go: orquestación y coreografía de transacciones con PostgreSQL y Redis

Ruslan Ismailov Publicado 18 min de lectura
I

Introducción: por qué las transacciones distribuidas son complejas

En una aplicación monolítica, una transacción es una operación atómica dentro de una única base de datos. Las garantías ACID de PostgreSQL la hacen confiable y predecible. Pero cuando el sistema se divide en microservicios, cada uno con su propia BD, las transacciones clásicas se vuelven imposibles: no es posible ejecutar BEGIN ... COMMIT simultáneamente en dos almacenes de datos aislados.

El commit en dos fases (2PC) resuelve este problema formalmente, pero en la práctica genera un acoplamiento rígido entre servicios, bloqueos de espera y un único punto de fallo en el coordinador. Para sistemas de alta carga, esto es inaceptable.

El patrón Saga es un enfoque alternativo: en lugar de una única transacción distribuida, se ejecuta una secuencia de transacciones locales en cada servicio. Si en algún paso ocurre un fallo, se lanzan transacciones compensatorias para revertir los pasos ya completados. Esto no es un rollback en el sentido tradicional: son acciones inversas explícitas implementadas en el código.

Dos enfoques: orquestación vs coreografía

El patrón Saga se implementa de dos formas fundamentalmente distintas. La elección entre ellas determina la arquitectura de todo el sistema.

Coreografía (Choreography)

Cada servicio publica eventos tras completar su transacción local. Los demás servicios se suscriben a estos eventos y reaccionan ante ellos. No existe un coordinador central: cada servicio sabe qué hacer en respuesta a un evento concreto.

  • Ventajas: descentralización total, bajo acoplamiento, facilidad para agregar nuevos participantes.
  • Desventajas: difícil rastrear el estado completo de la saga, la lógica queda dispersa entre servicios y la depuración se convierte en una investigación detectivesca.

Orquestación (Orchestration)

Un orquestador central (Saga Orchestrator) gestiona todos los pasos: invoca los servicios, espera respuestas y decide cuándo compensar. Almacena el estado de la saga, normalmente en una base de datos.

  • Ventajas: visibilidad completa del estado de la saga, lógica centralizada, orden de pasos claro y depuración sencilla.
  • Desventajas: el orquestador es un potencial punto de fallo, mayor acoplamiento con el coordinador y riesgo de que el orquestador se convierta en un servicio "gordo".

Comparación de enfoques

  • Estado de la saga: en coreografía, distribuido entre eventos y servicios; en orquestación, centralizado en la BD del orquestador.
  • Acoplamiento: la coreografía ofrece bajo acoplamiento entre servicios; la orquestación crea dependencia de todos los servicios hacia el orquestador.
  • Observabilidad: la coreografía requiere agregación de logs y trazas; la orquestación proporciona el estado completo desde un único lugar.
  • Cuándo elegir coreografía: sagas simples (2–3 pasos), el equipo ya trabaja con arquitectura event-driven, los servicios están débilmente acoplados.
  • Cuándo elegir orquestación: procesos de negocio complejos (5+ pasos), se requiere monitoreo explícito del estado, la compensación centralizada es importante.

Implementación de una Saga coreografiada en Go con Redis Streams

Para la coreografía, Redis Streams es una excelente opción: proporciona una cola de mensajes persistente con grupos de consumidores. A diferencia del simple Pub/Sub, los mensajes se almacenan y pueden releerse ante un fallo.

Estructura de eventos

// events.go
package saga

import "time"

type EventType string

const (
    OrderCreated     EventType = "order.created"
    PaymentProcessed EventType = "payment.processed"
    PaymentFailed    EventType = "payment.failed"
    InventoryReserved EventType = "inventory.reserved"
    InventoryFailed  EventType = "inventory.failed"
    OrderCompleted   EventType = "order.completed"
    OrderCancelled   EventType = "order.cancelled"
)

type SagaEvent struct {
    SagaID    string            `json:"saga_id"`
    EventType EventType         `json:"event_type"`
    Payload   map[string]string `json:"payload"`
    Timestamp time.Time         `json:"timestamp"`
}

Publicación de eventos en Redis Streams

// publisher.go
package saga

import (
    "context"
    "encoding/json"
    "fmt"
    "time"

    "github.com/redis/go-redis/v9"
)

type EventPublisher struct {
    rdb    *redis.Client
    stream string
}

func NewEventPublisher(rdb *redis.Client, stream string) *EventPublisher {
    return &EventPublisher{rdb: rdb, stream: stream}
}

func (p *EventPublisher) Publish(ctx context.Context, event SagaEvent) error {
    event.Timestamp = time.Now()
    data, err := json.Marshal(event)
    if err != nil {
        return fmt.Errorf("marshal event: %w", err)
    }

    return p.rdb.XAdd(ctx, &redis.XAddArgs{
        Stream: p.stream,
        Values: map[string]interface{}{
            "saga_id":    event.SagaID,
            "event_type": string(event.EventType),
            "data":       string(data),
        },
    }).Err()
}

Consumidor de eventos con Consumer Group

// consumer.go
package saga

import (
    "context"
    "encoding/json"
    "fmt"
    "log"
    "time"

    "github.com/redis/go-redis/v9"
)

type EventHandler func(ctx context.Context, event SagaEvent) error

type EventConsumer struct {
    rdb       *redis.Client
    stream    string
    group     string
    consumer  string
    handlers  map[EventType]EventHandler
}

func NewEventConsumer(rdb *redis.Client, stream, group, consumer string) *EventConsumer {
    return &EventConsumer{
        rdb:      rdb,
        stream:   stream,
        group:    group,
        consumer: consumer,
        handlers: make(map[EventType]EventHandler),
    }
}

func (c *EventConsumer) Register(eventType EventType, handler EventHandler) {
    c.handlers[eventType] = handler
}

func (c *EventConsumer) Start(ctx context.Context) error {
    // Creamos el grupo de consumidores si no existe
    c.rdb.XGroupCreateMkStream(ctx, c.stream, c.group, "0")

    for {
        select {
        case <-ctx.Done():
            return ctx.Err()
        default:
        }

        streams, err := c.rdb.XReadGroup(ctx, &redis.XReadGroupArgs{
            Group:    c.group,
            Consumer: c.consumer,
            Streams:  []string{c.stream, ">"},
            Count:    10,
            Block:    5 * time.Second,
        }).Result()

        if err != nil && err != redis.Nil {
            log.Printf("xreadgroup error: %v", err)
            continue
        }

        for _, stream := range streams {
            for _, msg := range stream.Messages {
                if err := c.processMessage(ctx, msg); err != nil {
                    log.Printf("process message %s error: %v", msg.ID, err)
                    continue // el mensaje permanece en PEL para reintento
                }
                // ACK solo tras procesamiento exitoso
                c.rdb.XAck(ctx, c.stream, c.group, msg.ID)
            }
        }
    }
}

func (c *EventConsumer) processMessage(ctx context.Context, msg redis.XMessage) error {
    dataStr, ok := msg.Values["data"].(string)
    if !ok {
        return fmt.Errorf("invalid message format")
    }

    var event SagaEvent
    if err := json.Unmarshal([]byte(dataStr), &event); err != nil {
        return fmt.Errorf("unmarshal event: %w", err)
    }

    handler, ok := c.handlers[event.EventType]
    if !ok {
        return nil // ignoramos eventos desconocidos
    }

    return handler(ctx, event)
}

Ejemplo: servicio de pago en una saga coreografiada

// payment_service.go
package payment

import (
    "context"
    "fmt"
    "log"

    "github.com/yourorg/saga"
)

type PaymentService struct {
    publisher *saga.EventPublisher
    repo      *PaymentRepository
}

func (s *PaymentService) HandleOrderCreated(ctx context.Context, event saga.SagaEvent) error {
    orderID := event.Payload["order_id"]
    amount := event.Payload["amount"]

    // Verificación de idempotencia — no procesar dos veces
    if exists, _ := s.repo.ExistsPayment(ctx, event.SagaID); exists {
        log.Printf("payment for saga %s already processed, skipping", event.SagaID)
        return nil
    }

    err := s.repo.ProcessPayment(ctx, event.SagaID, orderID, amount)
    if err != nil {
        // Publicamos el evento de fallo — otros servicios iniciarán la compensación
        return s.publisher.Publish(ctx, saga.SagaEvent{
            SagaID:    event.SagaID,
            EventType: saga.PaymentFailed,
            Payload:   map[string]string{"reason": err.Error()},
        })
    }

    return s.publisher.Publish(ctx, saga.SagaEvent{
        SagaID:    event.SagaID,
        EventType: saga.PaymentProcessed,
        Payload:   map[string]string{"order_id": orderID},
    })
}

func (s *PaymentService) HandleOrderCancelled(ctx context.Context, event saga.SagaEvent) error {
    // Transacción compensatoria: reembolso del pago
    if err := s.repo.RefundPayment(ctx, event.SagaID); err != nil {
        return fmt.Errorf("refund payment: %w", err)
    }
    log.Printf("payment refunded for saga %s", event.SagaID)
    return nil
}

Implementación de una Saga orquestada en Go con PostgreSQL

El orquestador almacena el estado de cada saga en PostgreSQL. Esto permite recuperarse tras un fallo: al reiniciarse, el orquestador lee las sagas incompletas y continúa la ejecución.

Esquema de base de datos

-- migrations/001_saga_state.sql
CREATE TYPE saga_status AS ENUM (
    'started',
    'processing',
    'completed',
    'compensating',
    'failed'
);

CREATE TABLE saga_instances (
    id          UUID PRIMARY KEY DEFAULT gen_random_uuid(),
    saga_type   VARCHAR(100) NOT NULL,
    status      saga_status NOT NULL DEFAULT 'started',
    current_step VARCHAR(100),
    payload     JSONB NOT NULL DEFAULT '{}',
    created_at  TIMESTAMPTZ NOT NULL DEFAULT NOW(),
    updated_at  TIMESTAMPTZ NOT NULL DEFAULT NOW()
);

CREATE TABLE saga_steps (
    id          UUID PRIMARY KEY DEFAULT gen_random_uuid(),
    saga_id     UUID NOT NULL REFERENCES saga_instances(id),
    step_name   VARCHAR(100) NOT NULL,
    status      VARCHAR(50) NOT NULL DEFAULT 'pending',
    result      JSONB,
    executed_at TIMESTAMPTZ,
    UNIQUE(saga_id, step_name)
);

CREATE INDEX idx_saga_instances_status ON saga_instances(status);

Modelo de estado y repositorio

// orchestrator/state.go
package orchestrator

import (
    "context"
    "database/sql"
    "encoding/json"
    "fmt"
    "time"

    "github.com/google/uuid"
)

type SagaStatus string

const (
    StatusStarted      SagaStatus = "started"
    StatusProcessing   SagaStatus = "processing"
    StatusCompleted    SagaStatus = "completed"
    StatusCompensating SagaStatus = "compensating"
    StatusFailed       SagaStatus = "failed"
)

type SagaInstance struct {
    ID          uuid.UUID      `db:"id"`
    SagaType    string         `db:"saga_type"`
    Status      SagaStatus     `db:"status"`
    CurrentStep string         `db:"current_step"`
    Payload     map[string]any `db:"payload"`
    CreatedAt   time.Time      `db:"created_at"`
    UpdatedAt   time.Time      `db:"updated_at"`
}

type SagaRepository struct {
    db *sql.DB
}

func (r *SagaRepository) Create(ctx context.Context, sagaType string, payload map[string]any) (*SagaInstance, error) {
    payloadJSON, err := json.Marshal(payload)
    if err != nil {
        return nil, err
    }

    instance := &SagaInstance{
        ID:       uuid.New(),
        SagaType: sagaType,
        Status:   StatusStarted,
        Payload:  payload,
    }

    _, err = r.db.ExecContext(ctx,
        `INSERT INTO saga_instances (id, saga_type, status, payload)
         VALUES ($1, $2, $3, $4)`,
        instance.ID, instance.SagaType, instance.Status, payloadJSON,
    )
    return instance, err
}

func (r *SagaRepository) UpdateStatus(ctx context.Context, sagaID uuid.UUID, status SagaStatus, step string) error {
    _, err := r.db.ExecContext(ctx,
        `UPDATE saga_instances
         SET status = $1, current_step = $2, updated_at = NOW()
         WHERE id = $3`,
        status, step, sagaID,
    )
    return err
}

func (r *SagaRepository) MarkStepDone(ctx context.Context, sagaID uuid.UUID, stepName string, result map[string]any) error {
    resultJSON, _ := json.Marshal(result)
    _, err := r.db.ExecContext(ctx,
        `INSERT INTO saga_steps (saga_id, step_name, status, result, executed_at)
         VALUES ($1, $2, 'completed', $3, NOW())
         ON CONFLICT (saga_id, step_name) DO UPDATE
         SET status = 'completed', result = $3, executed_at = NOW()`,
        sagaID, stepName, resultJSON,
    )
    return err
}

func (r *SagaRepository) IsStepDone(ctx context.Context, sagaID uuid.UUID, stepName string) (bool, error) {
    var count int
    err := r.db.QueryRowContext(ctx,
        `SELECT COUNT(*) FROM saga_steps
         WHERE saga_id = $1 AND step_name = $2 AND status = 'completed'`,
        sagaID, stepName,
    ).Scan(&count)
    return count > 0, err
}

func (r *SagaRepository) GetPendingSagas(ctx context.Context) ([]*SagaInstance, error) {
    rows, err := r.db.QueryContext(ctx,
        `SELECT id, saga_type, status, current_step, payload, created_at, updated_at
         FROM saga_instances
         WHERE status IN ('started', 'processing', 'compensating')
         AND updated_at < NOW() - INTERVAL '5 minutes'`)
    if err != nil {
        return nil, fmt.Errorf("query pending sagas: %w", err)
    }
    defer rows.Close()

    var instances []*SagaInstance
    for rows.Next() {
        var inst SagaInstance
        var payloadJSON []byte
        if err := rows.Scan(&inst.ID, &inst.SagaType, &inst.Status,
            &inst.CurrentStep, &payloadJSON, &inst.CreatedAt, &inst.UpdatedAt); err != nil {
            return nil, err
        }
        json.Unmarshal(payloadJSON, &inst.Payload)
        instances = append(instances, &inst)
    }
    return instances, rows.Err()
}

Orquestador central

// orchestrator/orchestrator.go
package orchestrator

import (
    "context"
    "fmt"
    "log"

    "github.com/google/uuid"
)

type Step struct {
    Name      string
    Execute   func(ctx context.Context, payload map[string]any) (map[string]any, error)
    Compensate func(ctx context.Context, payload map[string]any) error
}

type SagaOrchestrator struct {
    repo  *SagaRepository
    steps []Step
}

func NewSagaOrchestrator(repo *SagaRepository, steps []Step) *SagaOrchestrator {
    return &SagaOrchestrator{repo: repo, steps: steps}
}

func (o *SagaOrchestrator) Execute(ctx context.Context, sagaType string, payload map[string]any) error {
    instance, err := o.repo.Create(ctx, sagaType, payload)
    if err != nil {
        return fmt.Errorf("create saga: %w", err)
    }
    return o.run(ctx, instance)
}

func (o *SagaOrchestrator) run(ctx context.Context, instance *SagaInstance) error {
    executedSteps := []int{}

    for i, step := range o.steps {
        // Verificación de idempotencia: ¿ya fue ejecutado?
        done, err := o.repo.IsStepDone(ctx, instance.ID, step.Name)
        if err != nil {
            return fmt.Errorf("check step %s: %w", step.Name, err)
        }
        if done {
            executedSteps = append(executedSteps, i)
            continue
        }

        o.repo.UpdateStatus(ctx, instance.ID, StatusProcessing, step.Name)
        log.Printf("[saga %s] executing step: %s", instance.ID, step.Name)

        result, err := step.Execute(ctx, instance.Payload)
        if err != nil {
            log.Printf("[saga %s] step %s failed: %v — starting compensation", instance.ID, step.Name, err)
            o.repo.UpdateStatus(ctx, instance.ID, StatusCompensating, step.Name)
            return o.compensate(ctx, instance, executedSteps)
        }

        // Guardamos el resultado del paso para garantizar idempotencia
        if err := o.repo.MarkStepDone(ctx, instance.ID, step.Name, result); err != nil {
            return fmt.Errorf("mark step done: %w", err)
        }
        executedSteps = append(executedSteps, i)

        // Actualizamos el payload con los resultados del paso
        for k, v := range result {
            instance.Payload[k] = v
        }
    }

    return o.repo.UpdateStatus(ctx, instance.ID, StatusCompleted, "")
}

func (o *SagaOrchestrator) compensate(ctx context.Context, instance *SagaInstance, executedSteps []int) error {
    // Compensamos en orden inverso
    for i := len(executedSteps) - 1; i >= 0; i-- {
        stepIdx := executedSteps[i]
        step := o.steps[stepIdx]

        if step.Compensate == nil {
            continue
        }

        log.Printf("[saga %s] compensating step: %s", instance.ID, step.Name)
        if err := step.Compensate(ctx, instance.Payload); err != nil {
            // Registramos el error pero continuamos compensando los demás pasos
            log.Printf("[saga %s] compensation of %s failed: %v", instance.ID, step.Name, err)
        }
    }

    return o.repo.UpdateStatus(ctx, instance.ID, StatusFailed, "")
}

Uso del orquestador

// main.go (ejemplo de ensamblaje)
package main

import (
    "context"
    "database/sql"
    "log"

    _ "github.com/lib/pq"
    "github.com/yourorg/orchestrator"
)

func main() {
    db, _ := sql.Open("postgres", "postgres://user:pass@localhost/orders?sslmode=disable")
    repo := &orchestrator.SagaRepository{} // inicialización con db

    orderSaga := orchestrator.NewSagaOrchestrator(repo, []orchestrator.Step{
        {
            Name: "reserve_inventory",
            Execute: func(ctx context.Context, payload map[string]any) (map[string]any, error) {
                // llamada al inventory service
                return map[string]any{"reservation_id": "res-123"}, nil
            },
            Compensate: func(ctx context.Context, payload map[string]any) error {
                // cancelación de la reserva
                return nil
            },
        },
        {
            Name: "process_payment",
            Execute: func(ctx context.Context, payload map[string]any) (map[string]any, error) {
                // llamada al payment service
                return map[string]any{"payment_id": "pay-456"}, nil
            },
            Compensate: func(ctx context.Context, payload map[string]any) error {
                // reembolso del pago
                return nil
            },
        },
        {
            Name: "confirm_order",
            Execute: func(ctx context.Context, payload map[string]any) (map[string]any, error) {
                return map[string]any{"status": "confirmed"}, nil
            },
            Compensate: nil, // paso final, no requiere compensación
        },
    })

    ctx := context.Background()
    if err := orderSaga.Execute(ctx, "order_saga", map[string]any{
        "order_id": "ord-789",
        "amount":   "150.00",
        "user_id":  "usr-001",
    }); err != nil {
        log.Printf("saga failed: %v", err)
    }
}

Gestión de transacciones compensatorias e idempotencia

Las transacciones compensatorias no son un rollback. Son operaciones de negocio independientes que deben implementarse de forma explícita. Algunas reglas críticas:

La idempotencia es un requisito obligatorio

Cualquier operación en una saga (directa o compensatoria) puede ser invocada más de una vez debido a reintentos tras un fallo de red. Cada manejador debe funcionar correctamente ante llamadas repetidas.

El patrón estándar es almacenar un idempotency_key (normalmente saga_id + step_name) en la tabla de resultados. Antes de ejecutar la operación, se verifica si ya existe un registro con dicha clave.

// idempotency.go
func (r *PaymentRepository) ProcessPaymentIdempotent(
    ctx context.Context,
    idempotencyKey string,
    amount string,
) error {
    tx, err := r.db.BeginTx(ctx, nil)
    if err != nil {
        return err
    }
    defer tx.Rollback()

    // Verificación atómica + inserción mediante INSERT ... ON CONFLICT DO NOTHING
    result, err := tx.ExecContext(ctx,
        `INSERT INTO payment_operations (idempotency_key, amount, created_at)
         VALUES ($1, $2, NOW())
         ON CONFLICT (idempotency_key) DO NOTHING`,
        idempotencyKey, amount,
    )
    if err != nil {
        return fmt.Errorf("insert payment operation: %w", err)
    }

    rows, _ := result.RowsAffected()
    if rows == 0 {
        // La operación ya fue ejecutada — devolvemos éxito
        return tx.Commit()
    }

    // Ejecutamos el pago real
    if err := r.chargeCard(ctx, tx, amount); err != nil {
        return err
    }

    return tx.Commit()
}

Reintentos con backoff exponencial

Ante fallos temporales (red, sobrecarga de BD), use reintentos con jitter para evitar las "manadas de trueno":

// retry.go
package saga

import (
    "context"
    "math"
    "math/rand"
    "time"
)

func WithRetry(ctx context.Context, maxAttempts int, fn func() error) error {
    var lastErr error
    for attempt := 0; attempt < maxAttempts; attempt++ {
        if err := fn(); err == nil {
            return nil
        } else {
            lastErr = err
        }

        if attempt == maxAttempts-1 {
            break
        }

        backoff := time.Duration(math.Pow(2, float64(attempt))) * 100 * time.Millisecond
        jitter := time.Duration(rand.Intn(100)) * time.Millisecond
        select {
        case <-ctx.Done():
            return ctx.Err()
        case <-time.After(backoff + jitter):
        }
    }
    return lastErr
}

Monitoreo y depuración de transacciones Saga

Sin observabilidad, las sagas se convierten en una "caja negra". Para sistemas en producción es imprescindible:

  • Logging estructurado: cada evento de la saga se registra con los campos saga_id, step, status, duration_ms. Use log/slog o zerolog.
  • Métricas: contadores de sagas iniciadas, completadas y compensadas mediante Prometheus. Histogramas independientes para la duración de cada paso.
  • Trazado distribuido: propague el trace_id a través de todos los servicios vía contexto. OpenTelemetry + Jaeger ofrecen una visión completa del recorrido de la saga por los servicios.
  • Dead Letter Queue (DLQ): en Redis Streams, los mensajes con múltiples fallos deben trasladarse a un stream separado para análisis manual.
  • Interfaz administrativa: un endpoint HTTP sencillo que devuelva el estado de una saga por su ID desde PostgreSQL resulta invaluable durante incidentes.
// monitoring.go
package orchestrator

import (
    "log/slog"
    "time"

    "github.com/prometheus/client_golang/prometheus"
)

var (
    sagaTotal = prometheus.NewCounterVec(
        prometheus.CounterOpts{Name: "saga_total", Help: "Total saga executions"},
        []string{"type", "status"},
    )
    sagaDuration = prometheus.NewHistogramVec(
        prometheus.HistogramOpts{
            Name:    "saga_duration_seconds",
            Help:    "Saga execution duration",
            Buckets: prometheus.DefBuckets,
        },
        []string{"type"},
    )
)

func init() {
    prometheus.MustRegister(sagaTotal, sagaDuration)
}

func recordSagaMetrics(sagaType string, status SagaStatus, start time.Time) {
    sagaTotal.WithLabelValues(sagaType, string(status)).Inc()
    sagaDuration.WithLabelValues(sagaType).Observe(time.Since(start).Seconds())
    slog.Info("saga finished",
        "type", sagaType,
        "status", status,
        "duration_ms", time.Since(start).Milliseconds(),
    )
}

Errores típicos y antipatrones

  • Manejadores no idempotentes. El error más frecuente. Si un paso de la saga no es idempotente, una llamada repetida generará operaciones duplicadas: doble cobro, doble reserva. Use siempre idempotency_key.
  • Ausencia de timeouts. Un paso de la saga puede quedarse bloqueado esperando respuesta de un servicio no disponible. Cada llamada debe tener un contexto con deadline. Use context.WithTimeout.
  • Compensación sin reintentos. Si una transacción compensatoria falla, es una situación crítica. Se necesitan reintentos, alertas y posiblemente intervención manual. Nunca ignore silenciosamente los errores de compensación.
  • Sagas demasiado largas. Una saga con 10 o más pasos es señal de problemas en la descomposición de servicios. Una saga larga implica una ventana prolongada de inconsistencia de datos y un grafo de compensación enorme.
  • Leer datos no confirmados. Mientras la saga se ejecuta, los datos están en un estado intermedio. Otros servicios que los lean pueden obtener una imagen incorrecta. Diseñe explícitamente para la consistencia eventual.
  • Almacenar el estado de la saga solo en memoria. Tras reiniciar el servicio, todas las sagas incompletas se pierden. El estado debe persistirse en PostgreSQL antes de iniciar cada paso.
  • Ignorar el orden de compensación. La compensación debe realizarse estrictamente en el orden inverso a los pasos ejecutados. De lo contrario, pueden generarse nuevas inconsistencias.

Conclusión: cómo elegir el enfoque adecuado

El patrón Saga resuelve el problema real de las transacciones distribuidas en sistemas de microservicios con Go, pero exige un diseño cuidadoso y disciplina en la implementación.

Saga no hace el sistema más simple — hace que la complejidad sea explícita y manejable.

Elija coreografía si tiene 2–4 participantes en la saga, el equipo ya trabaja con el enfoque event-driven y Redis Streams, y está dispuesto a invertir en un buen trazado para la depuración. La coreografía escala bien horizontalmente.

Elija orquestación si la saga contiene lógica de negocio compleja con ramificaciones, se necesita visibilidad completa del estado en PostgreSQL, el equipo prefiere el control explícito del flujo y la velocidad de resolución de incidentes es crítica.

Independientemente de la elección: implemente idempotencia desde el primer día, añada trazado distribuido antes de salir a producción, prevea un mecanismo de recuperación para sagas bloqueadas mediante escaneo periódico de la BD y nunca ignore los errores de las transacciones compensatorias. Go, con su manejo explícito de errores y sus potentes primitivas de concurrencia, es una excelente elección para implementar ambos enfoques.

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í →