A distributed transaction needs every participant to reach the same decision, commit or abort, even when machines crash in the middle. Two-phase commit (2PC) achieves that at the price of blocking: a participant that has promised to commit can be stuck until the coordinator comes back. Three-phase commit (3PC) was designed to remove the block, and textbooks still teach it, yet almost no production system uses it.

This article treats both protocols as small state machines and then checks them mechanically. We enumerate every point at which the coordinator can crash, apply each protocol's termination rule, and count the outcomes, with and without network partitions. The script is about forty lines of Python, the results are printed below, and they show exactly why 3PC does not survive real networks. Implementation details of 2PC are covered in the hands-on 2PC guide and log-level optimisations in two-phase commit revisited.

Advertisement

Atomic commitment, stated precisely

The problem is called atomic commitment. Each participant votes yes (I can commit, and I have made my changes durable enough to do so) or no. The protocol must guarantee: agreement, no two participants decide differently; validity, commit is decided only if every participant voted yes, and if all vote yes and nothing fails the decision is commit; and termination, every correct participant eventually decides.

Validity is what makes this different from plain consensus. A consensus algorithm may choose any proposed value; atomic commit must choose abort if even one participant said no, so the decision is a function of every vote. That is why a participant that has voted yes cannot decide alone: it does not know the other votes, and it has given up the right to abort unilaterally.

2PC as a state machine, and where it blocks

In 2PC each participant starts in q. On a prepare request it either votes no and moves straight to a, or makes its writes durable, records a prepared entry and votes yes, moving to w. The coordinator collects votes, durably logs the decision, and sends it; participants move to c or a. A participant still in q that times out may abort on its own, because nobody can have committed without its vote.

Two-phase commit (participant)qinitialwvoted yesccommitaabortyesno / timeoutw is adjacent to both c and a: blocking stateThree-phase commit (participant)qinitialwvoted yesppre-commitccommitaabortno state is adjacent to both c and aCoordinator log2PC: decision record forces the outcome; 3PC: pre-commit round makes the decision known before anyone commits2PC under coordinator crashall participants in w: must wait (blocks)3PC under partitionone side sees p and commits, the other sees only w and aborts
Participant state machines for 2PC and 3PC. 3PC inserts a pre-commit state so that no state can lead directly to both commit and abort, which removes blocking when failures are crashes that timeouts detect reliably, and introduces disagreement when the network partitions.

The trouble is state w. From w a participant can move to c or to a, and which one depends on information only the coordinator holds. If the coordinator crashes after logging commit but before sending it, and every participant is in w, the survivors can ask each other and learn nothing: the vector of participant states, all w, is identical whether the coordinator decided commit or abort. Guessing either way could violate agreement when the coordinator recovers. They must wait, holding locks on every row the transaction touched.

Dale Skeen's analysis of commit protocols, published at SIGMOD in 1981, turned that observation into a rule: a protocol is non-blocking only if no state is adjacent to both a commit and an abort state, and no state from which the protocol cannot yet commit can coexist, in some execution, with a participant that has already committed. 2PC's w violates the first condition.

Advertisement

3PC: the pre-commit state and the termination rule

3PC splits the decision into two rounds. After all votes are yes, the coordinator sends pre-commit; participants move from w to p and acknowledge. Only after every acknowledgement does the coordinator send commit. Now w leads only to p or a, and p leads only to c. Seeing a p anywhere tells survivors that every participant voted yes; seeing only w tells them that nobody can have committed yet, because commit is sent only after all pre-commits are acknowledged.

When the coordinator fails, the surviving participants elect a new coordinator, which collects their states and applies a termination rule:

def terminate_3pc(states):          # states of the participants the new coordinator can reach
    if "c" in states: return "commit"
    if "a" in states or "q" in states: return "abort"   # q: someone never voted yes
    if "p" in states: return "commit"                   # first move everyone in w to p, then commit
    return "abort"                                       # all in w: nobody can have committed

def terminate_2pc(states):
    if "c" in states: return "commit"
    if "a" in states or "q" in states: return "abort"
    return "BLOCK"                                       # all in w: outcome unknowable

The all-w case is where the protocols differ. 2PC must block; 3PC may safely abort, because no participant can be in c while others are still in w. This reasoning is sound under two assumptions: a participant that does not answer has crashed rather than being slow, and every live participant can reach every other. In other words, a synchronous network with a perfect failure detector.

