PrepZone Logo
PrepZone

Saga Pattern — Orchestration

A central coordinator drives compensating steps when ShardPay cross-shard transfers fail mid-flight.

Why this matters

  • Cross-shard transfers cannot use 2PC without sacrificing availability — orchestrated sagas are ShardPay's production path.
  • A durable saga log gives on-call a visible state machine instead of inferring progress from scattered logs.
  • Every forward step must have a defined compensating action or be safely retryable — interviewers probe this relentlessly.
  • Orchestration trades a coordinator dependency for debuggability; ShardPay accepts that trade for its highest-value payment flows.
Orchestrator
Debit shard
Credit shard
Compensate debit
Central coordinator executes steps and triggers compensations on failure.

Orchestrated cross-shard transfer

ShardPay's transfer saga has four steps: validate fraud score, debit source shard, credit destination shard, emit settlement event. The orchestrator persists state after each step so a crash mid-saga resumes from the last completed step, not from scratch.

When merchant M-8812 sends $1,200 from account A-US-7 (shard 3) to A-EU-2 (shard 9), the orchestrator creates saga instance SAGA-9f3a in state STARTED, calls fraud service, then transitions to DEBIT_PENDING. Only after a successful debit does it attempt credit — never both in one atomic operation.

Key points

  • Orchestrator — a dedicated service (or module) that owns saga state and invokes participants in order. ShardPay's TransferOrchestrator runs on Java 17 with a Postgres-backed state table and leader election via etcd.
  • Saga instance — one row per in-flight business transaction with current state, step outputs, and timestamps. ShardPay indexes sagas by transferId so support can look up SAGA-9f3a from a merchant ticket.
  • Compensation — a semantic undo of a completed forward step, not a database rollback. If credit fails after debit succeeds, ShardPay runs CompensateDebit — a new transaction that credits the source account back with reason SAGA_ROLLBACK.
  • Timeout per step — if credit does not complete within 30 seconds, the orchestrator marks the step failed and triggers compensation. ShardPay alerts when >50 sagas exceed step timeout in 5 minutes.

Saga state machine

ShardPay models each transfer as an explicit state machine. Illegal transitions are rejected at the orchestrator — participants cannot skip ahead or commit out of order.

States and transitions

Java
STARTED → FRAUD_OK → DEBITED → CREDITED → COMPLETED
                ↓         ↓         ↓
             ABORTED  COMPENSATING → COMPENSATED

A failure at CREDITED is rare (credit succeeded but emit failed); the orchestrator retries the emit step idempotently rather than compensating a successful credit. Compensation is reserved for steps that must be undone — debiting without a matching credit.

Java
// ShardPay orchestrated saga — core state transitions
public void advance(SagaInstance saga) {
    switch (saga.state()) {
        case STARTED -> {
            fraudService.check(saga.transferId());
            saga.transition(FRAUD_OK);
            persist(saga);
            advance(saga); // tail-recurse to next step
        }
        case FRAUD_OK -> {
            DebitResult debit = shardClient.debit(saga.fromShard(), saga.request());
            saga.recordDebit(debit);
            saga.transition(DEBITED);
            persist(saga);
            advance(saga);
        }
        case DEBITED -> {
            try {
                shardClient.credit(saga.toShard(), saga.request());
                saga.transition(CREDITED);
                persist(saga);
            } catch (CreditException e) {
                compensateDebit(saga);
            }
        }
        case CREDITED -> {
            outbox.publish(SettlementEvent.from(saga));
            saga.transition(COMPLETED);
            persist(saga);
        }
        default -> { /* terminal or compensating — no-op */ }
    }
}

private void compensateDebit(SagaInstance saga) {
    saga.transition(COMPENSATING);
    persist(saga);
    shardClient.creditBack(saga.fromShard(), saga.debitRef());
    saga.transition(COMPENSATED);
    persist(saga);
}

Recovery and idempotency

Orchestrator crashes are expected. On restart, ShardPay scans sagas in non-terminal states and resumes from the last persisted step. Each participant call carries sagaId + stepName as an idempotency key so retries do not double-debit.

Walkthrough: orchestrator crash after debit

  1. Saga SAGA-9f3a reaches state DEBITED; debit of $1,200 confirmed on shard 3.
  2. Orchestrator pod killed by K8s node drain before credit attempt.
  3. New orchestrator leader loads SAGA-9f3a, sees DEBITED, retries credit with idempotency key SAGA-9f3a:CREDIT.
  4. Shard 9 applies credit once; saga moves to CREDITED, then COMPLETED.

If shard 9 already received the credit (orchestrator crashed after credit but before persist), the idempotency key returns the stored result and the saga advances without duplicate credit.

Key points

  • Durable saga log — every state transition written before the next external call. ShardPay uses UPDATE saga SET state = ?, version = version + 1 WHERE id = ? AND version = ? for optimistic locking.
  • Participant idempotency — debit and credit endpoints accept Idempotency-Key: {sagaId}:{step}. ShardPay rejects same key with different payload (HTTP 409).
  • Orchestrator HA — active-passive with etcd lease; only the leader advances sagas. Failover resumes within 8 seconds — well under ShardPay's 30-second step timeout.
  • Observability — each saga step creates a trace span linked to the merchant's transferId. On-call dashboards show in-flight saga count by state.

Quick recall

Everything you need if you only revisit this box.

  1. Orchestrator drives ordered steps with a durable state machine.
  2. Every forward step needs a compensating action or safe retry semantics.
  3. Participant idempotency keys prevent double-debit on orchestrator retry.

Test yourself

Answer these before moving on — recall is what makes it stick.