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.
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
TransferOrchestratorruns 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
transferIdso support can look upSAGA-9f3afrom 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 reasonSAGA_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
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.
// 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
- Saga
SAGA-9f3areaches stateDEBITED; debit of $1,200 confirmed on shard 3. - Orchestrator pod killed by K8s node drain before credit attempt.
- New orchestrator leader loads
SAGA-9f3a, seesDEBITED, retries credit with idempotency keySAGA-9f3a:CREDIT. - Shard 9 applies credit once; saga moves to
CREDITED, thenCOMPLETED.
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.
- Orchestrator drives ordered steps with a durable state machine.
- Every forward step needs a compensating action or safe retry semantics.
- Participant idempotency keys prevent double-debit on orchestrator retry.
Test yourself
Answer these before moving on — recall is what makes it stick.