Architecture

Implementing the Saga Pattern in Go Microservices: Orchestration and Choreography with PostgreSQL and Redis

Ruslan Ismailov Published 18 min read
I

Introduction: Why Distributed Transactions Are Hard

In a monolithic application, a transaction is an atomic operation within a single database. PostgreSQL's ACID guarantees make it reliable and predictable. But when a system is broken into microservices, each with its own database, classic transactions become impossible: you cannot issue a BEGIN ... COMMIT simultaneously across two isolated data stores.

Two-phase commit (2PC) formally addresses this problem, but in practice it introduces tight service coupling, blocking waits, and a single point of failure in the coordinator. For high-load systems, this is unacceptable.

The Saga pattern is an alternative approach: instead of one distributed transaction, a sequence of local transactions is executed in each service. If a failure occurs at any step, compensating transactions are triggered to undo the already-completed steps. This is not a rollback in the traditional sense — it consists of explicit reverse actions implemented in code.

Two Approaches: Orchestration vs. Choreography

The Saga pattern can be implemented in two fundamentally different ways. The choice between them shapes the architecture of the entire system.

Choreography

Each service publishes events after completing its local transaction. Other services subscribe to those events and react accordingly. There is no central coordinator — each service knows what to do in response to a specific event.

  • Pros: full decentralization, low coupling, easy to add new participants.
  • Cons: difficult to track the overall state of the saga, logic is spread across services, debugging becomes a detective exercise.

Orchestration

A central Saga Orchestrator manages all steps: it calls services, waits for responses, and decides when to trigger compensation. It stores the saga's state — typically in a database.

  • Pros: full visibility into saga state, centralized logic, clear step ordering, straightforward debugging.
  • Cons: the orchestrator is a potential single point of failure, tighter coupling to the coordinator, risk of the orchestrator becoming a "fat" service.

Comparing the Approaches

  • Saga state: with choreography — distributed across events and services; with orchestration — centralized in the orchestrator's database.
  • Coupling: choreography yields loose coupling between services; orchestration creates a dependency of all services on the orchestrator.
  • Observability: choreography requires aggregating logs and traces; orchestration provides complete state from a single place.
  • When to choose choreography: simple sagas (2–3 steps), the team already works with event-driven architecture, services are loosely coupled.
  • When to choose orchestration: complex business processes (5+ steps), explicit state monitoring is required, centralized compensation is important.

Implementing a Choreography Saga in Go with Redis Streams

Redis Streams are a great fit for choreography — they provide a persistent message queue with consumer groups. Unlike simple Pub/Sub, messages are stored and can be re-read after a failure.

Event Structure

// 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"`
}

Publishing Events to 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()
}

Event Consumer with 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 {
    // Create the consumer group if it doesn't exist
    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 // message remains in PEL for retry
                }
                // ACK only on successful processing
                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 // ignore unknown events
    }

    return handler(ctx, event)
}

Example: Payment Service in a Choreography Saga

// 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"]

    // Idempotency check — do not process twice
    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 {
        // Publish a failure event — other services will start compensation
        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 {
    // Compensating transaction: refund the payment
    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
}

Implementing an Orchestration Saga in Go with PostgreSQL

The orchestrator stores the state of each saga in PostgreSQL. This enables recovery after a failure: on restart, the orchestrator reads incomplete sagas and resumes execution.

Database Schema

-- 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);

State Model and Repository

// 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()
}

Central Orchestrator

// 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 {
        // Idempotency check: already completed?
        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)
        }

        // Save the step result for idempotency
        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)

        // Update the payload with the step's results
        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 {
    // Compensate in reverse order
    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 {
            // Log the error but continue compensating other steps
            log.Printf("[saga %s] compensation of %s failed: %v", instance.ID, step.Name, err)
        }
    }

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

Using the Orchestrator

// main.go (assembly example)
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{} // initialize with db

    orderSaga := orchestrator.NewSagaOrchestrator(repo, []orchestrator.Step{
        {
            Name: "reserve_inventory",
            Execute: func(ctx context.Context, payload map[string]any) (map[string]any, error) {
                // call inventory service
                return map[string]any{"reservation_id": "res-123"}, nil
            },
            Compensate: func(ctx context.Context, payload map[string]any) error {
                // cancel reservation
                return nil
            },
        },
        {
            Name: "process_payment",
            Execute: func(ctx context.Context, payload map[string]any) (map[string]any, error) {
                // call payment service
                return map[string]any{"payment_id": "pay-456"}, nil
            },
            Compensate: func(ctx context.Context, payload map[string]any) error {
                // refund payment
                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, // final step, no compensation needed
        },
    })

    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)
    }
}

