skip to content

A distributed system needs multiple services to agree on whether a single business transaction succeeds or fails, atomically. When would you reach for a two-phase-commit-style atomic commit protocol instead of running the operation through a full consensus protocol like Raft or Paxos, and what do you give up either way?

level: seniorimportance: should knowfreq 45%

answer

  1. 2PC = one-shot cross-service atomicity
  2. consensus = ongoing same-resource replication
  3. 2PC coordinator crash -> blocking (stuck holding locks)
  4. consensus majority quorum -> no single point of blocking
  5. Spanner/CockroachDB layer 2PC-like commit atop per-shard Raft

basics

~20 s

2PC is for getting several different services to agree on one all-or-nothing transaction right now; consensus (Raft/Paxos) is for keeping several copies of the same data in agreement over time, and it survives a coordinator crashing, while plain 2PC can freeze if its coordinator dies mid-transaction.

solid answer

~50 s

2PC solves a different problem than Raft/Paxos: it coordinates several distinct, independently-owned resources (e.g., two different databases, or a payment service and an inventory service) into one atomic all-or-nothing outcome for a single transaction, using a single coordinator to collect votes and broadcast the final decision. Raft/Paxos solve replicating one piece of state (a log) across multiple copies of the same logical resource so the group survives individual node failures. The key trade-off: classic 2PC's coordinator is a single point of failure - if it crashes after participants vote 'prepared' but before broadcasting the decision, participants are stuck holding locks, blocked, unable to unilaterally decide (the 'blocking problem'). Consensus protocols avoid that by using quorums, so the group keeps deciding even if some nodes (including any single leader) fail. In practice, teams use 2PC (or sagas) to coordinate heterogeneous services/databases in one transaction, and use Raft/Paxos underneath a single replicated store or coordination service.

go deeper

for a junior

Should know 2PC is about several systems agreeing on one transaction and that if the coordinator can die and get everyone stuck, that's a weakness - without deep mechanism detail.

for a middle

Should articulate that Raft/Paxos are for replicating one resource fault-tolerantly over time, while 2PC is for one-shot atomicity across separate resources.

for a senior

Should explain the blocking problem precisely (participants stuck holding locks on coordinator crash) and why majority-quorum consensus avoids the equivalent single-point stall.

for a principal

Should discuss real architectures that layer both (Spanner/CockroachDB-style per-shard consensus plus cross-shard commit), and alternatives like sagas, with their own trade-offs (eventual consistency, compensating actions).