Checking both protocols exhaustively

The model below has three participants and one coordinator. Votes may be any combination of yes and no; the coordinator delivers each message to one participant at a time and may crash after any delivery. For every reachable vector of participant states we apply the termination rules defined above, first to all three participants together and then to every way of splitting them into two groups that cannot communicate.

from itertools import product
N = 3

def runs(proto):
    for votes in product("yn", repeat=N):
        st = ["q"] * N
        for i in range(N):                       # phase 1: votes arrive one at a time
            yield tuple(st)
            st[i] = "w" if votes[i] == "y" else "a"
        yield tuple(st)
        if "n" in votes:
            for i in range(N):                   # abort fans out one at a time
                if st[i] == "w":
                    st[i] = "a"; yield tuple(st)
            continue
        if proto == "3pc":
            for i in range(N):
                st[i] = "p"; yield tuple(st)      # pre-commit fans out
        for i in range(N):
            st[i] = "c"; yield tuple(st)          # commit fans out

def splits():
    for mask in range(1, 2 ** N - 1):
        yield ([i for i in range(N) if mask >> i & 1],
               [i for i in range(N) if not mask >> i & 1])

for proto, term in (("2pc", terminate_2pc), ("3pc", terminate_3pc)):
    seen = set(runs(proto))
    whole = {"BLOCK": 0, "ok": 0}
    part, ex = {"blocked": 0, "split-brain": 0, "ok": 0}, None
    for s in sorted(seen):
        whole["BLOCK" if term(s) == "BLOCK" else "ok"] += 1
        for g1, g2 in splits():
            d1, d2 = term([s[i] for i in g1]), term([s[i] for i in g2])
            if "BLOCK" in (d1, d2): part["blocked"] += 1
            elif d1 != d2: part["split-brain"] += 1; ex = ex or (s, g1, g2, d1, d2)
            else: part["ok"] += 1
    print(proto, "distinct crash states", len(seen), "connected", whole, "partitioned", part)
    if ex: print("  first split", ex)

Each state vector is a snapshot of the participants when the coordinator dies; split-brain means the two halves decided differently. Running it prints (the output is deterministic because states are visited in sorted order):

2pc distinct crash states 18 connected {'BLOCK': 1, 'ok': 17} partitioned {'blocked': 50, 'split-brain': 0, 'ok': 58}
3pc distinct crash states 21 connected {'BLOCK': 0, 'ok': 21} partitioned {'blocked': 0, 'split-brain': 8, 'ok': 118}
  first split (('p', 'p', 'w'), [0, 1], [2], 'commit', 'abort')

Read it carefully. With everyone connected, 2PC has exactly one bad state, all three in w, and it blocks there. 3PC has none: it always decides. Under partition, 2PC blocks often (50 of 108 state-and-split combinations) but never disagrees. 3PC never blocks, and in 8 of 126 combinations the two halves decide differently. The first split it reports is the textbook one: the coordinator delivered pre-commit to participants 0 and 1 and crashed before reaching participant 2; the side holding 0 and 1 sees p and commits, while participant 2, alone, sees only w and aborts. That is a violation of agreement, the one property a commit protocol exists to provide.

The model is deliberately small: participants never crash, the coordinator never recovers and partitions are clean two-way splits. Adding those behaviours only adds more ways to fail. Exhaustive exploration of a tiny model is still the cheapest way to find protocol bugs, and the same approach scales up with TLA+ or a hand-written explorer over your real state machine.

Worked example: a transfer split in two

Make it concrete. A payment service moves 100 from an account on shard A to one on shard B and writes an audit row on shard C, using 3PC with 2-second timeouts.

  1. All three shards prepare, lock their rows and vote yes. Each is in w.
  2. The coordinator sends pre-commit to A. A moves to p and acknowledges.
  3. A network fault isolates A from B and C, and the coordinator process is killed during deployment.
  4. A times out waiting for commit, elects itself, sees only its own state p, and commits: the debit is applied.
  5. B and C time out, elect B, see w and w, and abort: no credit, no audit row.
  6. The partition heals. 100 has left one account and arrived nowhere, and no log anywhere says the transaction failed.

