learninfra · Linux · Networking · Kubernetes · System Design · AI Infrastructure · Exam blueprints · Drills

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

  1. 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.
  2. 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
  3. 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.
  4. 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.
  5. 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.
  6. 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.
  7. Recap & playground

    • Cheat sheet
    • Playground