Distributed systems
Part 1 of 6 · Two-Phase Commit — Protocol, Coordinator & ParticipantsTwo-Phase Commit — Protocol, Coordinator & Participants
Interview hub: 2PC coordinator, votes, forced logs, blocking vs 3PC, Paxos Commit, and sagas.
- 1Gist
- 2Maps
- 3Q&A
- 4Sandbox
Voice readout needs Web Speech Synthesis in this browser.
Question ladder
L1
Who decides commit or abort in classic 2PC?
Answer
The coordinator. The decision is its durable COMMIT or ABORT record, not a majority of participant guesses.
L2
What must be true before a participant sends YES?
Answer
A PREPARED record is force-logged. Undo and redo are durable, and locks stay held until the global decision.
L3
When does the uncertainty window open and close?
Answer
It opens when the participant sends YES and closes when that participant learns COMMIT or ABORT.
L4
A vote times out before any decision record exists. What is legal?
Answer
Abort. A missing vote is not YES. After COMMIT is durable, a timeout is a resend, not a new decision.
L5
Does 2PC make the distributed schedule serializable?
Answer
No. 2PC is atomic commit. Isolation is a separate protocol at each resource manager.
L6
What is the production symptom of a blocked prepare?
Answer
Prepared transactions aging in the database, lock waits climbing, vacuum or purge stuck, and a thread pool wedged on commit.
L7
When do you refuse 2PC in a design review?
Answer
The participant cannot durably prepare, the coordinator is one disk on a cross-region user path, or the product already accepts compensation. Offer one resource manager, or the existing saga series.
Failure modes
YES only in memory
A crash after the vote is sent and before the PREPARED record is stable lets one side commit and the other abort.
COMMIT sent before the force
A fast participant applies commit while the restarted coordinator finds no decision record and aborts.
Prepared participant guesses
Unilateral commit or abort inside the uncertainty window diverges from the coordinator log.
2PC sprayed across HTTP APIs
A third party that cannot prepare turns the protocol into a lock pileup or a fake atomic commit.
Misconceptions
A YES vote means the participant will probably commit.
YES is a contract. Both commit and abort are still possible, and the participant will not pick one alone.
3PC removes blocking and is safe to ship.
Under partition, 3PC can commit on one side and abort on the other. It is an interview contrast, not the production fix.
Raft replaces two-phase commit.
Raft can be the quorum log behind a coordinator. It does not by itself atomically commit two independent resource managers.
Interviewer traps
Retry the HTTP call and call that atomicity.
Retry is at-least-once. If the participant cannot prepare, say saga and point at the existing comparison page.
We use 2PC, so the system is serializable.
Separate atomic commit from isolation. Read-committed resource managers can still show distributed anomalies.
Design scenario
Same prompt for every reader.
Requirements
Both balances move or neither does. A coordinator crash must not invent a second outcome.
Traffic / scale
A few hundred transfers per second, two participants each.
Latency
The caller waits on two log forces and two round trips.
Consistency
Atomic commit. Each database keeps its own isolation level.
Availability
Prepared rows may stall until the decision log is reachable.
Failure assumptions
- The coordinator can die after both YES votes and before a decision record.
- A participant can restart with a PREPARED record and no decision.
Constraints
- Neither side is a third-party HTTP API.
- Do not replace the decision with a compensation.
Prompt
Debit one Postgres database and credit another in the same region. A partial transfer is an accounting incident.
API
Who is allowed to send PREPARE, COMMIT, and ABORT?
Data
What must be stable before YES, and before COMMIT?
Architecture
Where does the decision record live if one machine must not block transfers for minutes?
Two resource managers, one business action
Prefer
2PC when both sides can prepare
The coordinator force-logs one decision. Every updater has already promised it can still commit or still abort.
- Partial success is illegal, and both sides are resource managers you operate.
- The participant count stays small, often two.
- Lock hold across a round trip fits the latency budget.
Alternative
A saga, or one database
If a participant cannot prepare, or the product accepts pending and undo, atomic commit is the wrong contract.
- One database already has one commit record. Use it.
- Compensations are local commits plus undo, not a distributed rollback.
- 3PC does not fix a partition. It can decide twice.
Force the record before you announce it
A message that others will rely on leaves only after the matching log record is stable.
- 1
Name the cohort
The coordinator force-logs a start record that lists every participant. Recovery cannot finish a cohort it cannot name. - 2
Collect promises
A participant that can commit force-logs PREPARED, then sends YES. A participant that cannot commit aborts locally and sends NO. - 3
Write one decision
All required votes YES: force COMMIT, then send it. Any NO, or a timeout before that record exists: force ABORT, then send it. - 4
Apply, ack, forget
Prepared participants apply the decision, release locks, and ack. After the acks, the coordinator logs END and may forget the transaction.
Overview
Two-phase commit (2PC) is the classic atomic-commit protocol. A coordinator and a set of participants agree that a distributed transaction either commits everywhere or aborts everywhere. The price is an uncertainty window. After a participant force-logs a YES vote it cannot decide alone, so a dead coordinator blocks prepared participants and the locks they hold.
A single-node transaction is atomic because one log has one commit record. Across nodes there is no single commit record unless you build one. 2PC puts that record on the coordinator and makes every participant promise it can still commit or still abort when the record is revealed.
This hub is the protocol, the logging rule, and the choice against 3PC, Paxos Commit, a Raft-backed coordinator, and sagas. It does not rewrite the saga series.
Roles
| Role | Also called | Job |
|---|---|---|
| Coordinator | Transaction manager (TM) | Drive prepare, force the global decision, broadcast commit or abort, collect acks, forget the transaction |
| Participant | Resource manager (RM) | Do the work, vote, hold locks and undo or redo until the decision, apply it, ack |
| Application | AP in the X/Open model | Demarcates the transaction. It does not invent the outcome after a YES |
There is exactly one global decision. It is the coordinator's durable commit or abort record, not a majority of participant guesses and not whoever recovered first.
The protocol in one page
- The coordinator force-logs a start (or, under presumed commit, a collecting record naming participants) and sends PREPARE to every participant.
- Each participant that can commit force-logs PREPARED, with locks held and undo and redo durable, and only then sends YES. A participant that cannot commit logs abort, releases its locks, and sends NO. A read-only participant may send READ-ONLY and skip phase 2. That optimization has its own page.
- If every required vote is YES, the coordinator force-logs COMMIT and only then sends COMMIT. If any vote is NO, or a vote times out before a decision record exists, the coordinator force-logs ABORT and sends ABORT.
- A prepared participant applies that decision, logs it, releases locks, and acks.
- After acks, the coordinator logs END and may forget the transaction.
The rule interviews exist to test: a record that others will rely on is forced to stable storage before the message that announces it leaves the process. A YES that is only in memory is a lie after a crash. A COMMIT that is only in memory is a lie after a crash. Participants that acted on the lie and participants that did not will diverge.
Flow
- 1
1. Force log start
- next2. PREPARE every participant
- 2
2. PREPARE every participant
- next3. Force PREPARED, then YES
- 3
3. Force PREPARED, then YES
- next4. Force COMMIT or ABORT
- 4
4. Force COMMIT or ABORT
- next5. Apply decision, release locks
- 5
5. Apply decision, release locks
- next6. ACK, then log END
- 6
6. ACK, then log END
Lesson map
Two-phase commit
The coordinator is Preparing. Each store is locked and has voted YES.
Architecture. Coordinator Preparing. Store A Locked · YES. Store B Locked · YES
Select a node to see why it exists, or an edge to see the protocol, direction, effect, and consequence.
Mermaid export
flowchart TB coordinator["Coordinator Preparing"] store_a["Store A Locked YES"] store_b["Store B Locked YES"] coordinator -->|PREPARE| store_a coordinator -->|PREPARE| store_b store_a -->|YES| coordinator store_b -->|YES| coordinator coordinator -->|COMMIT| store_a coordinator -->|COMMIT| store_b store_a -->|ACK| coordinator store_b -->|ACK| coordinator store_a -->|Lock held| coordinator store_b -->|Lock held| coordinator coordinator -->|No decision| store_a coordinator -->|No decision| store_b store_a -->|Undo| store_b
What a YES vote really means
YES is not "I probably can." It is a contract:
- Local constraints passed, including disk space and serialization conflicts you are willing to lock through.
- Undo is durable, so abort is still possible.
- Redo is durable, so commit is still possible after a local crash.
- Locks, or the equivalent MVCC conflicts, stay held until the global decision.
- The participant will not unilaterally commit or abort while it is prepared.
NO is the opposite. The participant has already aborted locally and forgotten the transaction. It must not be asked to commit later. A timeout at the coordinator before the decision record is written is treated as NO: abort. A timeout after the decision record exists is not a new decision. Resend the recorded one.
Uncertainty window and blocking
The uncertainty window opens when a participant sends YES and closes when it learns COMMIT or ABORT. Inside that window:
- Commit locally, and the coordinator may have aborted because another participant voted NO or the coordinator timed out.
- Abort locally, and the coordinator may have committed. Its COMMIT record is durable and other participants may already have applied it.
So the prepared participant blocks. It waits, holds locks, and refuses new conflicting work. That is not an implementation bug. It is the atomic-commit safety condition with a single coordinator. Recovery cannot invent a decision the log does not contain. The crash matrix is on failures and recovery.
Flow
- 1
1. Both participants vote YES
- next2. Coordinator crashes
- 2
2. Coordinator crashes
- next3. No decision record on disk
- 3
3. No decision record on disk
- next4. Prepared participants block
- 4
4. Prepared participants block
- next5. Locks wait for the log
- 5
5. Locks wait for the log
The Failure view on the map above is this crash: both stores stay locked, and neither commit nor abort is legal.
Comparative choice
| Approach | Atomic? | Blocks if the coordinator is lost? | Partition | Happy-path cost | Use when |
|---|---|---|---|---|---|
| Classic 2PC | Yes, all or nothing | Yes, prepared RMs wait | Safe but stuck | 2 RTTs plus log forces | Few RMs, one admin domain, locks are acceptable |
| Presumed abort or commit | Same atomicity | Same blocking window | Same | Fewer forces or acks on one outcome | You already run 2PC and want the cheap path |
| 3PC (Skeen) | Tries to be non-blocking on crash | Reduced for fail-stop crashes | Can commit and abort | 3 phases | Almost never. Interview contrast only |
| Paxos Commit | Yes | No, if a majority of acceptors is up | Safe if quorums hold | More messages, disk delay can match 2PC | The coordinator itself must survive one failure |
| 2PC across Raft shards | Yes, if the decision is replicated | Blocks only if a quorum cannot learn it | Raft quorum rules | Shard consensus plus cross-shard 2PC | Multi-shard transactions inside one database |
| Sagas plus outbox | No. Compensations | No global lock hold | Business anomalies you design for | Per-step commits | Services you do not co-administer |
Picking 3PC to avoid blocking, without a partition story, is how you get two outcomes. Picking a saga when finance requires atomic debit and credit is how you get an incident. Picking 2PC across fifteen services on the public internet is how you get a lock pileup.
The long saga comparison already exists: Two-Phase Commit vs Sagas, inside Sagas and Distributed Transactions. Use those pages. This cluster does not rebuild them. Orchestration versus events is orchestration vs choreography. Deadlines are saga state machines. Poison steps are saga failure modes.
Raft is not a substitute lesson. When the transaction manager must survive one machine, a Raft group is the quorum log behind a coordinator. Election and log matching stay on Raft consensus.
When 2PC is the right tool
- Participants are resource managers you operate: shards of one database, or XA resource managers in one data center.
- The participant count is small, often two.
- Lock hold across a round trip is inside the latency budget.
- You can monitor prepared transactions and page a human before anyone makes a heuristic decision.
When it is the wrong tool
- A participant is a third-party HTTP API that cannot prepare.
- The coordinator is a single process with a single disk, and minutes of blocking are an outage.
- The business can accept "order placed, payment pending, compensate on failure." That is a saga, not a broken 2PC.
- You need progress through a coordinator crash. That is consensus on the decision, Paxos Commit or 2PC whose coordinator is a Raft group, not a hope that 3PC will survive a partition.
What this cluster covers
- Prepare and commit - votes, force-before-send, and the decision record.
- Failures and recovery - coordinator crash, participant crash, termination, heuristics.
- Presumed abort, presumed commit, and read-only votes - fewer forces, same uncertainty window.
- XA, Postgres, and MySQL -
PREPARE TRANSACTION,XA PREPARE, and heuristic decisions. - 2PC versus 3PC, consensus commit, and sagas - what to ship.
The ring returns here after the tradeoffs page.
Sandbox
The lists below stand in for stable logs. force means append, then the durability barrier, before any message that depends on that record. A crash is returning before the decision record exists.
Press Run. Snippets must be self-contained — no network, files, or native modules.
Press Run. Snippets must be self-contained — no network, files, or native modules.
Expected shape: the happy path prints commit, committed, committed. The crash path prints blocked, prepared, prepared. A NO vote aborts the cohort. The NO voter is already aborted before the coordinator decides, and the YES voter learns abort and leaves prepared.
Pitfalls
- Do not claim a retry of a half-finished HTTP call is atomic. It is at-least-once unless every step is idempotent and you already designed reconciliation.
- Do not treat read-committed inside each database as global serializability.
- Do not put an unreplicated coordinator log on the user-critical path across regions.
Interview Q&A
Why not just retry the HTTP call that failed halfway?
Answer
Retry gives you at-least-once effects, not atomicity. If the debit committed and the credit did not, a retry may debit twice unless every step is idempotent and you have a reconciliation protocol. 2PC is for the case where partial commit is illegal and both sides can prepare. If they cannot prepare, design a saga and say so. The existing comparison is Two-Phase Commit vs Sagas.
A participant voted YES and then its process crashed. On restart, what may it do?
Answer
It finds PREPARED in its log and no decision. It reclaims locks, stays prepared, and asks the coordinator, or runs a termination protocol. It must not commit or abort on its own. If it has no PREPARED record, it aborted. A later PREPARE gets NO.
The coordinator got one YES and then another participant timed out. Can it commit the YES voter?
Answer
No. A missing vote is not YES. It aborts, but only if it has not already force-logged COMMIT. If COMMIT is durable, the timeout is a delivery problem: resend COMMIT. The decision record wins over the network.
Why is the coordinator a single point of blocking even when every participant is healthy?
Answer
Because the only copy of the decision may sit on the coordinator. Healthy participants that voted YES still do not know which way the decision went if that disk is down. Replicating the coordinator, Paxos Commit or a Raft group as the transaction manager, removes that particular failure at the cost of more messages. Raft here is the quorum log behind a coordinator, taught on Raft consensus, not a substitute for this page.
Does READ COMMITTED inside each database make 2PC serializable globally?
Answer
No. 2PC gives atomic commit, not a global isolation level. Each resource manager can be read-committed and the distributed schedule can still show anomalies. Isolation is a separate protocol. Do not claim "we use 2PC so we are serializable."
Name the production symptom of a blocked 2PC.
Answer
Prepared transactions sitting in pg_prepared_xacts or XA RECOVER, lock waits climbing, vacuum or purge stuck behind an old xmin, and an application thread pool wedged on commit. The fix is to learn the real decision, not to kill sessions until you know whether a heuristic abort will diverge from the other resource manager.
When would you refuse 2PC in a design review?
Answer
The participant cannot durably prepare. The coordinator is unreplicated and the operation is on the user-critical path across regions. The participant count or lock footprint is large. Or the product explicitly allows compensation. Offer a saga and outbox, or a single-shard redesign, instead of sprinkling XA on HTTP handlers. Start from the saga hub.
What does force-before-send forbid?
Answer
Announcing a vote or a decision that is not yet stable. If you send COMMIT before the COMMIT record hits stable storage, a crash lets you recover as "no decision" and abort, while a fast participant already committed. Group commit is allowed. That transaction's record must be in a flushed batch before its message is sent.
Sketch two resource managers and a coordinator. Mark the first byte that is allowed to leave after each force. Then delete the coordinator disk after both YES votes and write down every outcome a participant is still allowed to choose alone. The legal list is empty.
Go deeper
- Jim Gray, Notes on Data Base Operating Systems covers logging and two-phase commit.
- Bernstein, Hadzilacos, and Goodman is the free atomic-commit reference.
- Gray and Lamport, Consensus on Transaction Commit treats classic 2PC as the F=0 case. The paper PDF is Paxos Commit, TODS 2006.
- Dale Skeen, Nonblocking Commit Protocols is the 3PC contrast.
- Mohan, Lindsay, and Obermarck on R-star and Adrian Colyer's notes cover presumed abort and presumed commit.
- Martin Fowler, Two-Phase Commit is the pattern summary.
- Engine surfaces: Postgres PREPARE TRANSACTION and MySQL XA.
- Lectures: CMU 15-445 and Kleppmann on two-phase commit.