Saga Pattern for Distributed Transactions: Orchestration, Compensation, and Production Pitfalls

Implementing the saga pattern for distributed transactions across 12 microservices with compensation logic, idempotency guarantees, and observability at 8K sagas/second.

#saga-pattern#microservices#transactions#async
Cover image for the article: Saga Pattern for Distributed Transactions: Orchestration, Compensation, and Production Pitfalls

Distributed transactions across microservices are the problem that never gets a clean solution — only managed complexity. Two-phase commit (2PC) provides strong consistency but creates tight coupling and availability problems. The saga pattern provides eventual consistency with explicit compensation logic, trading complexity in the protocol for independence in the services.

After implementing sagas across 12 microservices processing 8,000 distributed transactions per second, I can confirm: sagas work. They also introduce failure modes that will surprise you if your mental model comes from single-database transactions. This article covers the orchestration architecture, compensation strategies, and the operational tooling that makes sagas debuggable in production.

Why Sagas Over 2PC

We abandoned 2PC after it caused three outages in two months. The fundamental issue: 2PC requires all participants to be available simultaneously. When Service C is slow, Services A and B hold locks waiting for the coordinator's commit decision. At our scale (8K transactions/second), even 200ms of lock contention cascades into system-wide degradation.

Property2PCSaga (Orchestration)Saga (Choreography)
ConsistencyStrong (ACID)EventualEventual
AvailabilityLimited by slowest participantHigh (independent services)High
CouplingTight (all must be available)Loose (via orchestrator)Very loose
DebuggabilityModerateHigh (centralized state)Low (distributed state)
Compensation complexityNone (rollback)Explicit per stepDistributed

We chose orchestration-based sagas because debuggability matters more than theoretical purity in production. When a saga fails at step 7 of 12, the operations team needs to see the complete state machine in one place.

Architecture Overview

Our saga orchestrator is a stateless service backed by a durable state store (PostgreSQL) that manages the lifecycle of each saga instance.

Saga Pattern Architecture

Saga Definition

Each saga is defined as a sequence of steps with corresponding compensation actions. Steps execute forward; compensations execute backward when a step fails.

// Saga definition for order fulfillment (12-step distributed transaction)
interface SagaStep {
  name: string;
  service: string;
  action: string;
  compensation: string;
  timeout: number;
  retryPolicy: RetryPolicy;
  idempotencyKey: (ctx: SagaContext) => string;
}

interface RetryPolicy {
  maxAttempts: number;
  backoffMs: number;
  backoffMultiplier: number;
  maxBackoffMs: number;
  retryableErrors: string[];
}

const ORDER_FULFILLMENT_SAGA: SagaStep[] = [
  {
    name: 'validate_order',
    service: 'order-service',
    action: 'orders.validate',
    compensation: 'orders.invalidate',
    timeout: 5000,
    retryPolicy: { maxAttempts: 3, backoffMs: 100, backoffMultiplier: 2, maxBackoffMs: 2000, retryableErrors: ['TIMEOUT', 'UNAVAILABLE'] },
    idempotencyKey: (ctx) => `validate-${ctx.sagaId}`
  },
  {
    name: 'reserve_inventory',
    service: 'inventory-service',
    action: 'inventory.reserve',
    compensation: 'inventory.release',
    timeout: 3000,
    retryPolicy: { maxAttempts: 5, backoffMs: 200, backoffMultiplier: 2, maxBackoffMs: 5000, retryableErrors: ['TIMEOUT', 'UNAVAILABLE', 'CONFLICT'] },
    idempotencyKey: (ctx) => `reserve-${ctx.sagaId}-${ctx.data.itemId}`
  },
  {
    name: 'authorize_payment',
    service: 'payment-service',
    action: 'payments.authorize',
    compensation: 'payments.void',
    timeout: 10000,
    retryPolicy: { maxAttempts: 2, backoffMs: 500, backoffMultiplier: 2, maxBackoffMs: 2000, retryableErrors: ['TIMEOUT'] },
    idempotencyKey: (ctx) => `auth-${ctx.sagaId}`
  },
  {
    name: 'apply_discount',
    service: 'promotion-service',
    action: 'promotions.apply',
    compensation: 'promotions.reverse',
    timeout: 3000,
    retryPolicy: { maxAttempts: 3, backoffMs: 100, backoffMultiplier: 2, maxBackoffMs: 1000, retryableErrors: ['TIMEOUT', 'UNAVAILABLE'] },
    idempotencyKey: (ctx) => `discount-${ctx.sagaId}`
  },
  {
    name: 'calculate_tax',
    service: 'tax-service',
    action: 'tax.calculate',
    compensation: 'tax.void',
    timeout: 5000,
    retryPolicy: { maxAttempts: 3, backoffMs: 200, backoffMultiplier: 2, maxBackoffMs: 2000, retryableErrors: ['TIMEOUT', 'UNAVAILABLE'] },
    idempotencyKey: (ctx) => `tax-${ctx.sagaId}`
  },
  {
    name: 'capture_payment',
    service: 'payment-service',
    action: 'payments.capture',
    compensation: 'payments.refund',
    timeout: 15000,
    retryPolicy: { maxAttempts: 3, backoffMs: 1000, backoffMultiplier: 2, maxBackoffMs: 10000, retryableErrors: ['TIMEOUT'] },
    idempotencyKey: (ctx) => `capture-${ctx.sagaId}`
  }
];

