The problem: commit across two databases at once#
TL;DRthe 30-second version
- 2PC commits one transaction across several databases atomically: all of them commit, or all of them abort. One coordinator drives two rounds of messages to the participants (the databases doing the work).
- Phase 1, prepare: the coordinator asks every participant 'can you commit?'. Each does the work, durably records that it is ready, and votes yes or no. A yes is a binding promise. The participant holds its locks until it hears the outcome.
- Phase 2, commit: if every vote was yes, the coordinator durably records COMMIT (the point of no return) and tells everyone to commit. Any no, and it tells everyone to abort.
- Cost: two round trips, about 4n messages, and every participant holds its locks for both rounds. The slowest participant sets the pace for all of them.
- Fatal flaw: 2PC blocks. If the coordinator crashes after the votes are in but before it announces the decision, the participants are stuck holding locks until it recovers. Spanner's answer is to replicate the coordinator's decision with consensus. A Saga gives up atomicity instead.
Picture two PostgreSQL databases at a bank. One holds savings balances, the other holds checking. Moving $50 from savings to checking means two writes: subtract 50 in savings, add 50 in checking. They live on separate machines, with separate logs and separate locks. Both writes must happen, or neither. If savings commits the debit and checking never commits the credit, $50 vanishes. The other way round, $50 appears from nothing.
The naive fix is to send commit to savings, then send commit to checking. It doesn't work. The instant savings commits, its change is permanent. Checking can still crash or reject its write a moment later, and now you have a half-done transfer you can't undo. There is no single instant where you know both will succeed before either one has made its change permanent.
The mechanism: two phases#
There is one coordinator and a set of participants. The coordinator is usually the node that began the transaction, here the application server running the transfer. Each participant is an independent database that holds some of the data. The coordinator drives the protocol. Participants only respond.
Phase 1 is prepare and vote. The coordinator sends PREPARE to every participant. Each participant does all the work it needs to be able to commit: it takes its locks, checks its constraints, and writes its changes. Then it force-writes a 'prepared' record to its write-ahead log (the on-disk log a database uses to survive its own crashes) and replies yes. If it can't, it replies no. A participant that voted yes is now prepared. It has given up its freedom to choose. It may no longer commit or abort on its own, and it holds its locks until it hears the outcome.
Phase 2 is commit or abort. The coordinator collects the votes. If every vote was yes, it force-writes a COMMIT record to its own log. That write is the commit point: the instant the outcome becomes final. Then it sends COMMIT to everyone. If even one participant voted no, or never answered, it logs ABORT and sends ABORT instead. Each participant applies the decision, releases its locks, and acknowledges. Once all acks are in, the coordinator forgets the transaction.
In PostgreSQL the participant's side is a real feature. PREPARE TRANSACTION stores a transaction in the voted-yes state, and a later COMMIT PREPARED or ROLLBACK PREPARED finishes it. In the $50 transfer, savings runs the debit (500 to 450), checks the balance stays positive, and prepares. Checking runs the credit (200 to 250) and prepares. Both vote yes, the coordinator logs COMMIT, and both run COMMIT PREPARED. A reader who hits a locked row mid-flight either waits or reads the last committed value. Nobody ever sees 450 next to 200.
PredictAt exactly which message does the transaction's fate become irreversible, the moment after which it will commit no matter what crashes next?
Hint: It is not when a participant votes yes, and not when a participant receives COMMIT.
When the coordinator force-writes the COMMIT record to its own log, before it sends a single COMMIT message. Once that record is durable, the outcome is decided. Even if the coordinator crashes right after, recovery reads that record and drives every participant to commit. A participant voting yes only makes commit possible. The coordinator's durable decision makes it certain.
Why everyone logs first, and what each crash does#
2PC's correctness rests on log records that are forced to disk (fsync'd, physically flushed, not just handed to the operating system) before the matching message is sent. The rule on both sides is the same: make your intention durable, then act. That is what lets a crashed node rejoin a transaction on restart instead of guessing.
- A participant force-logs 'prepared' before voting yes. If it crashes and restarts, it reads that record, knows it is bound, re-takes its locks, and asks the coordinator for the outcome.
- The coordinator force-logs COMMIT or ABORT before sending it. If it crashes mid-broadcast, recovery replays the log and re-sends the same decision to everyone. No two participants can ever see different answers.
- Each participant force-logs the final outcome before acking, so a second crash cannot undo a commit it already applied.
With those records in place, here is what each crash does.
- Participant crashes before voting: no harm. It promised nothing. The coordinator times out the missing vote and aborts. On restart the participant finds no 'prepared' record and rolls back its partial work.
- Participant crashes after voting yes: on restart it reads its 'prepared' record and asks the coordinator for the outcome. It must not decide on its own.
- Coordinator crashes before the commit point: recovery finds no decision record and treats the transaction as aborted. This is called presumed-abort, and it is the default in the XA standard. Participants that ask are told to abort.
- Coordinator crashes after the commit point: recovery reads the COMMIT record and re-sends COMMIT until every participant acks. The outcome was already sealed.
- Coordinator crashes after the votes are in, and stays down: the prepared participants block. This is the unavoidable case.
Cost: messages, round trips, latency, durability#
For n participants, a successful commit costs two round trips and a predictable amount of I/O. All of it sits on the critical path of the transaction.
| Resource | Cost for n participants | Why |
|---|---|---|
| Round trips | 2 (prepare→vote, then commit→ack) | One per phase; each is a full request/response to every participant |
| Messages | ~4n (2n requests + 2n replies) | PREPARE+vote in phase 1, COMMIT/ABORT+ack in phase 2 |
| Latency | ≈ 2 × (slowest participant round trip) | Each phase waits for the last reply; one slow node stalls the whole transaction |
| Forced log writes | 1 (coordinator) + 2 per participant | fsync on the decision, and on each participant's 'prepared' and 'outcome' records |
| Lock hold time | From local work until the phase-2 message arrives | Locks span both round trips, far longer than a local transaction |
Trade-offs: safety bought with liveness
2PC's central bargain is atomicity at the cost of availability. It guarantees that all participants agree on commit or abort (safety). It cannot guarantee that they always make progress (liveness).
- Pro: true atomic commit across different systems, with a simple protocol and broad standard support (XA/JTA).
- Pro: no partial commits, ever, even across crashes, thanks to forced logging.
- Con: blocking. A coordinator failure in the in-doubt window stalls prepared participants that are holding locks.
- Con: the coordinator is a single point of failure and a bottleneck. Latency tracks the slowest participant.
- Con: long lock hold times kill concurrency on contended data, and it scales poorly as the participant count grows.
2PC vs 3PC vs Paxos Commit vs Saga
| Protocol | Atomicity | Blocking? | Round trips | Partition-safe? | Used in practice |
|---|---|---|---|---|---|
| 2PC | Yes (all-or-nothing) | Blocks on coordinator crash | 2 | Safe but unavailable (CP) | XA/JTA, MSDTC, Postgres prepared txns |
| 3PC | Yes (no partitions) | Non-blocking (synchronous net only) | 3 | No, can split-brain | Almost never (textbook) |
| Paxos/Raft Commit | Yes | Non-blocking (survives node loss) | 2 + consensus | Safe and available | Spanner, CockroachDB, modern systems |
| Saga | No, eventual via compensation | Non-blocking | n local txns + compensations | Available (AP) | Long-lived microservice workflows |
Three-phase commit (3PC) inserts a pre-commit phase between vote and commit, so a participant that loses the coordinator can use a timeout to decide safely. That only works if a timeout reliably means 'the other side is dead'. A network partition breaks that. One group of participants times out into commit while the other times out into abort, and atomicity is gone. Real data center networks do partition, so 3PC stays in textbooks.
Paxos Commit (Gray & Lamport) keeps the two phases but replaces the single coordinator with a consensus group. The decision is a consensus value, so losing any one node no longer blocks the transaction. This is the modern answer, and it is what Spanner does.
A Saga is the odd one out. It does not give atomicity. Each step is an independent local commit, and if a later step fails, earlier steps are undone with compensating transactions. Other observers can briefly see partial state. You pick a Saga when you cannot afford to block and can live with that.
In the wild
- X/Open XA and JTA: the dominant standard. XA defines the interface between a transaction manager and resource managers (databases, message brokers). Java's JTA exposes it to applications. The xa_prepare, xa_commit, and xa_rollback calls are the 2PC phases.
- PostgreSQL and MySQL prepared transactions: PREPARE TRANSACTION writes the transaction in the voted-yes state, and COMMIT PREPARED or ROLLBACK PREPARED finishes it. Postgres warns that orphaned prepared transactions hold locks. That is the blocking problem as an operational fact.
- Microsoft MSDTC: the Distributed Transaction Coordinator on Windows, running 2PC across SQL Server, message queues, and other resource managers.
- Google Spanner: runs 2PC across its shards, but the coordinator and each participant is itself a Paxos group. The decision is consensus-replicated, so one machine failing no longer blocks.
- Apache Kafka transactions: the transaction coordinator uses a two-phase commit to write atomically to several partitions and the consumer-offsets topic. That is what gives exactly-once processing across a read-process-write loop.
Common misconceptions and gotchas
What happens if the coordinator dies right after everyone votes yes?
The prepared participants are stuck in-doubt. They have promised to commit and are holding locks, but they don't know the decision and may not invent one. They block until the coordinator recovers. On recovery it consults its log. If it had force-logged COMMIT, it re-sends COMMIT. If not, presumed-abort makes it abort. The participants ask, get the answer, finish, and release their locks.
Is 2PC the same as consensus (Paxos/Raft)?
No. Consensus gets a group to agree on a value while tolerating a minority of failures. It stays available as long as a majority is up. 2PC needs unanimity (every participant must vote yes) and cannot survive a coordinator failure without blocking. They solve different problems. That is why modern systems use consensus to replicate the 2PC coordinator's decision rather than replacing 2PC.
2PC vs Saga: when do I use which?
Use 2PC when you need genuine atomicity over a small, stable set of resources you control and can accept the blocking risk, like databases in one data center. Use a Saga for a long-lived, multi-service workflow where blocking is unacceptable. Each step commits locally and failures are undone by compensating transactions, and observers may briefly see partial results.
Why is 2PC called a 'blocking' protocol?
A prepared participant that loses the coordinator cannot safely make progress. It can neither commit nor abort on its own, since either choice might contradict the coordinator's actual decision. So it waits, holding locks, until contact is restored. The FLP result (Fischer, Lynch and Paterson) says no protocol can guarantee both safety and progress in an asynchronous network with even one crash. 2PC chooses safety.
In an interview
Lead with the guarantee and the cost. 2PC gives atomic commit across several participants, but it blocks: a coordinator failure between the phases leaves prepared participants stuck in-doubt, holding locks. Name the durability detail. Both sides force-log before acting, and the coordinator's COMMIT record is the commit point.
Be ready to walk the two phases, point at the commit point, and explain why a prepared participant cannot decide on its own. State the cost: 2 round trips, about 4n messages, latency set by the slowest participant, long lock holds. Then name the mitigations. Replicate the coordinator's decision with Raft or Paxos so it isn't a single point of failure. Know that 3PC is non-blocking only on a synchronous network and breaks under partitions. Know when a Saga is the better tool.
Two common follow-ups. '2PC vs consensus': 2PC needs unanimity and blocks on coordinator failure; consensus needs only a majority and stays available. '2PC vs Saga': atomic but blocking, versus available but only eventually consistent through compensation. Drop a real system to show depth: Spanner runs 2PC across Paxos groups, PostgreSQL exposes the participant role as PREPARE TRANSACTION, and Kafka uses 2PC for exactly-once.
References & further reading
- Gray & Lamport — Consensus on Transaction Commit (2006) — frames atomic commit as consensus; introduces Paxos Commit to remove the coordinator as a single point of failure
- Bernstein, Hadzilacos & Goodman — Concurrency Control and Recovery in Database Systems — the classic textbook; Ch. 7 covers 2PC, presumed-abort/commit, and recovery (free PDF)
- Corbett et al. — Spanner: Google's Globally-Distributed Database (OSDI 2012) — 2PC across Paxos groups so the coordinator decision is consensus-replicated
- Martin Kleppmann — Designing Data-Intensive Applications, Ch. 9 — the clearest book-length treatment of 2PC, its blocking flaw, and the alternatives