Skip to main content

Saga Pattern (Distributed Workflows)

Who this guide is for

1. The Problem: Why Can't We Just Use a Database Transaction?

In a monolith with a single database, @Transactional is sufficient. Any failure rolls everything back atomically, and the database is always consistent.

In a microservices architecture, each service owns its own database. There is no shared database connection and no distributed transaction coordinator in most modern stacks. When an e-commerce checkout must:

  1. Create an order in order-service (PostgreSQL)
  2. Reserve stock in inventory-service (MySQL)
  3. Charge the customer via payment-service (Stripe API)
  4. Send a confirmation via notification-service (SendGrid API)

...there is no BEGIN TRANSACTION that spans all four. If Step 3 fails after Steps 1 and 2 have succeeded, the system is partially committed β€” the customer's stock is reserved, an order exists, but no payment was taken.

Why Not Two-Phase Commit (2PC)?

2PC is the classical solution: a coordinator asks all participants to "prepare" (lock resources), then issues a global "commit" or "abort". It provides true distributed ACID semantics.

Two-Phase Commit (2PC) β€” Why It Fails at Scale

Coordinator(Single Point of Failure)Order DBπŸ”’ LOCKED (holding lock)"Ready"Inventory DBπŸ”’ LOCKED (holding lock)"Ready"Payment APIπŸ”’ LOCKED (holding lock)"Ready"← Prepare? "Lock your resources"
❌ Blocking locks
All DBs hold row locks for the entire 2PC duration β†’ throughput collapses under concurrency.
❌ Coordinator SPOF
If coordinator crashes between phases β†’ participants hold locks indefinitely (no timeout to unlock).
❌ External APIs incompatible
Stripe, SendGrid have no prepare/commit support β€” cannot join a 2PC protocol.
❌ Availability sacrifice
CAP theorem: 2PC chooses consistency over availability. One unreachable participant blocks all.

πŸ’‘ Toggle phases above to see how 2PC causes locking across all participants simultaneously.

Why 2PC fails at microservices scale:

ProblemImpact
Blocking locksAll participants hold locks for the entire 2PC duration β€” throughput collapses under concurrency
Coordinator SPOFIf the coordinator crashes between phases, participants are stuck holding locks indefinitely
External API incompatibilityStripe, SendGrid, and most external APIs are not 2PC-aware β€” they cannot participate in a prepare/commit protocol
Cross-database unavailabilityPostgres + MySQL + Stripe cannot share a single distributed transaction manager
Availability vs consistency2PC sacrifices availability β€” if any participant is unreachable, the whole transaction blocks

The Saga pattern accepts eventual consistency instead of ACID consistency, trading distributed locking for independent service progress and explicit compensation logic.


2. The Travel Booking Analogy

Imagine booking a trip: flight + hotel + car rental. Each is booked independently with a different company.

