Distributed systems
Part 1 of 6 · Sagas & Distributed TransactionsSagas & Distributed Transactions — Orchestration, Choreography & Compensations
Interview hub on sagas versus 2PC, orchestration versus choreography, compensations, deadlines, and reconciliation for multi-service business transactions.
- 1Gist
- 2Maps
- 3Q&A
- 4Sandbox
Voice readout needs Web Speech Synthesis in this browser.
Checkout spans inventory, payments, and shipping
Prefer
Local commits plus a saga
Each service commits its own short transaction. If a later step fails, earlier work is undone by a new business action: release the hold, void or refund, cancel the shipment.
- Locks stay inside one local transaction.
- Services keep their own schema and deploy.
- The outcome is eventual and observable, including intermediate states.
Alternative
One XA transaction across every store
Prepare holds locks until the coordinator commits or aborts. A coordinator blip becomes an availability incident, and many stores do not speak XA cleanly.
- Participants block while prepared.
- Cross-region round trips sit on the user path.
- Heuristic commit or abort still needs reconciliation.
What must be true if step N fails
Steps 1 through N-1 already committed in their own databases. There is no shared undo log.
- 1
Name the completed steps
The saga log lists which local commits succeeded. That list is the compensation plan. - 2
Stop the forward path
Do not start step N+1. Do not keep retrying a business reject. - 3
Undo in reverse order
Each undo is a new transaction. It must be idempotent under at-least-once delivery. - 4
Reach a terminal state
COMPLETED, FAILED after undo, or COMPENSATION_FAILED for a human or reconciler.
Overview
A business action often crosses services: reserve inventory, charge the card, then ship. Those services do not share BEGIN and COMMIT. Sagas sequence local transactions and pair each step with a compensation so partial failure still converges to a known business outcome.
That is eventual consistency with an explicit undo policy. It is not serializability across databases, and it is not a lock held for the whole checkout.
Interviews probe four decisions:
- Why 2PC / XA stalls under partition and operator pain.
- How semantic rollback (cancel a hold, refund) stays idempotent.
- When a coordinator is clearer than an event dance.
- How timeouts, retries, poison steps, and reconciliation attach to saga state.
Ask this out loud: if step N fails after steps 1 through N-1 have committed locally, what is already visible to other systems, and which new transactions put the business back to a safe state?
Consistency choices
| Approach | Consistency | Availability under partition | Prefer when |
|---|---|---|---|
| Single-DB ACID | Strong across the tables in that database | One failure domain | A monolith or modular monolith is still the right boundary |
| 2PC / XA | Atomic commit across resource managers | Coordinator and participants block | Rare: short, co-located, low churn |
| Saga, orchestration | Eventual, with compensations | High when undo and deadlines are real | Multi-service workflows with branching |
| Saga, choreography | Eventual, with local reactions | High, harder to observe | Stable event contracts and few steps |
| Outbox plus inbox | At-least-once delivery glue | High | Pair with a saga; it does not define multi-step undo |
Rule: Transactional Outbox & Inbox Patterns and Change Data Capture teach reliable event emission. This cluster teaches the multi-step business outcome and its undo. Use both. Do not treat CDC or the outbox as a substitute for compensations.
Raft is a different tool. Raft Consensus — Leader Election, Log Replication & Safety replicates one log inside a cluster. It does not make inventory and payments one atomic transaction. The next lesson draws that line.
Decision path
Flow
- 1
1. Multi-service business action
- next2. Strict cross-DB atomicity needed
- 2
2. Strict cross-DB atomicity needed
- next3. Shared DB or rare short 2PC
- 3
3. Shared DB or rare short 2PC
- next4. Else name who owns the flow
- 4
4. Else name who owns the flow
- next5. Orchestrate or choreograph
- 5
5. Orchestrate or choreograph
- next6. Write compensations and deadlines
- 6
6. Write compensations and deadlines
- next7. Idempotent steps plus outbox
- 7
7. Idempotent steps plus outbox
- next8. Timeouts retries poison path
- 8
8. Timeouts retries poison path
- next9. Reconcile sagas stuck past SLA
- 9
9. Reconcile sagas stuck past SLA
Lesson map
Local commits, then undo
Inventory and payments have committed locally. Shipping has not started. The saga log is the only list of what must be undone.
Architecture. Orchestrator RUNNING. Inventory Committed. Payments Captured. Shipping Not started
Select a node to see why it exists, or an edge to see the protocol, direction, effect, and consequence.
Mermaid export
flowchart TB orchestrator["Orchestrator RUNNING"] inventory["Inventory Committed"] payments["Payments Captured"] shipping["Shipping Not started"] orchestrator -->|Reserve| inventory orchestrator -->|Charge| payments orchestrator -->|Ship| shipping shipping -->|Step failed| orchestrator orchestrator -->|Refund| payments orchestrator -->|Release| inventory inventory -->|Outbox| orchestrator orchestrator -->|Not a saga| inventory
The map is the checkout, not a checklist you always finish. If the action fits in one database, stop at a local transaction. If you truly need cross-store atomicity and can pay the blocking cost, 2PC is a niche. Everyone else picks an owner for the flow, writes compensations before the first production charge, and budgets timeouts.
Flow
- 1
1. Reserve inventory locally
- next2. Charge the card locally
- 2
2. Charge the card locally
- next3. A later step fails
- 3
3. A later step fails
- next4. Compensate in reverse order
- 4
4. Compensate in reverse order
- next5. Terminal FAILED with an audit
- 5
5. Terminal FAILED with an audit
Rule of thumb: sagas for user-facing multi-service workflows. 2PC only for tightly coupled stores you can operate. A distributed lock held for the whole checkout is not a transaction.
What this cluster covers
- Two-Phase Commit vs Sagas — Why 2PC Breaks at Scale — prepare, commit, blocking, heuristic outcomes.
- Orchestration vs Choreography — Central Coordinator vs Event Dance — who owns the happy path.
- Compensating Transactions — Idempotent Undo & Semantic Rollback — refund, void, release, reversing ledger entries.
- Saga State Machines — Timeouts, Retries & Deadlines — durable states and budgets.
- Saga Failure Modes — Poison Steps, Partial Failure & Reconciliation — quarantine and invariant checks.
The cycle returns here after failure modes. Each page is a full lesson. This hub does not re-teach them.
A distributed lock is not a saga
Why it shows up on whiteboards. Serializing the world feels like a transaction. One lock, then every step, then unlock.
Why it fails.
- A lock is not an atomic commit across resources. The holder can crash after inventory committed and before payment committed.
- Lease expiry and clock skew allow a second holder. That family of bugs lives in Distributed Locks — Correctness, Leases & Fencing Tokens. Do not rebuild it here.
- Throughput collapses, and the failure story is “who holds the lock,” not “which business steps must undo.”
Better default: short local transactions, saga compensations, and idempotency keys. Use a lock only for a short critical section inside one service.
Orchestrated saga (run this)
No network and no databases. Each step mutates a context dict. A declined card throws, the saga compensates the reserve, and a duplicate conceptual undo is safe because release overwrites the same field.
Press Run. Snippets must be self-contained — no network, files, or native modules.
Production differs in three ways this toy hides: each do and undo is its own database transaction, the saga row is durable before the call returns, and the “step completed” event is published with the outbox, not with a second write after commit.
Choreographed reactions (run this)
Services do not call each other. They react. An order id ending in X fails payment and releases stock. There is still no central brain in the handlers.
Press Run. Snippets must be self-contained — no network, files, or native modules.
The release itself is a local transaction inside the compensation handler. In production that handler also writes an outbox row in the same commit. The orchestration lesson is where to decide whether this dance is enough.
Interview Q&A
What is a saga?
Answer
A sequence of local transactions. Each step has a compensating action so the business process can undo earlier work when a later step fails. The guarantee is eventual consistency of the business outcome, not ACID atomicity across services.
How is a saga different from two-phase commit?
Answer
2PC tries an atomic commit with prepare and commit phases and can block participants that have voted yes. A saga accepts intermediate states and compensates. Sagas fit autonomous services. 2PC fits rare, tightly coupled stores. Depth is the next lesson.
Orchestration or choreography?
Answer
Orchestration: a coordinator drives steps and stores saga state. Choreography: services react to events and nobody holds the whole script. Orchestration is easier to debug. Choreography avoids a central workflow service and can turn into event spaghetti. Choose from branching and ownership, not from fashion.
Are compensations the same as a database rollback?
Answer
No. The local commits already happened and other systems may have seen them. Compensation is a new business transaction: refund, cancel reservation, reversing ledger entry. It must be idempotent, and it may only achieve semantic undo. You cannot unsend an email.
Where do the outbox and inbox fit?
Answer
They stop you from losing the “step completed” or “step undone” event after a local commit, and they stop a duplicate delivery from running the step twice. They do not choose saga topology or define the undo. Pair them with this cluster. The lesson is Transactional Outbox & Inbox Patterns.
What do timeouts change?
Answer
Every step and the whole saga need a deadline. A timeout is not an automatic undo. You need an explicit compensate-or-reconcile policy. Budget propagation and jitter live on Timeouts, Budgets & Deadline Propagation and Retry Storms, Backoff & Jitter. Saga states live on the state-machine lesson.
Can a saga guarantee exactly-once side effects?
Answer
Treat messaging as at-least-once and handlers as idempotent. A stable key of saga id, step id, and direction (do or undo) stops a retry from double-charging. Exactly-once is an effect you build, not a bus flag you flip.
When is a shared database still the right answer?
Answer
When team boundaries do not need independent schemas and deploys. A modular monolith with one ACID transaction beats a premature split plus a saga. Say that before you draw coordinators.
What is a poison step?
Answer
A step that will not succeed without a code or data fix: the same error on every retry. Quarantine it, dead-letter it, and reconcile. Infinite retry hides the bug and amplifies the outage. The failure-modes lesson is the catalog.
Which metric shows saga health?
Answer
Count of sagas stuck past SLA by state, compensation failure rate, and time-to-terminal for COMPLETED and FAILED. Money at risk in non-terminal payment states belongs on the same dashboard.
Pitfalls
Draw three boxes: inventory, payments, shipping. Commit inventory. Fail the charge. List what a customer can already observe, the exact undo call, the idempotency key, and the terminal saga state. Then say whether a coordinator or an event owns that undo, and why you rejected 2PC.