Chuyển đến nội dung chính

Lesson 10: Saga Pattern — Distributed Transactions

Why 2PC is not suitable for microservices, Saga Pattern (Choreography vs Orchestration), compensating transactions, Saga Orchestrator implementation, error handling and dead letter queue.

🏗️ Architecture — Lesson 10 Lesson 10: Saga Pattern — Distributed Transactions

Cloud Native Microservices Architecture

Part 3: Data Management in Microservices

xdev.asia

Lesson 10: Saga Pattern — Distributed Transactions

Introduction

In monolith, a business transaction can include multiple database operations in the same ACID transaction. In microservices, each service has its own database — cannot use traditional distributed transactions (2PC) because it creates tight coupling and affects performance. Saga Pattern is an alternative.


1. Problem with Distributed Transactions

1.1 Two-Phase Commit (2PC) — Why not suitable

2PC Flow:
┌──────────────┐    ┌──────────────┐    ┌──────────────┐
│ Transaction  │    │  Service A   │    │  Service B   │
│ Coordinator  │    │  (Order)     │    │  (Payment)   │
└──────┬───────┘    └──────┬───────┘    └──────┬───────┘
       │                   │                   │
       │── Phase 1: PREPARE ──────────────────▶│
       │◀── VOTE YES ────────────────────────── │
       │── Phase 1: PREPARE ──▶│               │
       │◀── VOTE YES ──────── │               │
       │                      │               │
       │── Phase 2: COMMIT ──▶│               │
       │── Phase 2: COMMIT ──────────────────▶│
       │                      │               │

2PC problem in microservices:

ProblemExplanation
Single Point of FailureCoordinator dies → all participants are locked
Performance bottleneckLock resources during 2PC
Tight couplingAll services must be available at the same time
Not scalableLock contention increases with the number of participants
Network partitionIf the network breaks between Phase 1 and Phase 2 → inconsistent state

1.2 CAP Theorem repeats

In a distributed system, only 2 out of 3 can be achieved:

  • Consistency: Every read returns the latest write
  • Availability: Every request receives a response
  • Partition tolerance: The system continues to operate when the network is divided

Microservices select AP (Availability + Partition tolerance) → accept Eventual Consistency.


2. Saga Pattern

2.1 Concepts

Saga is a sequence of local transactions, each transaction performed by a service. If a step fails, the saga performs compensating transactions to undo the previous steps.

Saga = T1 → T2 → T3 → ... → Tn

Nếu Ti fail:
  Compensate: C(i-1) → C(i-2) → ... → C1

Ví dụ Order Saga:
  T1: Create Order (status: PENDING)
  T2: Reserve Payment
  T3: Reserve Inventory
  T4: Confirm Order (status: CONFIRMED)

Nếu T3 fail:
  C2: Refund Payment
  C1: Cancel Order (status: CANCELLED)

2.2 Compensating Transactions

Each transaction in the saga must have a corresponding compensating transaction:

StepActionCompensating Action
T1Create OrderCancel Order
T2Reserve PaymentRefund Payment
T3Reserve InventoryRelease Inventory
T4Shipping ScheduleCancel Shipping
T5Send Confirmation EmailSend Cancellation Email

Important note:

  • Compensating transactions must be idempotent (run multiple times with the same result)
  • Compensating transactions cannot fail (must retry until successful)
  • Some actions cannot be compensated (for example, send a sent email) → use "semantic undo" (send a canceled email)

3. Choreography Saga

3.1 Concepts

Each service publish event when completing the local transaction, the next service subscribes and processes:

┌──────────┐    ┌──────────────┐    ┌──────────────┐    ┌──────────────┐
│  Order   │    │   Payment    │    │  Inventory   │    │  Shipping    │
│  Service │    │   Service    │    │  Service     │    │  Service     │
└────┬─────┘    └──────┬───────┘    └──────┬───────┘    └──────┬───────┘
     │                 │                   │                   │
     │── OrderCreated ▶│                   │                   │
     │                 │── PaymentReserved ▶│                   │
     │                 │                   │── InventoryReserved ▶│
     │                 │                   │                   │── ShipmentScheduled
     │◀───────────────────────────────────────────────────────── │
     │  OrderConfirmed │                   │                   │
     │                 │                   │                   │
     │ === FAILURE === │                   │                   │
     │                 │                   │── InventoryFailed ▶│
     │                 │◀─ CompensatePayment─                  │
     │◀─ PaymentRefunded─                  │                   │
     │  OrderCancelled │                   │                   │

3.2 Implementation Example

// Order Service — publishes OrderCreated
@Service
public class OrderService {