Step 1: Book flight βœ… (Confirmed: BA-123)
Step 2: Book hotel βœ… (Confirmed: Marriott #456)
Step 3: Book car ❌ No cars available at that price

What do you do?
β†’ Call the hotel β†’ Cancel booking #456
β†’ Call the airline β†’ Cancel booking BA-123

There is no single "undo" button across all three.
Each cancellation is a NEW forward-moving action β€” not a rollback.

This is exactly a Saga: a sequence of local transactions, where each step has a compensating action that semantically undoes its effect if a later step fails.

The fundamental insight: you cannot rewind time in a distributed system. The order was created, the stock was reserved, the email was queued. Compensation acknowledges this β€” it creates new, observable business events that undo the effect, rather than pretending the earlier steps never happened.


3. Saga Pattern β€” Core Mechanics

A Saga decomposes a distributed business transaction into a sequence of local ACID transactions, one per service. Each local transaction:

  1. Updates its own database (local ACID guarantee β€” only this service's data)
  2. Signals the next step (via event or command)

Happy Path & Failure Path

Saga Core Mechanics β€” Happy Path

T1: Create Orderorder-serviceT2: Reserve Stockinventory-serviceT3: Process Paymentpayment-serviceT4: Send Confirmationnotification-serviceCOMPLETED βœ…

πŸ’‘ Switch between paths to see how compensating transactions execute in reverse when a step fails.

What Are Compensating Transactions?

Compensations are NOT database rollbacks

A database rollback erases changes as if they never happened. A compensating transaction is a new, committed, observable business operation that semantically reverses the effect of a prior step. The prior step's commit is permanent β€” compensation creates new state on top of it.

StepForward TransactionCompensation
T1: Create OrderINSERT INTO orders (status='PENDING')UPDATE orders SET status='CANCELLED'
T2: Reserve StockUPDATE inventory SET reserved = reserved + 1UPDATE inventory SET reserved = reserved - 1
T3: Charge CardPOST /stripe/chargesPOST /stripe/refunds
T4: Send EmailPOST /sendgrid/send(No compensation β€” email already delivered)

Notice T1's compensation does NOT DELETE the order β€” it marks it CANCELLED. This preserves the audit trail, which is mandatory in financial systems. T4 has no technical compensation β€” an email cannot be unsent. This is a pivot transaction (see Section 9).


4. Saga Coordination Styles

Two fundamentally different patterns exist for coordinating a Saga:

Choreography β€” Coordination Style

OrderServiceInventoryServicePaymentServiceβ†’ OrderCreated← event replyβ†’ StockReserved← event replyChoreography Properties:β€’ Each service reacts to events from the previous β†’ workflow is implicitβ€’ No central controller β€” distributed decision-makingβ€’ Decoupled in theory, but coupled by event contracts in practiceβ€’ Simple 2–4 step flows: fine. Complex branching: workflow becomes invisible.

πŸ’‘ Toggle to compare how the same saga plays out under each coordination style.

DimensionChoreographyOrchestration
ControlDistributed β€” each service reacts to eventsCentralized β€” one orchestrator drives all steps
CommunicationAsync eventsDirected commands (sync or async)
Workflow visibility❌ Implicit β€” must read every serviceβœ… Explicit β€” one place to understand the flow
CouplingServices coupled to event contractsServices coupled to orchestrator's command API
Failure handlingDistributed β€” each service emits compensation eventsCentralized β€” orchestrator coordinates all compensations
DebuggingRequires distributed tracing across all servicesSingle orchestrator log shows full saga history
ScalingEach service scales independentlyOrchestrator must be HA
Best forSimple linear flows, high team autonomyComplex branching flows, compliance-heavy systems

5. Choreography β€” Deep Dive

In choreography, each service subscribes to events from the previous service and publishes events that trigger the next service. The workflow emerges from event propagation β€” no central controller exists.

Full Event Flow

Choreography Flow Map

Client πŸ’»Order Service πŸ“¦Inventory Service πŸ—„οΈPayment Service πŸ’³Notification Service βœ‰οΈPOST /ordersEmit: OrderCreatedReserve StockEmit: StockReservedCharge CustomerEmit: PaymentProcessedEmit: PaymentProcessedUpdate Order β†’ CONFIRMEDSend Confirmation Email
Step 1 of 9
CLIENT β†’ OS

POST /orders

Client submits a new order. Order Service writes order to database with PENDING status.

πŸ’‘ Click any step name directly inside the sequence map or use the stepper to explore step-by-step.

Choreography Implementation β€” Inventory Service

@Component
@RequiredArgsConstructor
@Slf4j
public class InventoryChoreographer {

private final InventoryRepository inventoryRepository;
private final OutboxRepository outboxRepository;

// Forward step: triggered by OrderCreated event
@KafkaListener(topics = "order-events", groupId = "inventory-service")
@Transactional // Local ACID: reserve stock AND write outbox event atomically
public void onOrderCreated(OrderCreatedEvent event) {
log.info("Reserving stock for order {}", event.getOrderId());

// Idempotency guard: skip if already processed (retry safety)
if (inventoryRepository.reservationExists(event.getOrderId())) {
log.info("Duplicate OrderCreated for order {} β€” skipping", event.getOrderId());
return;
}

try {
inventoryRepository.reserve(event.getOrderId(), event.getItems());

// Write success event to outbox IN SAME TRANSACTION as business state
// If this transaction commits: reservation AND event are both durable
// If it fails: neither happens β€” safe to retry
outboxRepository.save(OutboxEvent.builder()
.aggregateId(event.getOrderId())
.eventType("StockReserved")
.payload(new StockReservedEvent(event.getOrderId(), event.getAmount()))
.createdAt(Instant.now())
.build());

} catch (InsufficientStockException e) {
log.warn("Insufficient stock for order {}: {}", event.getOrderId(), e.getMessage());

// Failure event also written atomically with any partial state changes
outboxRepository.save(OutboxEvent.builder()
.aggregateId(event.getOrderId())
.eventType("StockReservationFailed")
.payload(new StockReservationFailedEvent(event.getOrderId(), e.getMessage()))
.build());
}
}

// Compensation step: triggered by PaymentFailed event
@KafkaListener(topics = "payment-events", groupId = "inventory-service")
@Transactional
public void onPaymentFailed(PaymentFailedEvent event) {
log.info("Releasing stock for order {} β€” payment failed", event.getOrderId());

// Idempotency guard: skip if already compensated
if (!inventoryRepository.reservationExists(event.getOrderId())) {
log.info("No reservation found for order {} β€” already compensated or never reserved",
event.getOrderId());
return;
}

inventoryRepository.release(event.getOrderId());

outboxRepository.save(OutboxEvent.builder()
.aggregateId(event.getOrderId())
.eventType("StockReleased")
.payload(new StockReleasedEvent(event.getOrderId()))
.build());
}
}

Choreography: The Hidden Coupling Problem

Choreography appears decoupled, but the services are implicitly coupled through the event contract:

OrderService publishes: OrderCreated { orderId, items, customerId, amount }
InventoryService expects: orderId, items
PaymentService expects: orderId, amount, customerId

If OrderCreated drops the "amount" field:
β†’ InventoryService still works (doesn't use amount)
β†’ PaymentService BREAKS (cannot charge without amount)
β†’ The coupling is invisible until runtime

With orchestration:
The orchestrator's PaymentCommand explicitly includes amount
The contract is visible in one place β€” the orchestrator

When Choreography Becomes Unmanageable

When Choreography Becomes Unmanageable

Simple (2 services) β€” Choreography is fine βœ…OrderServicePaymentServiceβœ… DONEOrderCreated β†’ PaymentProcessed β†’ Done. Simple, traceable, manageable.βœ… Why choreography works here:β€’ Only 2 services β€” easy to trace the event flow in logsβ€’ No conditional branches β€” one happy path onlyβ€’ Small team owns both services β€” event schema changes are visibleβ€’ Debugging: check 2 service logs max

πŸ’‘ Toggle to see how choreography degrades from simple linear flows to complex conditional branching.


6. Orchestration β€” Deep Dive

An Orchestrator is a dedicated service that explicitly commands each participant, waits for replies, and drives the saga forward β€” including coordinating all compensations from one place.

Full Orchestration Flow

Orchestration Flow Map (Payment Failure Exception)

Saga Orchestrator πŸ€–Inventory Service πŸ—„οΈPayment Service πŸ’³Order Service πŸ“¦State: STARTEDReserveStock commandStockReserved reply βœ…State: STOCK_RESERVEDProcessPayment commandPaymentFailed reply ❌State: COMPENSATINGReleaseStock command (Comp)StockReleased reply βœ…CancelOrder command (Comp)OrderCancelled reply βœ…State: CANCELLED βœ…
Step 1 of 12Saga State: STARTED
LOCAL_STATE

State: STARTED

Saga state initialized and persisted locally as STARTED before executing downstream network calls.

πŸ’‘ Play through step-by-step to see how the orchestrator drives command & compensation flows, tracking state transitions in its local database.

Saga State Machine

A production orchestrator must model the saga lifecycle as a formal state machine to prevent invalid transitions, enable recovery, and provide auditability.

Saga State Machine β€” Order Saga Lifecycle

STARTEDSTOCK_RESERVEDPAYMENT_PROCESSEDNOTIFIEDCOMPLETED βœ…STOCK_FAILPAYMENT_FAILEDCOMPENSATINGCANCELLED βœ…MANUAL REQ ⚠️

πŸ’‘ Click any state to highlight its transitions and see description.

public enum SagaStatus {
STARTED,
STOCK_RESERVED,
STOCK_RESERVATION_FAILED,
PAYMENT_PROCESSED,
PAYMENT_FAILED,
NOTIFIED,
COMPLETED,
COMPENSATING,
CANCELLED,
MANUAL_INTERVENTION_REQUIRED;

// Valid forward transitions β€” enforced by the state machine
private static final Map<SagaStatus, Set<SagaStatus>> VALID_TRANSITIONS = Map.of(
STARTED, Set.of(STOCK_RESERVED, STOCK_RESERVATION_FAILED),
STOCK_RESERVED, Set.of(PAYMENT_PROCESSED, PAYMENT_FAILED),
PAYMENT_PROCESSED, Set.of(NOTIFIED, COMPENSATING),
NOTIFIED, Set.of(COMPLETED),
STOCK_RESERVATION_FAILED, Set.of(CANCELLED),
PAYMENT_FAILED, Set.of(COMPENSATING),
COMPENSATING, Set.of(CANCELLED, MANUAL_INTERVENTION_REQUIRED)
);

public boolean canTransitionTo(SagaStatus next) {
return VALID_TRANSITIONS.getOrDefault(this, Set.of()).contains(next);
}

public boolean isTerminal() {
return switch (this) {
case COMPLETED, CANCELLED, MANUAL_INTERVENTION_REQUIRED -> true;
default -> false;
};
}
}
@Entity
@Table(name = "saga_state")
@Data
@Builder
public class SagaState {

@Id
private String id; // sagaId β€” globally unique

@Enumerated(EnumType.STRING)
private SagaStatus status;

private String orderId;
private String customerId;
private BigDecimal amount;

private String failureReason;
private int compensationAttempts;

@Version
private Long version; // Optimistic locking β€” prevents concurrent updates

private Instant createdAt;
private Instant updatedAt;
private Instant completedAt;

public void transition(SagaStatus next) {
if (!status.canTransitionTo(next)) {
throw new InvalidSagaTransitionException(
String.format("Saga %s: invalid transition %s β†’ %s", id, status, next));
}
this.status = next;
this.updatedAt = Instant.now();
if (next.isTerminal()) {
this.completedAt = Instant.now();
}
}
}

Full Orchestrator Implementation

@Service
@RequiredArgsConstructor
@Slf4j
public class OrderSagaOrchestrator {

private final SagaStateRepository sagaRepository;
private final InventoryClient inventoryClient;
private final PaymentClient paymentClient;
private final NotificationClient notificationClient;
private final OrderRepository orderRepository;
private final SagaMetrics metrics;
private final SagaEscalationService escalationService;

// Entry point β€” called when a new order is placed
@Transactional
public void startSaga(CreateOrderCommand cmd) {
String sagaId = UUID.randomUUID().toString();
Instant startTime = Instant.now();

// Persist initial state BEFORE any network call
// If we crash after this line, recovery job can resume the saga
SagaState state = SagaState.builder()
.id(sagaId)
.status(SagaStatus.STARTED)
.orderId(cmd.getOrderId())
.customerId(cmd.getCustomerId())
.amount(cmd.getAmount())
.createdAt(startTime)
.build();
sagaRepository.save(state);
metrics.recordSagaStarted("order-saga");

log.info("[Saga:{}] Started for order {}", sagaId, cmd.getOrderId());
advanceSaga(sagaId, cmd);
}

// State machine driver β€” called on each step (and by recovery jobs)
@Transactional
public void advanceSaga(String sagaId, CreateOrderCommand cmd) {
SagaState state = sagaRepository.findById(sagaId)
.orElseThrow(() -> new SagaNotFoundException(sagaId));

if (state.getStatus().isTerminal()) {
log.info("[Saga:{}] Already terminal ({}), skipping", sagaId, state.getStatus());
return;
}

try {
switch (state.getStatus()) {

case STARTED -> {
log.info("[Saga:{}] Step 1: Reserving stock", sagaId);
// Persist BEFORE calling downstream β€” idempotency key ensures safe retry
inventoryClient.reserve(ReserveStockCommand.of(sagaId, cmd));
state.transition(SagaStatus.STOCK_RESERVED);
sagaRepository.save(state);
log.info("[Saga:{}] Stock reserved successfully", sagaId);
advanceSaga(sagaId, cmd); // Drive to next step
}

case STOCK_RESERVED -> {
log.info("[Saga:{}] Step 2: Processing payment", sagaId);
paymentClient.charge(ProcessPaymentCommand.of(sagaId, cmd));
state.transition(SagaStatus.PAYMENT_PROCESSED);
sagaRepository.save(state);
log.info("[Saga:{}] Payment processed successfully", sagaId);
advanceSaga(sagaId, cmd);
}

case PAYMENT_PROCESSED -> {
log.info("[Saga:{}] Step 3: Sending notification", sagaId);
notificationClient.sendConfirmation(NotifyCommand.of(sagaId, cmd));
state.transition(SagaStatus.NOTIFIED);
sagaRepository.save(state);
advanceSaga(sagaId, cmd);
}

case NOTIFIED -> {
log.info("[Saga:{}] All steps complete β€” marking order confirmed", sagaId);
orderRepository.markConfirmed(cmd.getOrderId());
state.transition(SagaStatus.COMPLETED);
sagaRepository.save(state);
metrics.recordSagaCompleted("order-saga",
Duration.between(state.getCreatedAt(), Instant.now()));
log.info("[Saga:{}] COMPLETED successfully", sagaId);
}
}

} catch (InsufficientStockException e) {
// Non-retryable: domain failure β€” no stock available
log.warn("[Saga:{}] Stock unavailable β€” saga fails without compensation needed", sagaId);
state.transition(SagaStatus.STOCK_RESERVATION_FAILED);
state.setFailureReason("Insufficient stock: " + e.getMessage());
sagaRepository.save(state);
orderRepository.markFailed(cmd.getOrderId(), "Stock unavailable");
state.transition(SagaStatus.CANCELLED);
sagaRepository.save(state);
metrics.recordSagaFailed("order-saga", "STOCK_RESERVATION_FAILED");

} catch (PaymentException e) {
// Payment failed β€” must compensate steps that already succeeded
log.error("[Saga:{}] Payment failed β€” initiating compensation. Reason: {}",
sagaId, e.getMessage());
state.transition(SagaStatus.COMPENSATING);
state.setFailureReason("Payment failed: " + e.getMessage());
sagaRepository.save(state);
metrics.recordSagaFailed("order-saga", "PAYMENT_FAILED");
metrics.recordCompensationStarted("order-saga");
compensate(sagaId, cmd);

} catch (TransientException e) {
// Transient failure (network blip, upstream timeout) β€” will be retried by recovery job
log.warn("[Saga:{}] Transient error in step {} β€” will retry: {}",
sagaId, state.getStatus(), e.getMessage());
// Do NOT transition state β€” let the recovery job retry from current step
throw e;
}
}

// Compensation: executed in reverse order of forward steps
@Transactional
private void compensate(String sagaId, CreateOrderCommand cmd) {
SagaState state = sagaRepository.findById(sagaId).orElseThrow();

try {
// Compensate Step 2 (stock reservation) β€” the only completed forward step
// that needs reversal (payment failed β†’ payment never committed β†’ no refund needed)
log.info("[Saga:{}] Compensation: releasing stock reservation", sagaId);
inventoryClient.releaseReservation(ReleaseStockCommand.of(sagaId, cmd));

log.info("[Saga:{}] Compensation: cancelling order", sagaId);
orderRepository.markCancelled(cmd.getOrderId(), state.getFailureReason());

state.transition(SagaStatus.CANCELLED);
sagaRepository.save(state);
log.info("[Saga:{}] Compensation complete β€” saga CANCELLED", sagaId);

} catch (Exception compensationEx) {
// Compensation itself failed β€” requires human intervention
log.error("[Saga:{}] COMPENSATION FAILED at step: {}", sagaId, compensationEx.getMessage());
state.setCompensationAttempts(state.getCompensationAttempts() + 1);
sagaRepository.save(state);
escalationService.escalate(sagaId, "compensation", compensationEx);
}
}
}

7. Async Orchestration β€” Kafka-Based Commands and Replies

The synchronous orchestrator above blocks while waiting for each downstream service. Under high load, this ties up threads and does not gracefully handle long-running downstream operations. A production-grade orchestrator uses async messaging.

Async Kafka Orchestration β€” Command & Reply Topics

Saga OrchestratorNon-blocking Β· publishes commands Β· handles repliesinventory-commands← ReserveStock / ReleaseStockinventory-repliesβ†’ StockReserved / StockFailedpayment-commands← ProcessPaymentpayment-repliesβ†’ PaymentProcessed / FailedKAFKA BROKERSagaId used as message key β†’ partition affinityOrchestrator thread never blocksSaga can span minutes / hours / human approvalAt-least-once delivery + idempotency = exactly-once effectnotification-commands / notification-replies (not shown)

πŸ’‘ Click a step chip above to see what happens at that point in the async Kafka-based orchestration.

Async Orchestrator

@Service
@RequiredArgsConstructor
@Slf4j
public class AsyncOrderSagaOrchestrator {

private final SagaStateRepository sagaRepository;
private final KafkaTemplate<String, Object> kafka;
private final OrderRepository orderRepository;

// Step 1: Start the saga β€” publish first command
@Transactional
public void startSaga(CreateOrderCommand cmd) {
String sagaId = UUID.randomUUID().toString();

SagaState state = SagaState.builder()
.id(sagaId).status(SagaStatus.STARTED)
.orderId(cmd.getOrderId()).amount(cmd.getAmount())
.createdAt(Instant.now()).build();
sagaRepository.save(state);

// Publish reserve-stock command to Kafka
kafka.send("inventory-commands",
sagaId, // Use sagaId as Kafka key β†’ all messages for this saga go to same partition
ReserveStockCommand.of(sagaId, cmd));

log.info("[Saga:{}] Started β€” ReserveStock command published", sagaId);
}

// Reply handler β€” called when inventory service responds
@KafkaListener(topics = "inventory-replies", groupId = "saga-orchestrator")
@Transactional
public void onInventoryReply(SagaReply reply) {
SagaState state = sagaRepository.findById(reply.getSagaId()).orElseThrow();

if (state.getStatus().isTerminal()) return; // Already resolved

if (reply instanceof StockReservedEvent event) {
state.transition(SagaStatus.STOCK_RESERVED);
sagaRepository.save(state);

// Trigger next step: payment
kafka.send("payment-commands", state.getId(),
ProcessPaymentCommand.of(state.getId(), state.getOrderId(), state.getAmount()));
log.info("[Saga:{}] StockReserved β€” ProcessPayment command published", state.getId());

} else if (reply instanceof StockReservationFailedEvent event) {
state.transition(SagaStatus.STOCK_RESERVATION_FAILED);
state.setFailureReason(event.getReason());
sagaRepository.save(state);
orderRepository.markCancelled(state.getOrderId(), event.getReason());
state.transition(SagaStatus.CANCELLED);
sagaRepository.save(state);
log.warn("[Saga:{}] StockReservationFailed β€” saga CANCELLED", state.getId());
}
}

// Reply handler β€” called when payment service responds
@KafkaListener(topics = "payment-replies", groupId = "saga-orchestrator")
@Transactional
public void onPaymentReply(SagaReply reply) {
SagaState state = sagaRepository.findById(reply.getSagaId()).orElseThrow();

if (state.getStatus().isTerminal()) return;

if (reply instanceof PaymentProcessedEvent event) {
state.transition(SagaStatus.PAYMENT_PROCESSED);
sagaRepository.save(state);

// Trigger next step: notification
kafka.send("notification-commands", state.getId(),
SendConfirmationCommand.of(state.getId(), state.getOrderId()));
log.info("[Saga:{}] PaymentProcessed β€” SendConfirmation command published", state.getId());

} else if (reply instanceof PaymentFailedEvent event) {
state.transition(SagaStatus.COMPENSATING);
state.setFailureReason("Payment: " + event.getReason());
sagaRepository.save(state);

// Trigger compensation: release stock
kafka.send("inventory-commands", state.getId(),
ReleaseStockCommand.of(state.getId(), state.getOrderId()));
log.warn("[Saga:{}] PaymentFailed β€” ReleaseStock command published", state.getId());
}
}

// Compensation reply handler
@KafkaListener(topics = {"inventory-replies"}, groupId = "saga-orchestrator-compensation")
@Transactional
public void onCompensationReply(SagaReply reply) {
SagaState state = sagaRepository.findById(reply.getSagaId()).orElseThrow();

if (state.getStatus() != SagaStatus.COMPENSATING) return;

if (reply instanceof StockReleasedEvent) {
orderRepository.markCancelled(state.getOrderId(), state.getFailureReason());
state.transition(SagaStatus.CANCELLED);
sagaRepository.save(state);
log.info("[Saga:{}] Compensation complete β€” CANCELLED", state.getId());
}
}
}

Key architectural property of async orchestration: the orchestrator's thread is never blocked waiting for a downstream response. It publishes a command and immediately handles other replies. The saga can take minutes or hours and no thread is held. This is essential for sagas involving human approval steps, external partner callbacks, or long-running processes.


8. Parallel Saga Steps

Not all saga steps are sequential. If two steps are independent (don't need each other's output), they can run in parallel β€” reducing saga completion time.

Parallel Saga Steps β€” Reducing Saga Completion Time

Sequential: T1 β†’ T2 β†’ T3 β†’ T4 β€” Time = T1 + T2 + T3 + T4T1: Create Order~50msT2: Reserve Stock~200msT3: Fraud Check~300msT4: Payment~400msTotal time: ~50 + ~200 + ~300 + ~400 = ~950ms (worst case β€” all serial)T3 (fraud check) must complete before T4 (payment) can even start

πŸ’‘ Toggle to see the time saving when independent saga steps run in parallel.

// Async orchestrator: trigger T2 and T3 simultaneously
@Transactional
public void afterOrderCreated(String sagaId) {
SagaState state = sagaRepository.findById(sagaId).orElseThrow();
state.transition(SagaStatus.PARALLEL_VALIDATION_STARTED);
state.setPendingParallelSteps(Set.of("stock-reserve", "fraud-check")); // Track completions
sagaRepository.save(state);

// Publish both commands simultaneously
kafka.send("inventory-commands", sagaId, ReserveStockCommand.of(sagaId, ...));
kafka.send("fraud-commands", sagaId, FraudCheckCommand.of(sagaId, ...));
}

// Only advance when ALL parallel steps have completed
@Transactional
public void onParallelStepComplete(String sagaId, String completedStep) {
SagaState state = sagaRepository.findById(sagaId).orElseThrow();
state.getPendingParallelSteps().remove(completedStep);

if (state.getPendingParallelSteps().isEmpty()) {
// All parallel steps done β€” advance to payment
state.transition(SagaStatus.PARALLEL_VALIDATION_COMPLETE);
sagaRepository.save(state);
kafka.send("payment-commands", sagaId, ProcessPaymentCommand.of(sagaId, ...));
} else {
sagaRepository.save(state); // Save updated pending set
}
}

9. Failure Taxonomy & Compensation Design

Not all failures in a saga are equal. The correct response depends on the failure type.

Failure Types

Failure TypeExamplesCorrect Response
Transient infrastructureNetwork timeout, connection reset, temporary DB unavailabilityRetry with exponential backoff + jitter
Transient domainLock contention, optimistic locking conflict, rate limit (429)Retry with backoff
Permanent domainInsufficient stock, invalid payment method, fraud blockedFail immediately, trigger compensation β€” retrying will always fail
Compensation failureInventory service down while releasing reservationRetry compensation; escalate to ops if exhausted
Pivot transactionEmail delivered, SMS sent, payment wire sent to bankCannot compensate technically β€” create business adjustment

Distinguishing Transient from Permanent Errors

@Component
public class SagaErrorClassifier {

public boolean isRetryable(Exception ex) {
return switch (ex) {
case ConnectException e -> true; // Network β€” transient
case SocketTimeoutException e -> true; // Network β€” transient
case OptimisticLockException e -> true; // DB conflict β€” transient
case HttpServerErrorException e
when e.getStatusCode().is5xxServerError()
&& e.getStatusCode().value() != 501 -> true; // 5xx except "Not Implemented"
case InsufficientStockException e -> false; // Domain β€” permanent
case PaymentDeclinedException e -> false; // Domain β€” permanent
case FraudBlockedException e -> false; // Domain β€” permanent
case HttpClientErrorException e -> false; // 4xx β€” permanent
default -> false; // Unknown β€” treat as permanent; don't retry blindly
};
}
}

Pivot Transactions β€” The Irreversibility Boundary

A pivot transaction is the point of no return in a saga β€” the step after which compensation becomes impossible or only possible through business-level adjustment:

E-commerce order saga:

T1: Reserve stock ← Compensable (release reservation)
T2: Charge payment card ← Compensable (refund)
T3: *** PIVOT ***
T4: Dispatch to warehouse ← Compensable only if not yet picked (narrow window)
T5: Ship parcel ← NOT compensable (return process required)
T6: Deliver to customer ← NOT compensable (return + refund required)

Design implications of pivot transactions:

  1. Place pivot transactions as late as possible β€” maximize the window for technical compensation.
  2. Handle post-pivot failures as business processes β€” create return requests, refund records, or adjustment entries rather than expecting a technical undo.
  3. Pivot transitions are always MANUAL_INTERVENTION_REQUIRED if something goes wrong after them.
// After T4 (warehouse dispatch) β€” if payment later fails (e.g., chargeback discovered)
// Technical compensation is no longer possible; create a business adjustment instead
private void handlePostPivotFailure(SagaState state, String reason) {
log.error("[Saga:{}] Failure after pivot point β€” creating business adjustment record", state.getId());

// Create a return/refund request (new business process, not a technical undo)
returnRequestService.createReturnRequest(ReturnRequest.builder()
.orderId(state.getOrderId())
.reason(reason)
.requestedBy("saga-orchestrator")
.build());

state.transition(SagaStatus.MANUAL_INTERVENTION_REQUIRED);
state.setFailureReason("Post-pivot failure: " + reason + " β€” return request created");
sagaRepository.save(state);

alertingService.sendCritical("Post-pivot saga failure requiring ops review",
Map.of("sagaId", state.getId(), "orderId", state.getOrderId()));
}

10. Idempotency & Exactly-Once Semantics

Every saga step will be retried under failure scenarios. Every step must be idempotent β€” executing it twice must produce the same observable result as executing it once.

Why Retries Are Inevitable

Scenario 1: Orchestrator sends ReserveStock command, inventory reserves stock,
reply is lost (network partition).
Orchestrator retries β†’ ReserveStock command sent again.
Without idempotency: stock reserved TWICE β†’ inventory inconsistency.

Scenario 2: Orchestrator crashes after publishing command but before persisting
STOCK_RESERVED state.
Recovery job re-reads STARTED state β†’ re-publishes ReserveStock command.
Without idempotency: second reservation attempt.

Conclusion: every participant must be idempotent with respect to sagaId + stepName.

Strategy 1 β€” sagaId as Idempotency Key

// Inventory service: idempotent reserve endpoint
@PostMapping("/inventory/reserve")
@Transactional
public ResponseEntity<Void> reserve(@RequestBody ReserveStockCommand cmd) {
String idempotencyKey = "reserve:" + cmd.getSagaId();

// Atomic check-then-insert with unique constraint
try {
processedCommandRepository.insert(
new ProcessedCommand(idempotencyKey, Instant.now())
);
} catch (DataIntegrityViolationException e) {
// Already processed this saga step β€” return success (idempotent)
log.info("Duplicate reserve command for saga {} β€” idempotent skip", cmd.getSagaId());
return ResponseEntity.ok().build();
}

// Only reached if this is the first time we've seen this sagaId + step
inventoryRepository.reserve(cmd.getOrderId(), cmd.getItems());
return ResponseEntity.ok().build();
}

Strategy 2 β€” Outbox Pattern (Atomic Event Publishing)

The single most important reliability pattern in saga step implementations:

@Transactional
public void reserveStock(ReserveStockCommand cmd) {
// Problem without Outbox Pattern:
// 1. Reserve stock in DB ← DB commit
// 2. Publish StockReserved to Kafka ← might fail!
//
// If step 2 fails: stock is reserved but orchestrator never hears about it
// β†’ Orchestrator times out β†’ retries β†’ double reservation
//
// With Outbox Pattern: both DB write and event are in the SAME transaction

// Business operation
inventoryRepository.reserve(cmd.getOrderId(), cmd.getItems());

// Event written to outbox IN THE SAME TRANSACTION
// Outbox relay (Debezium CDC or polling) publishes it to Kafka reliably
outboxRepository.save(OutboxEvent.builder()
.sagaId(cmd.getSagaId())
.aggregateType("Inventory")
.aggregateId(cmd.getOrderId())
.eventType("StockReserved")
.payload(objectMapper.writeValueAsString(
new StockReservedEvent(cmd.getSagaId(), cmd.getOrderId())))
.build());

// If this transaction commits:
// β†’ stock reserved βœ…, outbox event written βœ…
// β†’ Outbox relay will publish to Kafka reliably (at-least-once)
// If this transaction rolls back:
// β†’ neither happens βœ… β€” safe to retry
}

Strategy 3 β€” Saga-Step-Aware Upsert

For steps that represent full entity state (not incremental actions):

// Order status update β€” idempotent by design
@Transactional
public void updateOrderStatus(String orderId, OrderStatus targetStatus) {
jdbcTemplate.update("""
UPDATE orders
SET status = ?, updated_at = now()
WHERE order_id = ?
AND status != ? -- Only update if not already in target status
""",
targetStatus.name(), orderId, targetStatus.name()
);
// Re-running this with the same targetStatus is a no-op β€” idempotent by SQL semantics
}

11. Production Failure Handling

Saga Timeout and Recovery

Sagas that are stuck (waiting for a reply that never comes) must be detected and re-driven:

@Service
@RequiredArgsConstructor
@Slf4j
public class SagaRecoveryJob {

private final SagaStateRepository sagaRepository;
private final AsyncOrderSagaOrchestrator orchestrator;

@Scheduled(fixedDelay = 60_000) // Run every 60 seconds
public void recoverStuckSagas() {
Instant stuckThreshold = Instant.now().minus(Duration.ofMinutes(5));

// Find sagas that haven't progressed in 5 minutes and aren't terminal
List<SagaState> stuckSagas = sagaRepository
.findByStatusNotInAndUpdatedAtBefore(
Set.of(SagaStatus.COMPLETED, SagaStatus.CANCELLED,
SagaStatus.MANUAL_INTERVENTION_REQUIRED),
stuckThreshold
);

log.info("Recovery job found {} stuck sagas", stuckSagas.size());

for (SagaState state : stuckSagas) {
try {
log.warn("[Saga:{}] Recovering stuck saga in state {}", state.getId(), state.getStatus());
// Re-drive from current persisted state
orchestrator.advanceSaga(state.getId(), buildCommandFromState(state));
} catch (Exception e) {
log.error("[Saga:{}] Recovery failed: {}", state.getId(), e.getMessage());
}
}
}
}

Retry with Exponential Backoff + Jitter

@Component
@RequiredArgsConstructor
public class SagaStepExecutor {

private static final int MAX_ATTEMPTS = 5;
private static final long BASE_DELAY_MS = 500L;
private static final long MAX_DELAY_MS = 30_000L;

private final SagaErrorClassifier errorClassifier;

public <T> T executeWithRetry(String sagaId, String step, Supplier<T> action) {
Exception lastException = null;

for (int attempt = 1; attempt <= MAX_ATTEMPTS; attempt++) {
try {
return action.get();
} catch (Exception e) {
if (!errorClassifier.isRetryable(e)) {
// Permanent failure β€” do not retry, propagate immediately
log.error("[Saga:{}] Step {} permanently failed: {}", sagaId, step, e.getMessage());
throw e;
}

lastException = e;
if (attempt == MAX_ATTEMPTS) break;

long delay = calculateBackoff(attempt);
log.warn("[Saga:{}] Step {} attempt {}/{} failed (transient). Retrying in {}ms. Error: {}",
sagaId, step, attempt, MAX_ATTEMPTS, delay, e.getMessage());

try { Thread.sleep(delay); }
catch (InterruptedException ie) { Thread.currentThread().interrupt(); throw new RuntimeException(ie); }
}
}

throw new SagaStepExhaustedException(
String.format("Saga %s step %s failed after %d attempts", sagaId, step, MAX_ATTEMPTS),
lastException);
}

private long calculateBackoff(int attempt) {
// Exponential backoff with full jitter
long exponential = BASE_DELAY_MS * (1L << (attempt - 1)); // 500ms, 1s, 2s, 4s, 8s
long capped = Math.min(exponential, MAX_DELAY_MS);
return (long) (Math.random() * capped); // Full jitter: random(0, capped)
}
}

Escalation Playbook

Escalation Playbook β€” Saga Failure Response Levels

L1Automatic Retry0–5 minutes Β· 5 attemptsTransient errors only Exponential backoff + Logged at WARN levelL2Recovery Job5–30 minutesRuns every 60s, finds Re-reads persisted sagLogged at WARN levelL3DLQ Routing30 minutesSaga moved to DLQ tablIsolated from healthy Logged at ERROR levelL4Manual Intervention RequiredImmediately on permanent failureSagaStatus β†’ MANUAL_INPagerDuty / alerting fAdmin dashboard shows L5Business ResolutionPost-pivot failuresPost-pivot: technical Create return request Finance / ops team not

πŸ’‘ Click any level to see what happens at that escalation stage.

@Service
@RequiredArgsConstructor
public class SagaEscalationService {

private final SagaStateRepository sagaRepository;
private final AlertingService alertingService;
private final DlqPublisher dlqPublisher;

public void escalate(String sagaId, String failedStep, Exception cause) {
SagaState state = sagaRepository.findById(sagaId).orElseThrow();
int attempts = state.incrementCompensationAttempts();
sagaRepository.save(state);

if (attempts <= 5) {
// Level 1: let retry executor handle it
return;
}

if (attempts == 6) {
// Level 3: publish to DLQ for isolation
dlqPublisher.publish(FailedSagaMessage.of(sagaId, failedStep, cause));
log.error("[Saga:{}] Published to DLQ after {} compensation attempts", sagaId, attempts);
}

// Level 4: alert ops immediately
state.transition(SagaStatus.MANUAL_INTERVENTION_REQUIRED);
state.setFailureReason(failedStep + ": " + cause.getMessage());
sagaRepository.save(state);

alertingService.sendCritical("Saga requires manual intervention",
Map.of(
"sagaId", sagaId,
"failedStep", failedStep,
"compensationAttempts", String.valueOf(attempts),
"error", cause.getMessage()
));
}
}

12. Observability β€” Metrics, Tracing, and Alerting

Saga Metrics

@Component
@RequiredArgsConstructor
public class SagaMetrics {

private final MeterRegistry registry;

public void recordSagaStarted(String sagaType) {
registry.counter("saga.started", "type", sagaType).increment();
}

public void recordSagaCompleted(String sagaType, Duration duration) {
registry.timer("saga.duration", "type", sagaType, "result", "completed")
.record(duration);
}

public void recordSagaFailed(String sagaType, String failedStep) {
registry.counter("saga.failed", "type", sagaType, "step", failedStep).increment();
}

public void recordCompensationStarted(String sagaType) {
registry.counter("saga.compensation.started", "type", sagaType).increment();
}

public void recordManualIntervention(String sagaType, String step) {
registry.counter("saga.manual_intervention", "type", sagaType, "step", step).increment();
}

public void recordStepDuration(String sagaType, String step, Duration duration) {
registry.timer("saga.step.duration", "type", sagaType, "step", step)
.record(duration);
}
}

Prometheus Alert Rules

groups:
- name: saga-alerts
rules:
# Critical: manual intervention required
- alert: SagaManualInterventionRequired
expr: increase(saga_manual_intervention_total[5m]) > 0
labels:
severity: critical
annotations:
summary: "Saga requires manual intervention: {{ $labels.type }} step {{ $labels.step }}"

# Warning: compensation rate elevated
- alert: SagaCompensationRateHigh
expr: >
rate(saga_compensation_started_total[5m])
/ rate(saga_started_total[5m]) > 0.05
for: 5m
labels:
severity: warning
annotations:
summary: "Saga compensation rate >5% β€” downstream service may be degrading"

# Warning: saga failure rate elevated
- alert: SagaFailureRateHigh
expr: >
rate(saga_failed_total[5m])
/ rate(saga_started_total[5m]) > 0.01
for: 5m
labels:
severity: warning
annotations:
summary: "Saga failure rate >1% for {{ $labels.type }}"

# Warning: saga p99 duration elevated
- alert: SagaDurationHigh
expr: histogram_quantile(0.99, rate(saga_duration_seconds_bucket[5m])) > 30
for: 10m
labels:
severity: warning
annotations:
summary: "Saga p99 duration >30s β€” downstream bottleneck suspected"

# Critical: sagas stuck in COMPENSATING
- alert: SagasStuckInCompensating
expr: >
count(saga_state == "COMPENSATING" and saga_updated_at < now() - 300) > 0
for: 5m
labels:
severity: critical
annotations:
summary: "Sagas stuck in COMPENSATING state β€” compensation service may be down"

Structured Logging for Distributed Tracing

// Every saga log line must include sagaId and orderId for correlation
// Use MDC to automatically include in all log output within the thread
MDC.put("sagaId", sagaId);
MDC.put("orderId", orderId);
MDC.put("sagaStep", currentStep);

// Trace propagation: pass traceId/spanId through Kafka message headers
// so distributed traces span all services involved in the saga
ProducerRecord<String, Object> record = new ProducerRecord<>(topic, sagaId, command);
record.headers()
.add("X-Trace-Id", traceId.getBytes())
.add("X-Span-Id", spanId.getBytes())
.add("X-Saga-Id", sagaId.getBytes())
.add("X-Order-Id", orderId.getBytes());

13. Choreography vs Orchestration β€” Full Decision Guide

Choreography vs Orchestration β€” Full Decision Guide

Start
Question 1 of 4
Simple linear flow? (≀4 services, no branching)

Interview Questions

Q: What is the difference between a Saga and a distributed transaction (2PC)?

A distributed transaction (2PC) achieves strong ACID atomicity across databases by holding locks during a two-phase prepare/commit protocol. All participants commit or abort together β€” strongly consistent but blocking, slow, and impossible with external APIs. A Saga achieves eventual consistency through a sequence of local ACID transactions, with explicit compensating transactions for failure recovery. Sagas accept that the system may be temporarily inconsistent during execution, in exchange for availability, scalability, and compatibility with any API.

Q: Can a Saga leave the system in a temporarily inconsistent state?

Yes β€” by design. While a saga executes, the system is in an intermediate state (e.g., stock reserved but payment not yet processed). External observers can see this intermediate state. This is the price of eventual consistency. The saga guarantees the system will eventually reach a consistent terminal state β€” either fully completed or fully compensated. Systems using sagas must account for this exposure window in their read models (e.g., don't show the order as "confirmed" until all saga steps complete).

Q: How do you handle a compensation transaction that itself fails?

Compensation failures are the hardest case in saga design. The response depends on the failure type: transient failures (network, timeout) are retried with exponential backoff. If retries are exhausted, the saga transitions to MANUAL_INTERVENTION_REQUIRED and alerts operations. If the compensation is for a post-pivot transaction (e.g., a payment already sent to a bank), no technical compensation is possible β€” a new business process is created (refund request, return order) and finance teams are involved. The key principle is: never silently discard a failed compensation β€” always escalate to human oversight.

Q: How do you ensure idempotency in saga steps?

Every saga step must be idempotent because retries are guaranteed. The primary mechanism is a processed_commands table keyed by sagaId + stepName with a unique constraint. Before processing, insert this key β€” if the insert fails with a unique constraint violation, the step was already processed and we return success. The business state change and the outbox event must both be written in the same local database transaction (Outbox Pattern) to prevent the case where state is changed but the event is lost, or the event is published but the state change is rolled back.

Q: What is a pivot transaction and why does it matter?

A pivot transaction is the point in a saga after which compensation becomes technically impossible β€” for example, dispatching an order to a warehouse or sending a wire transfer to a bank. After the pivot, failure handling becomes a business process (returns, refunds, adjustments) rather than a technical undo. Good saga design places pivot transactions as late as possible in the flow, validates all preconditions before reaching the pivot, and transitions to MANUAL_INTERVENTION_REQUIRED on post-pivot failures rather than attempting futile technical compensation.

Q: Choreography or Orchestration β€” which do you prefer and when?

Orchestration for any production system with more than trivial complexity. The orchestrator's persisted state machine provides auditability, debuggability, and a single place to understand and modify the workflow. Choreography's apparent decoupling is often illusory β€” services are still coupled through event contracts, but that coupling is invisible. The first time you need to debug a cross-service failure at 2 AM and have to correlate events across 6 services simultaneously, you appreciate having one orchestrator log that shows the full saga history. I would choose choreography only for truly simple, well-bounded flows where the team is disciplined about distributed tracing and the workflow is unlikely to grow in complexity.


15. See Also

πŸ“–
Track Page Progress0 / 635 Read
Knowledge Base Completion0%