Saga Orchestrator Engine

// Saga orchestrator with compensation and idempotency
import { v4 as uuidv4 } from 'uuid';

enum SagaStatus {
  RUNNING = 'running',
  COMPLETED = 'completed',
  COMPENSATING = 'compensating',
  COMPENSATED = 'compensated',
  FAILED = 'failed'  // Compensation also failed — requires manual intervention
}

interface SagaInstance {
  sagaId: string;
  sagaType: string;
  status: SagaStatus;
  currentStep: number;
  completedSteps: string[];
  data: Record<string, any>;
  error?: string;
  startedAt: Date;
  completedAt?: Date;
}

class SagaOrchestrator {
  constructor(
    private stateStore: SagaStateStore,
    private messageBroker: MessageBroker,
    private metrics: MetricsCollector
  ) {}

  async startSaga(sagaType: string, data: Record<string, any>): Promise<string> {
    const sagaId = uuidv4();
    const instance: SagaInstance = {
      sagaId,
      sagaType,
      status: SagaStatus.RUNNING,
      currentStep: 0,
      completedSteps: [],
      data,
      startedAt: new Date()
    };

    await this.stateStore.save(instance);
    this.metrics.increment('saga.started', { type: sagaType });
    
    await this.executeNextStep(instance);
    return sagaId;
  }

  async handleStepResult(sagaId: string, stepName: string, result: StepResult): Promise<void> {
    const instance = await this.stateStore.load(sagaId);
    
    if (result.success) {
      instance.completedSteps.push(stepName);
      instance.data = { ...instance.data, ...result.data };
      instance.currentStep++;
      
      const steps = this.getSteps(instance.sagaType);
      if (instance.currentStep >= steps.length) {
        instance.status = SagaStatus.COMPLETED;
        instance.completedAt = new Date();
        this.metrics.increment('saga.completed', { type: instance.sagaType });
        this.metrics.histogram('saga.duration_ms', 
          Date.now() - instance.startedAt.getTime(),
          { type: instance.sagaType });
      } else {
        await this.executeNextStep(instance);
      }
    } else {
      instance.error = result.error;
      instance.status = SagaStatus.COMPENSATING;
      this.metrics.increment('saga.compensation_started', { 
        type: instance.sagaType,
        failed_step: stepName 
      });
      await this.startCompensation(instance);
    }
    
    await this.stateStore.save(instance);
  }

  private async startCompensation(instance: SagaInstance): Promise<void> {
    const steps = this.getSteps(instance.sagaType);
    
    // Compensate in reverse order
    for (let i = instance.completedSteps.length - 1; i >= 0; i--) {
      const stepName = instance.completedSteps[i];
      const step = steps.find(s => s.name === stepName)!;
      
      await this.messageBroker.publish(step.service, {
        action: step.compensation,
        sagaId: instance.sagaId,
        idempotencyKey: `compensate-${step.idempotencyKey({ 
          sagaId: instance.sagaId, 
          data: instance.data 
        } as SagaContext)}`,
        data: instance.data
      });
    }
  }
}

Idempotency: The Non-Negotiable Requirement

