Distributed Transactions: 2PC & Sagas
One order touches four services with four databases. Two-phase commit makes it atomic and blocks when the coordinator dies; a saga never blocks and pays with compensations, idempotent steps and the outbox. Watch both fail, and learn which to choose.
An interactive System Design lesson: 19 steps, about 34 minutes, on a live simulation in your browser.
Alice buys one item for 30. The order service must take her money in payment-db, take one unit of stock in inventory-db and book a shipment in shipping-db. Each service owns its database; there is no single BEGIN ... COMMIT that spans them. Yet a half-done order (money taken, no shipment) is a support ticket.
Two-phase commit (2PC) makes the three one transaction. Phase 1: order, the coordinator, sends PREPARE; each participant does its work and ends with Postgres's PREPARE TRANSACTION 'tx-1' instead of COMMIT: the change is durable but invisible, its row stays locked, and the participant votes yes. Phase 2: all yes, so order writes COMMIT to its own log and sends COMMIT PREPARED 'tx-1'.
What you will learn
Two-phase commit
- One order, three databases: In 2PC the commit point is one line in the coordinator's log; everything before it can still be undone, everything after it must be finished.
- One participant says no: 2PC is unanimous: every participant can veto, and a single no rolls back everyone who said yes.
When the coordinator dies
- The coordinator dies after PREPARE: A participant that has voted yes is in doubt until it hears the decision: it can neither commit nor roll back alone, and it keeps its locks.
- Someone else needs the row: An in-doubt transaction blocks unrelated work: its locks outlive the coordinator and stay until someone delivers a decision.
- Drill: find the in-doubt ones
- Recovery: presumed abort: The coordinator writes its decision before announcing it, so on recovery "no decision in the log" safely means abort.
- The coordinator dies after deciding: Once COMMIT is in the coordinator's log the transaction has committed, even if no participant knows yet; recovery only delivers the news.
- Drill: resolve one by hand
Sagas: commit now, undo later
- An orchestrated saga: A saga is a chain of local transactions that each commit at once; isolation is replaced by visible states such as PENDING.
- Shipping says no: A saga cannot roll back; it moves forward by compensating each committed step, newest first.
- A compensation that keeps failing: Compensations must not fail for good: retry them, make them idempotent, and send the ones that still fail to people, never to the floor.
Retrying a step
- A lost reply, an idempotent step: Every saga step and compensation is retried, so every participant must be idempotent on the saga id.
- A naive participant charges twice: An at-least-once retry plus a participant that is not idempotent is a duplicate action, reported as a success.
Choreography and the outbox
- A choreographed saga: Orchestration puts the saga's flow in one service; choreography spreads it across events, which decouples services and hides the flow.
- Commit, then publish: A dual write (commit, then publish) loses the event whenever the process dies between the two; nothing will ever notice.
- The transactional outbox: The outbox makes the event part of the local transaction: commit both or neither, then publish at least once.
Choosing
- 2PC or saga?: Use 2PC for short transactions between databases you control, sagas for business flows across services, and one database when you can.
Recap & playground
- Cheat sheet
- Playground