Two-phase commit is usually taught as two rounds of messages: ask everyone to prepare, then tell everyone the outcome. That picture is correct and nearly useless for building or operating it, because the protocol's correctness and its cost live in the log records, not in the messages. Which records must be forced to disk, and in what order, decides whether a crash loses a transaction, how long locks are held, and how many synchronous writes each commit pays for.
This article revisits 2PC at that level. It walks through the basic protocol as a sequence of log writes, then the optimisation family that every serious implementation uses (presumed abort, presumed commit, read-only and one-phase commit), the precise reason the protocol blocks, why three-phase commit does not fix that on real networks, and how replicating the decision removes the blocking. It ends with a compact coordinator implementation and an operating checklist. The XA and PostgreSQL mechanics, and how distributed SQL databases make commit cheap, are covered in distributed transactions in practice; this page is about the protocol itself.
The protocol as log records
Atomic commit asks for one property: either every participant commits a transaction or none does, even if machines crash and restart. Each participant can always abort on its own until it votes yes. Voting yes gives that right away: a prepared participant has promised to commit if told to, and must keep that promise across crashes. That is why the vote must be durable before it is sent.
- The coordinator sends PREPARE to every participant.
- Each participant does whatever is needed to make commit possible later: writes its redo and undo information, keeps its locks, then forces a PREPARED record. Only then does it reply YES. If it cannot commit, it aborts locally and replies NO.
- If every vote is YES, the coordinator forces a COMMIT record listing the participants. That forced write is the commit point. Any NO or timeout means ABORT.
- The coordinator sends the decision. Each participant forces its COMMITTED or ABORTED record, releases locks and acknowledges.
- When every acknowledgement is in, the coordinator writes an END record, lazily, and may forget the transaction.
What it costs
The classic analysis, from Mohan, Lindsay and Obermarck's work on the R* system (ACM TODS, 1986), counts forced log writes and messages per transaction. The figures below are textbook accounting for the basic variants with N participants, all voting yes; real systems batch writes with group commit and pipeline messages, so treat them as relative, not absolute.
| Variant, outcome | Messages | Forced writes | Notes |
|---|---|---|---|
| Basic 2PC, commit | 4N | 2N + 1 | Prepare and commit forced at each participant, commit at the coordinator |
| Presumed abort, commit | 4N | 2N + 1 | Same as basic on the commit path |
| Presumed abort, abort | 3N | N | Coordinator forces nothing and collects no acks |
| Read-only participant (PA) | 2 | 0 | Votes READ-ONLY and drops out of phase two |
| Presumed commit, commit | 3N | N + 2 | No acks; participants need not force COMMITTED |
The commit latency the client sees is two message round trips plus the participant's forced PREPARED write and the coordinator's forced COMMIT write on the critical path. Phase two is off the critical path for the client but not for locks: participants hold locks until the decision arrives, so slow phase two delivery is directly visible as contention.
Presumed abort and presumed commit
Both optimisations exploit one idea: if the coordinator has no record of a transaction, the answer to "what happened?" can be a fixed default, so the cases that match the default need fewer writes and messages.
Presumed abort makes no information mean abort. The coordinator does not force anything for an aborting transaction and does not wait for abort acknowledgements; if a participant later asks about a transaction the coordinator has forgotten, the answer is abort. Participants need not force their ABORTED records either. Since aborts are the case where something already went wrong, making them cheap and quick to forget keeps the coordinator's log small. It is widely adopted, including in the X/Open DTP model.
Presumed commit makes no information mean commit, which makes the common case cheaper: no acknowledgements for commits and no forced COMMITTED record at participants. The catch is that the coordinator must now remember aborts, including transactions that were in progress when it crashed. So it forces a collecting record listing the participants before sending PREPARE, and on restart aborts anything that has a collecting record and no decision. That extra forced write is why presumed commit pays off mainly with many participants.
Read-only participants are the cheapest win of all. A participant that changed nothing votes READ-ONLY, releases its locks immediately and is left out of phase two. If every participant is read-only, there is no phase two. One-phase commit covers the single-participant case: the coordinator just tells it to commit, with no prepare. The last agent optimisation extends this: prepare everyone else, then hand the decision to the final participant, which commits in one step, making it the commit point.
Why it blocks, precisely
Consider participant A, prepared and waiting, when the coordinator crashes after sending some COMMIT messages. A cannot abort: the coordinator may already have committed and told B. A cannot commit: the coordinator may have aborted because some C voted no. A can ask other participants, and cooperative termination helps if anyone has heard the decision or anyone voted no. But if every reachable participant is also prepared and uninformed, the decision exists only in the coordinator's log, and A must wait for it to recover, holding locks the whole time.
This is not an implementation flaw. Atomic commit needs agreement on the outcome, and in 2PC a single process holds the decision. Three-phase commit adds a pre-commit round so that no participant can be left uncertain while another has committed, and with bounded message delays and no partitions it lets survivors decide. On a real network those assumptions fail: if the participants split into two groups that each see a different set of states, each can apply the termination rule and reach different outcomes. 3PC trades blocking for possible inconsistency under partition, which is why production systems do not use it.
Removing the single point: replicate the decision
The fix that works is to stop storing the decision in one place. Gray and Lamport's Paxos Commit ("Consensus on Transaction Commit", ACM TODS, 2006) runs a consensus instance for each participant's vote, all sharing a set of 2F+1 acceptors. The transaction commits if every instance chooses Prepared. Any F acceptors can fail without blocking, and with F = 0 and the coordinator as the only acceptor the protocol reduces to 2PC, which shows exactly what 2PC is: consensus with no fault tolerance on the decision.
Modern systems take the same step in a more practical form. Each participant is itself a replicated group, so a participant does not vanish with one machine, and the coordinator's decision is written into a replicated log or stored as a transaction record inside the database, so any node can read it. Spanner runs 2PC across Paxos groups; Percolator-style and parallel-commit designs put the record in the data. Paxos explains the consensus layer these designs rely on. The message pattern is still two-phase; what changed is that no single crash can hide the decision.
A coordinator you can reason about
The sketch below implements presumed abort. The interesting parts are the order of log writes, idempotent delivery and the answer to inquiries. Note the race in on_inquiry: under presumed abort, a coordinator that answers abort for a transaction still collecting votes must then never commit it, so the inquiry itself decides the outcome.
import threading
class Coordinator:
def __init__(self, log, rpc):
self.log, self.rpc = log, rpc # log.force() is a durable fsync'd append
self.decided = {} # txid -> "COMMIT" | "ABORT"
self.mu = threading.Lock()
def commit(self, txid, parts):
votes = self.rpc.broadcast(parts, "PREPARE", txid, timeout=2.0)
yes = all(votes.get(p) in ("YES", "READ_ONLY") for p in parts) # missing vote = NO
writers = [p for p, v in votes.items() if v == "YES"]
with self.mu:
if txid in self.decided: # an inquiry already forced abort
yes = False
if yes and writers:
self.log.force({"type": "COMMIT", "txid": txid, "parts": writers})
self.decided[txid] = "COMMIT" if yes else "ABORT" # PA: abort not logged
self.deliver(txid, self.decided[txid], writers)
return self.decided[txid]
def deliver(self, txid, decision, parts):
pending = set(parts)
while pending: # retry forever; participants are idempotent
acks = self.rpc.broadcast(pending, decision, txid, timeout=2.0)
pending -= {p for p, a in acks.items() if a == "ACK"}
if decision == "COMMIT":
self.log.append({"type": "END", "txid": txid}) # lazy, not forced
def on_inquiry(self, txid):
with self.mu:
# Undecided or unknown: presume abort and make it stick.
return self.decided.setdefault(txid, "ABORT")
def recover(self): # runs before inquiries are served
ended = {r["txid"] for r in self.log.scan() if r["type"] == "END"}
todo = [r for r in self.log.scan()
if r["type"] == "COMMIT" and r["txid"] not in ended]
for r in todo: # pass 1: load every decision
self.decided[r["txid"]] = "COMMIT"
for r in todo: # pass 2: redeliver (may block on a down participant)
self.deliver(r["txid"], "COMMIT", r["parts"])Participants mirror this. On PREPARE they write everything needed for commit, force PREPARED, then vote. On restart they scan for PREPARED records with no outcome, reacquire the locks those transactions held, and keep asking the coordinator. Applying COMMIT twice must be harmless, because delivery retries until it sees an acknowledgement. The same race applies at restart: a recovering coordinator must load every logged decision before it answers a single inquiry, or it will presume abort for a transaction its log says committed. A real implementation also forgets decided entries once they are safely ended, and stores the coordinator's log on the same durability footing as the databases it coordinates.
Worked example: a crash at the worst moment
Three shards hold an order, its inventory reservation and a payment row. All three vote YES. The coordinator forces COMMIT, sends the decision to shard 1, and crashes. Shard 1 commits and releases its locks. Shards 2 and 3 are prepared and in doubt, holding row locks on the inventory and payment rows; every checkout touching those rows now waits.
Shard 2 times out and asks shard 3, which knows nothing more. Neither can decide, so both wait. When the coordinator restarts, recover finds the COMMIT record with no END, redelivers COMMIT to all three (shard 1 acknowledges again, harmlessly), and the locks clear. The outage for those rows lasted exactly as long as the coordinator was down. That duration, not the message count, is the real cost of 2PC, and it is why the coordinator must restart fast or be replicated.
Operating 2PC
- Alert on in-doubt age. Any transaction prepared for longer than a few seconds is a lock-holding incident. Expose the count and oldest age per participant.
- Make the coordinator's state durable and highly available. An application process that keeps decisions in memory or on a local disk is a coordinator that can block every participant indefinitely.
- Avoid heuristic decisions. Many resource managers let an operator commit or roll back an in-doubt transaction by hand. If the guess disagrees with the real decision, atomicity is broken and someone has to reconcile data manually. Look up the coordinator's record first.
- Keep transactions short before PREPARE. Do the work, then prepare quickly; never call an external service while participants are prepared.
- Use globally unique transaction ids that encode the coordinator, so recovery and inquiries never mix two coordinators' transactions.
- Prepared transactions pin resources. In MVCC databases they hold locks and also hold back cleanup of old row versions, so a forgotten one slowly degrades the whole database. Group commit is what keeps the forced writes affordable under load.
When to use it, and when not
| Situation | Better choice |
|---|---|
| Several shards of one database product | The database's built-in distributed commit |
| Two databases you operate, rare cross-writes | 2PC with a durable coordinator and in-doubt alerting |
| Independent services owned by different teams | Sagas with compensations, see the saga pattern |
| Cross-shard writes on the hot path | Redesign keys so the common transaction stays on one shard |
The architecture of the basic protocol inside a single database engine is covered in two-phase commit architecture.
What to do next
- Draw your system's commit path as log writes, not messages, and mark every forced write and where it sits on the critical path.
- Identify where the coordinator's decision is stored and what happens to prepared participants if that store is unavailable for an hour.
- Confirm you use presumed abort and that read-only participants are excluded from phase two.
- Add in-doubt count and age metrics for every participant, with an alert at a threshold measured in seconds.
- Write a runbook for in-doubt transactions that starts with looking up the coordinator's decision, never with a heuristic guess.
- List cross-shard transactions by frequency and redesign the top ones to stay on one shard where possible.