Every saga step and compensation must be idempotent. Network failures, message redelivery, and orchestrator restarts will cause duplicate invocations. If your steps are not idempotent, you will corrupt data.

// Idempotent step execution with deduplication store
class IdempotentStepExecutor {
  constructor(private deduplicationStore: DeduplicationStore) {}

  async execute(
    idempotencyKey: string,
    operation: () => Promise<StepResult>,
    ttlSeconds: number = 86400
  ): Promise<StepResult> {
    // Check if this exact operation was already executed
    const existing = await this.deduplicationStore.get(idempotencyKey);
    if (existing) {
      return existing; // Return cached result — no re-execution
    }

    const result = await operation();
    
    // Store result for deduplication
    await this.deduplicationStore.set(idempotencyKey, result, ttlSeconds);
    
    return result;
  }
}

Compensation Patterns

Not all compensations are simple reversals. Some require semantic compensation that accounts for side effects.

StepForward ActionNaive CompensationSemantic Compensation
Reserve inventoryDecrement available countIncrement countRelease specific reservation (handles concurrent reservations)
Authorize paymentHold funds on cardVoid authorizationVoid if same day, refund if next day (card network rules)
Send notificationEmail sent to customerCannot unsendSend correction/apology email
Update analyticsIncrement conversion counterDecrement counterEmit compensation event for analytics pipeline
Generate invoiceCreate invoice documentDelete invoiceVoid invoice with audit trail (legal requirement)

Observability and Debugging

Saga debugging requires end-to-end tracing across all participating services. We built a saga dashboard that shows the complete state machine for any saga instance.

Key Metrics

MetricCurrent ValueAlert Threshold
Sagas started/second8,200N/A
Saga completion rate99.7%< 99%
Saga average duration1.2s> 5s
Compensation triggered rate0.3%> 2%
Failed compensation (manual intervention)0.002%> 0.01%
Saga p99 duration4.8s> 15s

Dead Letter Handling

When compensation fails (the service is down, the data is inconsistent), sagas enter a FAILED state requiring manual intervention. We process approximately 15 failed sagas per day out of 700M+ total.

-- Query for sagas requiring manual intervention
SELECT 
    saga_id,
    saga_type,
    status,
    current_step,
    error,
    started_at,
    (NOW() - started_at) as age,
    data->>'orderId' as order_id,
    data->>'customerId' as customer_id,
    completed_steps
FROM saga_instances
WHERE status = 'failed'
AND started_at > NOW() - INTERVAL '24 hours'
ORDER BY started_at DESC;

Saga Pattern Anti-Patterns

Anti-pattern 1: Saga within a saga. Nested sagas create exponential complexity in compensation paths. If step 3 of saga A triggers saga B, and saga B fails at step 5, you must compensate saga B then continue compensating saga A. Instead, flatten the workflow.

Anti-pattern 2: Non-idempotent compensations. If compensating "deduct $50" means "add $50" without checking the current state, a retry of the compensation will add $50 twice. Always check the current state before compensating.

Anti-pattern 3: Timeout-based completion. Assuming a step completed because it did not respond within the timeout. Always require explicit acknowledgment. A timed-out step may have succeeded — you need to check before compensating.

Performance at Scale

ConfigurationThroughputAvg LatencyP99 Latency
1 orchestrator instance2,400 sagas/s1.4s6.2s
4 orchestrator instances8,200 sagas/s1.2s4.8s
8 orchestrator instances14,800 sagas/s1.3s5.1s

Orchestrator instances are stateless — they read saga state from the database and publish messages. Horizontal scaling is linear up to the database write capacity (which we shard by saga_id).

Conclusion

The saga pattern is the pragmatic answer to distributed transactions in microservices. It trades the simplicity of ACID rollbacks for the availability and loose coupling that microservices demand. The cost is explicit: you must implement compensation logic for every step, ensure idempotency throughout, and build observability tooling that makes the distributed state machine debuggable.

Start with the orchestration approach if your team is new to sagas. The centralized state machine is dramatically easier to debug than choreography-based sagas where state is scattered across services. Only move to choreography when the orchestrator becomes a bottleneck — and at 8K sagas/second on 4 instances, that threshold is higher than most teams will reach.

Comments

    No comments yet. Be the first to share your thoughts.