Topic 10.3
Distributed Transactions: Two-Phase Commit and Sagas
In one line
Two-phase commit makes several resources commit atomically: a coordinator asks all participants to prepare, then to commit. It blocks if the coordinator fails after prepare. Sagas split a business transaction into local transactions with compensating actions, orchestrated by a coordinator or choreographed by events. They trade isolation for availability.
Think of it like this
A wedding. The officiant asks both people "do you?" (prepare); only after two yeses is the marriage declared (commit). If the officiant faints between the yeses and the declaration, everyone waits. A group holiday booked piece by piece, where you cancel the hotel if the flight fails, is a saga.
Key ideas
- 01
2PC phases: prepare (each participant writes everything durably, holds locks, votes yes/no), then commit or abort (the coordinator logs its decision, then tells everyone). A participant that voted yes can't decide alone; it's in doubt until the coordinator returns.
- 02
PostgreSQL supports it with
PREPARE TRANSACTION 'id'andCOMMIT PREPARED 'id'(max_prepared_transactions> 0). Orphaned prepared transactions hold locks and block VACUUM, so monitorpg_prepared_xacts. XA does the same for JTA. - 03
Distributed SQL databases (Spanner, CockroachDB, YugabyteDB) run 2PC internally over consensus groups, so participants and the coordinator are themselves replicated, which removes the blocking problem at the cost of latency.
- 04
Sagas: T1, T2, T3 with compensations C1, C2. If T3 fails, run C2 then C1. Compensations are semantic (refund, release stock), not rollbacks. Sagas lack isolation: other transactions may see intermediate states, so use pending states, semantic locks, and commutative updates.
- 05
Orchestration: a central saga orchestrator (a state machine stored in a database; Temporal, Camunda) sends commands and tracks progress. Choreography: services react to each other's events. Simpler for 2–3 steps, harder to follow as it grows. Both need idempotent steps and the outbox.
Code & diagrams
-- on each participant database
BEGIN;
UPDATE account SET balance = balance - 100 WHERE id = 1;
PREPARE TRANSACTION 'transfer-7781-db1'; -- durable, locks still held
-- coordinator logs decision "commit", then:
COMMIT PREPARED 'transfer-7781-db1';
-- in-doubt transactions left behind by a crashed coordinator
SELECT gid, prepared, owner FROM pg_prepared_xacts;Interview problem
The problem
Travel booking across three providers
Book a flight, hotel and car from three external providers as one trip. Providers offer reserve and cancel APIs, no 2PC. Design it so a user never ends up with a partial trip without being told, and nothing is charged twice.
When it breaks
Coordinator crash after participants prepared
What you see
Prepared transactions hold row locks indefinitely; dependent writes queue, and VACUUM can't advance past them.
Fix & prevent
Replicate or persist the coordinator's decision log and recover automatically; alert on old entries in pg_prepared_xacts; prefer sagas or a distributed SQL database.
Non-idempotent compensation
What you see
A retried refund issues two refunds.
Fix & prevent
Give every step and compensation an idempotency key recorded by the receiver; make compensations check state before acting.
Explain it without notes
Why is 2PC called a blocking protocol?
Practice
Write the compensations for an order saga with steps: reserve stock, create shipment label, charge payment, send confirmation email.
Trade-offs
- ↔
2PC gives atomicity and blocks on failure; sagas stay available and give up isolation, needing compensations and careful state design.
Done when you can
I can explain 2PC, its blocking failure, and design an orchestrated saga with idempotent compensations.