## Two protocols that sound alike Two-phase commit (2PC) and consensus protocols like Raft/Paxos solve superficially similar-sounding problems - 'get multiple parties to agree' - but they're aimed at different situations, and conflating them leads to picking the wrong tool. ## What 2PC is for 2PC's job is to make a single, one-shot transaction spanning multiple independent, heterogeneous resources atomic: all of them commit, or all of them abort, with nothing in between. 1. A coordinator collects a 'yes, I can commit' or 'no' vote from each participant. 2. Once every participant has said yes, the coordinator broadcasts the final commit decision (or abort, if any participant said no). This is the right tool when the participants are genuinely different systems that don't share a replicated log - e.g., a payment processor and an inventory database owned by different teams - and you need a single transaction outcome across both, rather than long-run replicated agreement on an evolving stream of operations. ## What Raft and Paxos are for Raft and Paxos, by contrast, aren't about coordinating an atomic transaction across different resources - they're about keeping many copies of the same logical resource (typically an append-only log, or a single value in classic Paxos) consistent as new entries are continuously proposed over time, while tolerating the failure of individual replicas, including the current leader. A consensus protocol doesn't ask 'should transaction X commit across these five different databases' - it asks 'what is the next entry in this one replicated log, and has the group durably agreed on it.' | Dimension | Atomic commit (2PC) | Consensus (Raft/Paxos) | |---|---|---| | Agreement scope | one transaction across separate, independently-owned resources | one replicated log across copies of the same resource | | Duration | one-shot | continuously, over the system's whole lifetime | | That one node crashes | the coordinator's crash can freeze a transaction | leader failure triggers automatic re-election | | Machinery | a coordinator and participants for the duration of one transaction | an ongoing cluster of replicas, leader election infrastructure, and log persistence | ## The trade-off when the coordinator dies The trade-off that most sharply separates the two is fault tolerance under coordinator/leader failure. Classic 2PC has a well-known **blocking problem**: after participants vote 'yes' (prepared) and are holding resources locked awaiting the final word, if the coordinator crashes before broadcasting commit-or-abort, participants are stuck - they can't unilaterally decide to commit (another participant might have voted no and never told them) or abort (another might already have been told to commit), so they must simply wait, holding locks, until the coordinator recovers. This is a genuine availability hazard: a single coordinator failure can freeze a transaction indefinitely. Consensus protocols are built specifically to avoid this class of problem: because agreement only ever requires a majority quorum, not unanimous participation, and because leader failure triggers automatic re-election rather than an indefinite stall, the group as a whole keeps making progress as long as a majority of replicas are healthy - no single node's crash, including the leader's, can freeze the system the way a 2PC coordinator's crash can. ## The flip side: scope and cost The flip side of that fault tolerance is scope and cost. - Consensus protocols assume a relatively small, fixed set of replicas of the same resource - they don't naturally generalize to 'agree across N independently-owned, heterogeneous services,' because those services don't share a common replicated log or state machine to converge on. - And running full consensus is heavier machinery than 2PC for a single simple transaction: it requires an ongoing cluster of replicas, leader election infrastructure, and log persistence, whereas 2PC just needs a coordinator and participants for the duration of one transaction. That's why 2PC (or its non-blocking cousins, and more commonly today, sagas with compensating actions) remains the practical choice for cross-service transactional coordination, while consensus protocols are the practical choice underneath a single replicated store, lock service, or metadata catalog. ## Where they show up together In production, the two often show up together rather than as a strict either/or: - A distributed database might use Raft internally to replicate each shard's write-ahead log for fault tolerance, while using a 2PC-like protocol across shards to make a single multi-shard transaction atomic - CockroachDB and Google Spanner both do variants of this, layering cross-shard atomic commit on top of per-shard Raft/Paxos replication. - Similarly, some systems harden the classic 2PC blocking problem by making the coordinator itself a Raft/Paxos-replicated role rather than a single physical process, so 'the coordinator' can survive individual machine failure - borrowing consensus's fault tolerance to patch 2PC's single-point-of-failure weakness rather than replacing 2PC's transaction semantics outright. ## The practical decision rule The practical decision rule for an engineer: - **Reach for an atomic-commit-style protocol** when you need one-shot, all-or-nothing agreement across genuinely separate systems or resources for a single transaction, and be honest that a naive implementation introduces a blocking risk you'll need to mitigate (timeouts, a replicated coordinator, or a saga-based redesign that avoids distributed locking altogether). - **Reach for a consensus protocol** when the real requirement is keeping multiple replicas of the same piece of state consistent and available despite individual node failures over the system's whole lifetime, not just for one transaction.

  • Why doesn't a team just replace 2PC entirely with Raft to get better fault tolerance for cross-service transactions?
    Raft replicates a single logical log across copies of the same resource; it doesn't natively express 'these five independently owned services must all commit or all abort together,' since those services don't share a state machine to converge on. Making cross-service atomicity fault-tolerant usually means making the 2PC coordinator itself Raft-replicated, or avoiding distributed transactions altogether with sagas, not swapping in plain consensus.
  • What's the practical symptom of the 2PC blocking problem in a live production system?
    Participants that voted 'prepared' hold their locks (and often database resources) open indefinitely while waiting for a coordinator that has crashed and not yet recovered, which can cascade into resource exhaustion or unrelated transactions timing out on those same locked rows - operators typically notice via a spike in long-held locks or transaction timeouts correlated with a coordinator outage.
  • How do sagas avoid the blocking problem that plain 2PC has?
    Sagas break a multi-step transaction into a sequence of local, independently-committed steps, each with a predefined compensating action to undo it if a later step fails, rather than holding all resources locked and waiting for a single atomic decision - so no participant ever blocks indefinitely waiting on a coordinator, at the cost of only eventual (not immediate) consistency and needing carefully designed, idempotent compensations.

2PC is like a wedding officiant asking both partners 'do you agree' and only pronouncing them married once both say yes - if the officiant collapses after both said yes but before the pronouncement, the couple is stuck in limbo, unable to declare themselves married or call it off. Consensus is more like a large jury that only needs a majority to reach a verdict - if a few jurors (even the jury foreman) step out, the rest can still reach and record a valid verdict.

saying these in an interview costs you the question

  • treats 2PC and Raft/Paxos as interchangeable solutions to the same problem
  • doesn't know 2PC has a coordinator single-point-of-failure/blocking hazard
  • thinks consensus protocols can directly coordinate atomic commit across unrelated heterogeneous services without adaptation
  • unaware that real systems (Spanner, CockroachDB) layer cross-shard commit atop per-shard consensus rather than picking just one
  • suggests 2PC is fault-tolerant to coordinator crashes without any additional replication of the coordinator

context