    @Transactional
    public Order createOrder(CreateOrderCommand cmd) {
        Order order = new Order(cmd.getCustomerId(), cmd.getItems());
        order.setStatus(OrderStatus.PENDING);
        orderRepository.save(order);

        // Publish event
        eventPublisher.publish(new OrderCreatedEvent(
            order.getId(),
            order.getCustomerId(),
            order.getItems(),
            order.getTotalAmount()
        ));

        return order;
    }

    // Compensating handler
    @EventHandler
    public void on(PaymentFailedEvent event) {
        Order order = orderRepository.findById(event.getOrderId());
        order.setStatus(OrderStatus.CANCELLED);
        order.setFailureReason(event.getReason());
        orderRepository.save(order);
    }
}

// Payment Service — listens to OrderCreated
@Service
public class PaymentService {

    @EventHandler
    public void on(OrderCreatedEvent event) {
        try {
            Payment payment = paymentGateway.reserve(
                event.getCustomerId(),
                event.getTotalAmount()
            );
            eventPublisher.publish(new PaymentReservedEvent(
                event.getOrderId(),
                payment.getId()
            ));
        } catch (InsufficientFundsException e) {
            eventPublisher.publish(new PaymentFailedEvent(
                event.getOrderId(),
                "Insufficient funds"
            ));
        }
    }
}

3.3 Advantages and disadvantages

AdvantagesDisadvantages
Loose coupling between servicesDifficult to track complex flows
Simple for less steps sagaCyclic dependencies can occur
There is no single point of failureDifficult to debug when there are errors
Natural with event-driven architectureComplex Testing

4. Orchestration Saga

4.1 Concepts

A central Saga Orchestrator coordinates the entire workflow, sending commands to each service and processing the response:

                    ┌──────────────────────┐
                    │  Order Saga          │
                    │  Orchestrator        │
                    │                      │
                    │  State Machine:      │
                    │  CREATED             │
                    │  → PAYMENT_PENDING   │
                    │  → INVENTORY_PENDING │
                    │  → SHIPPING_PENDING  │
                    │  → CONFIRMED         │
                    │  or → COMPENSATING   │
                    │  → CANCELLED         │
                    └────────┬─────────────┘
                             │
            ┌────────────────┼────────────────┐
            │                │                │
     ┌──────▼──────┐  ┌─────▼──────┐  ┌──────▼──────┐
     │  Payment    │  │ Inventory  │  │  Shipping   │
     │  Service    │  │ Service    │  │  Service    │
     └─────────────┘  └────────────┘  └─────────────┘

4.2 State Machine Implementation

public class OrderSaga {

    public enum State {
        CREATED,
        PAYMENT_PENDING,
        PAYMENT_RESERVED,
        INVENTORY_PENDING,
        INVENTORY_RESERVED,
        SHIPPING_PENDING,
        CONFIRMED,
        COMPENSATING_INVENTORY,
        COMPENSATING_PAYMENT,
        CANCELLED
    }

    @Autowired
    private SagaRepository sagaRepository;

    public void start(CreateOrderCommand cmd) {
        SagaState saga = new SagaState(cmd.getOrderId(), State.CREATED);
        sagaRepository.save(saga);
        
        // Step 1: Reserve Payment
        saga.setState(State.PAYMENT_PENDING);
        commandGateway.send(new ReservePaymentCommand(
            cmd.getOrderId(), cmd.getAmount()
        ));
    }

    @SagaEventHandler
    public void on(PaymentReservedEvent event) {
        SagaState saga = sagaRepository.findByOrderId(event.getOrderId());
        saga.setState(State.INVENTORY_PENDING);
        
        // Step 2: Reserve Inventory
        commandGateway.send(new ReserveInventoryCommand(
            event.getOrderId(), saga.getItems()
        ));
    }

    @SagaEventHandler
    public void on(InventoryReservedEvent event) {
        SagaState saga = sagaRepository.findByOrderId(event.getOrderId());
        saga.setState(State.SHIPPING_PENDING);
        
        // Step 3: Schedule Shipping
        commandGateway.send(new ScheduleShippingCommand(
            event.getOrderId(), saga.getAddress()
        ));
    }

    @SagaEventHandler
    public void on(ShipmentScheduledEvent event) {
        SagaState saga = sagaRepository.findByOrderId(event.getOrderId());
        saga.setState(State.CONFIRMED);
        
        commandGateway.send(new ConfirmOrderCommand(event.getOrderId()));
    }

    // === COMPENSATION ===
    