Handling Compensating Transactions and Idempotency

Compensating transactions are not rollbacks. They are separate business operations that must be implemented explicitly. Several critically important rules apply:

Idempotency Is a Hard Requirement

Any operation in a saga — both forward and compensating — may be called more than once due to retries after network failures. Every handler must behave correctly when invoked multiple times.

The standard pattern is to store an idempotency_key (typically saga_id + step_name) in a results table. Before executing the operation, check whether a record with that key already exists.

// 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()

    // Atomic check + insert via 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 {
        // Operation already completed — return success
        return tx.Commit()
    }

    // Execute the actual payment
    if err := r.chargeCard(ctx, tx, amount); err != nil {
        return err
    }

    return tx.Commit()
}

Retries with Exponential Backoff

For transient failures (network issues, database overload), use retries with jitter to avoid thundering herds:

// 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
}

Monitoring and Debugging Saga Transactions

Without observability, sagas become a black box. For production systems, you need:

  • Structured logging: every saga event is logged with fields saga_id, step, status, duration_ms. Use log/slog or zerolog.
  • Metrics: counters for started, completed, and compensated sagas via Prometheus. Separate histograms for the duration of each step.
  • Distributed tracing: propagate trace_id across all services through the context. OpenTelemetry + Jaeger provide a complete picture of how a saga flows through services.
  • Dead Letter Queue (DLQ): in Redis Streams, messages that have failed repeatedly should be moved to a separate stream for manual analysis.
  • Admin endpoint: a simple HTTP endpoint that returns the saga state by its ID from PostgreSQL — invaluable during incidents.
// 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(),
    )
}

Common Mistakes and Anti-Patterns

  • Non-idempotent handlers. The most frequent mistake. If a saga step is not idempotent, a repeated call will create duplicate operations — double charges, double reservations. Always use an idempotency_key.
  • Missing timeouts. A saga step can hang indefinitely waiting for a response from an unavailable service. Every call must have a context with a deadline. Use context.WithTimeout.
  • Compensation without retries. If a compensating transaction fails, that is a critical situation. Retries, alerts, and possibly manual intervention are required. Never silently ignore compensation errors.
  • Overly long sagas. A saga with 10+ steps signals problems in service decomposition. A long saga means a prolonged window of data inconsistency and a massive compensation graph.
  • Reading uncommitted data. While a saga is executing, data is in an intermediate state. Other services reading that data may get an incorrect picture. Design for eventual consistency explicitly.
  • Storing saga state only in memory. After a service restart, all in-flight sagas are lost. State must be persisted to PostgreSQL before each step begins.
  • Ignoring compensation order. Compensation must happen in the strict reverse order of the executed steps. Otherwise, new inconsistencies can be introduced.

Conclusion: Choosing the Right Approach

The Saga pattern addresses a real problem of distributed transactions in Go microservice systems, but it requires careful design and disciplined implementation.

Saga doesn't make the system simpler — it makes complexity explicit and manageable.

Choose choreography if you have 2–4 participants in the saga, your team already works with an event-driven approach and Redis Streams, and you are willing to invest in good tracing for debugging. Choreography scales horizontally with ease.

Choose orchestration if the saga involves complex business logic with branching, you need full state visibility in PostgreSQL, the team prefers explicit flow control, and fast incident debugging is critical.

Regardless of your choice: implement idempotency from day one, add distributed tracing before going to production, provide a mechanism for recovering stalled sagas through periodic database scanning, and never ignore errors from compensating transactions. Go, with its explicit error handling and powerful concurrency primitives, is an excellent choice for implementing either approach.

Technologies

Tags

Ruslan Ismailov

Senior Web / Backend Developer. Senior web/backend developer with 9 years of experience. Stack: PHP, Laravel, PostgreSQL, Redis, Docker, Kubernetes, REST, microservices, CI/CD. More about me →