Under 2PC the same timeline leaves all three shards in w holding locks until the coordinator restarts and reads its log. Payments stall for minutes and an on-call engineer is paged, but the books stay correct. This is the trade every real system has made: a liveness problem is visible and recoverable; an agreement violation is silent and permanent.

What fixes it: quorums and replicated coordinators

The partition problem has two known answers. The first keeps the three-phase shape but requires a majority to move: Skeen's quorum-based commit protocol and the later enhanced three-phase commit (E3PC) of Keidar and Dolev only let a group commit or abort if it holds a quorum, so a minority side waits instead of deciding. That restores agreement and blocks only when no quorum can talk.

The second, used in practice, is to keep 2PC and make the coordinator's decision itself highly available. If the decision is written to a replicated log, built on Paxos or Raft, a crashed coordinator is replaced by a new leader that reads the decision and finishes the job. Google Spanner runs 2PC across participant groups that are each Paxos-replicated, and Kafka's transaction coordinator keeps its state in a replicated internal topic. Gray and Lamport's Paxos Commit generalises this so each participant's vote is itself agreed by consensus. The window in which a transaction is stuck shrinks to a leader election.

Distributed databases add further tricks to make the common case cheap, such as committing in one round when participants can be inferred from writes; those are covered in distributed transactions in practice.

Operating commit protocols

If you run 2PC, directly through XA or prepared transactions or indirectly through a database that uses it, the operational job is to keep the blocking window short and visible.

  • Replicate or quickly restart the coordinator. Blocking lasts as long as the coordinator is unavailable, so its recovery time is your worst-case lock hold.
  • Alert on in-doubt transactions. Monitor prepared-but-undecided transactions by age (for PostgreSQL, rows in pg_prepared_xacts; for MySQL, XA RECOVER) and page when any is older than a few coordinator restart times.
  • Never resolve by guessing. Manual commit or rollback of an in-doubt transaction (a heuristic decision) can break agreement exactly as 3PC does under partition. Look up the coordinator's log first and record who decided.
  • Fence old coordinators. A coordinator that was paused, not dead, can wake up and send a stale decision. Use epochs or fencing tokens so participants reject messages from a deposed leader.
  • Test crash points. Inject a crash at every step of the protocol in staging, as the explorer does in miniature, and verify that recovery reaches the same outcome every time.

Failure modes

Failure2PC3PC
Coordinator crash, participants connectedBlocks if all are in wSurvivors decide correctly
Coordinator crash plus partitionMinority and majority may both block; never disagreeHalves can decide differently
Slow coordinator mistaken for deadParticipants keep waiting; safeBackup coordinator may decide while the original sends a different decision
Participant crashes after voting yesRecovers into w and asks the coordinatorRecovers into w or p; must join termination
Operator forces an outcomeHeuristic damage possibleSame
Message cost on commit2 round trips3 round trips

Trade-offs

2PC is safe under any pattern of crashes, delays and partitions, and pays with blocking. 3PC is live under crash failures in a synchronous network, pays an extra round trip on every commit, and gives up safety when its timing assumptions break, which on real networks they do. Quorum variants and consensus-replicated coordinators recover both properties at the cost of more machines and more messages. The order of preference for most systems is: one partition if the data model allows, a database that runs replicated 2PC internally, then sagas, and never classic 3PC.

What to do next

  1. Write out the participant state machine of every distributed transaction in your system, and mark which states can lead to both commit and abort.
  2. Run an exhaustive crash-point explorer, like the one above, against that state machine, including partitions.
  3. Make the coordinator's decision durable and replicated, or measure and accept its restart time as your maximum lock hold.
  4. Add an in-doubt transaction age metric and an alert, and a runbook that forbids guessing.
  5. Fence coordinator messages with epochs so a paused leader cannot send a stale decision.
  6. Where a transaction spans services rather than shards of one database, evaluate sagas or an outbox before reaching for 2PC.
Key takeaway: Two-phase commit blocks when the coordinator fails while participants are prepared, but never lets participants disagree. Three-phase commit removes blocking only if every slow node is truly dead and every live node can reach every other; under a partition an exhaustive check shows it committing on one side and aborting on the other. Production systems therefore keep 2PC and replicate the coordinator's decision with consensus, monitor in-doubt transactions, fence stale coordinators and never resolve by guessing.