    @SagaEventHandler
    public void on(InventoryReservationFailedEvent event) {
        SagaState saga = sagaRepository.findByOrderId(event.getOrderId());
        saga.setState(State.COMPENSATING_PAYMENT);
        
        // Compensate: Refund Payment
        commandGateway.send(new RefundPaymentCommand(
            event.getOrderId(), saga.getPaymentId()
        ));
    }

    @SagaEventHandler
    public void on(PaymentRefundedEvent event) {
        SagaState saga = sagaRepository.findByOrderId(event.getOrderId());
        saga.setState(State.CANCELLED);
        
        commandGateway.send(new CancelOrderCommand(
            event.getOrderId(), "Inventory not available"
        ));
    }
}

4.3 Advantages and disadvantages

AdvantagesDisadvantages
Easy to understand workflow (centralized)Orchestrator is single point of failure
Easy to debug and monitorRisk of god class (too much logic)
Easy to add/remove stepsCoupling between orchestrator and services
Compensation logic clearlyNeed persistent saga state

5. Compare Choreography vs Orchestration

CriteriaChoreographyOrchestration
CouplingVery looseMedium (via orchestrator)
ComplexityIncreases by number of servicesFocus in orchestrator
VisibilityDifficult to see overall flowExplicitly in state machine
Failure handlingDisperseFocus
TestingDifficult (distributed)Easier (test orchestrator)
ScalabilityGoodGood (orchestrator stateless)
Recommended≤ 4 services in saga> 4 services or complex flow

6. Error Handling & Dead Letter Queue

6.1 Dead Letter Queue (DLQ)

Messages that cannot be processed after multiple retries are passed into DLQ:

                    ┌─────────────┐
                    │  Main Queue │
                    │ (order.cmds)│
                    └──────┬──────┘
                           │
                    ┌──────▼──────┐
                    │  Consumer   │
                    │  (Service)  │
                    └──────┬──────┘
                           │
                  Success? ─┼─ No (after max retries)
                  │         │
                  ▼         ▼
               ┌─────┐  ┌──────────┐
               │ ACK │  │   DLQ    │
               └─────┘  │(order.   │
                         │ cmds.dlq)│
                         └──────────┘
                              │
                         ┌────▼────┐
                         │ Alert + │
                         │ Manual  │
                         │ Review  │
                         └─────────┘

6.2 Idempotency

Make sure to process the message multiple times for the same result:

@Service
public class PaymentService {

    @EventHandler
    public void on(ReservePaymentCommand cmd) {
        // Idempotency check
        String idempotencyKey = "payment:" + cmd.getOrderId();
        if (processedStore.exists(idempotencyKey)) {
            log.info("Already processed payment for order {}", cmd.getOrderId());
            return; // Skip duplicate
        }

        Payment payment = processPayment(cmd);

        // Mark as processed
        processedStore.save(idempotencyKey, payment.getId());
    }
}

6.3 Saga Timeout

public class OrderSaga {

    @SagaTimeout(duration = "5m")
    public void onTimeout(SagaState saga) {
        log.warn("Saga timeout for order {}", saga.getOrderId());
        
        // Compensate based on current state
        switch (saga.getState()) {
            case INVENTORY_PENDING:
                compensatePayment(saga);
                break;
            case SHIPPING_PENDING:
                compensateInventory(saga);
                compensatePayment(saga);
                break;
        }
        
        saga.setState(State.CANCELLED);
        saga.setFailureReason("Saga timeout");
    }
}

7. Best Practices

7.1 Saga Design Guidelines

1. Mỗi step phải có compensating action
2. Compensating actions phải idempotent
3. Sử dụng correlation ID (saga ID) xuyên suốt
4. Persist saga state (survive service restart)
5. Set timeout cho mỗi saga instance
6. Monitor saga metrics (success rate, duration, failure reasons)
7. DLQ cho messages không xử lý được
8. Tránh saga quá nhiều steps (> 7 steps → xem lại design)

7.2 Choose Choreography or Orchestration?

Flow đơn giản (2-4 services)?
  └── Choreography

Flow phức tạp (> 4 services, conditional logic)?
  └── Orchestration

Cần visibility cao vào business process?
  └── Orchestration

Muốn minimize coupling?
  └── Choreography

Team experience với event-driven?
  ├── Nhiều → Choreography
  └── Ít → Orchestration

Summary

  • 2PC is not suitable for microservices due to tight coupling and lock contention
  • Saga Pattern uses chain of local transactions + compensating transactions
  • Choreography: Services communicate via events, decentralized
  • Orchestration: Saga Orchestrator coordinates centrally through the state machine
  • Compensating transactions must be idempotent and must not fail
  • Use DLQ, timeout and idempotency for error handling