Table of Contents

    2PC

    SYSTEM DESIGN • CHAPTER 15.6

    Two-Phase Commit

    Understand how atomic commitment works across independent systems, why the protocol blocks when a coordinator fails, and why its blocking behaviour is a structural property rather than an implementation flaw.

    Learning objective: By the end of this article, you will understand the prepare and commit phases, the irrevocable nature of a vote, why 2PC blocks, the in-doubt window, how three-phase commit and consensus-backed coordinators address it, and when to avoid distributed transactions entirely.

    Prerequisites

    Recommended Knowledge

    • Transactions, atomicity, and durability
    • Write-ahead logging and recovery
    • Partial failure and timeout ambiguity
    • Consensus and quorum intersection
    • Leader election and coordinator failure
    • Locking and isolation levels
    • Idempotency and retry semantics
    • Latency budgets across services

    The Atomic Commitment Problem

    A single database makes atomicity straightforward: one log, one commit point, one recovery procedure. The problem changes entirely when a transaction spans several independent systems, each with its own log and its own failures.

    THE REQUIREMENT
    \[ \text{All commit} \;\vee\; \text{All abort} \]
    Property Statement
    Agreement No two participants reach different outcomes
    Validity Commit only if every participant voted yes
    Stability A decision once made is never reversed
    Termination Every participant eventually decides

    Simple Analogy

    Several parties sign a contract. Each first confirms they are willing and able, then a notary declares it binding. The difficulty arises when the notary collapses after collecting signatures but before announcing the outcome.

    Two-phase commit does not make distributed transactions safe. It makes their failure modes explicit, and those failure modes are the reason most systems avoid it.

    The Two Phases

    PROTOCOL FLOW
    Prepare Collect Votes Decide Commit or Abort
    1

    Prepare Phase

    The coordinator asks each participant whether it can commit. A participant answering yes must durably record that promise and hold every resource needed to honour it.

    2

    Decision Point

    If all votes are yes, the coordinator durably logs commit. If any vote is no, or any participant is unreachable, it logs abort. This log write is the moment of truth.

    3

    Commit Phase

    The coordinator informs every participant of the decision. Participants apply it and release their held resources, acknowledging completion.

    async function twoPhaseCommit(coordinator, participants, transaction) {
        const txId = coordinator.newTransactionId();
    
        // Phase one: solicit votes
        const votes = await Promise.allSettled(
            participants.map(p => p.prepare(txId, transaction.workFor(p.id)))
        );
    
        const allAgreed = votes.every(
            v => v.status === "fulfilled" && v.value.vote === "yes"
        );
    
        // Decision point: this log write is irreversible
        const decision = allAgreed ? "commit" : "abort";
        await coordinator.logDecision(txId, decision);
    
        // Phase two: propagate until every participant acknowledges
        const outcomes = await Promise.allSettled(
            participants.map(p =>
                decision === "commit" ? p.commit(txId) : p.abort(txId)
            )
        );
    
        const unacknowledged = outcomes
            .map((o, i) => ({ outcome: o, participant: participants[i] }))
            .filter(x => x.outcome.status !== "fulfilled");
    
        if (unacknowledged.length > 0) {
            await coordinator.scheduleRecovery(txId, decision, unacknowledged);
        }
    
        return { txId, decision, pendingRecovery: unacknowledged.length };
    }
    The decision log write is the commit point: Once it is durable, the transaction is committed regardless of what happens next. Phase two is delivery of a decision already made, not a further opportunity to reconsider.

    A Yes Vote Cannot Be Withdrawn

    The property that makes 2PC work is also what makes it painful. Voting yes surrenders autonomy: the participant can no longer decide its own outcome.

    After Voting Yes Participant Must
    Hold locks Retain them until the decision arrives
    Survive restart Recover the prepared state from its log
    Refuse conflicting work Block anything touching held resources
    Wait indefinitely Never unilaterally abort
    Guarantee commit capability Reserve every resource commit will require
    Why It Cannot Simply Time Out A prepared participant does not know whether the coordinator decided commit before failing. Aborting unilaterally risks diverging from participants that already committed.
    THE LOCK HOLD WINDOW
    \[ T_{\text{locked}} = T_{\text{prepare}} + T_{\text{decide}} + T_{\text{deliver}} + T_{\text{recovery}} \]

    The Blocking Problem

    This is the defining weakness of 2PC and the reason it is classified as a blocking protocol. It is not a bug to be fixed but a consequence of the protocol's structure.

    THE IN-DOUBT WINDOW
    Voted Yes Coordinator Fails Blocked Indefinitely
    Failure Timing Participant State Resolution
    Before prepare sent Nothing prepared Safe to abort locally
    During voting Some prepared, some not Unprepared abort, prepared block
    After votes, before decision logged All prepared, no decision exists Blocked until coordinator recovers
    After decision logged Decision exists but undelivered Blocked until coordinator recovers
    During phase two Some informed, some not Uninformed block, others proceed
    THE STRUCTURAL LIMIT
    A prepared participant cannot determine the outcome by asking peers, because a peer that has not heard either knows nothing more, and a peer that committed may be unreachable.

    Cascading Impact

    Stage Effect
    Locks held Rows or ranges become inaccessible
    Conflicting transactions queue Latency rises on unrelated work
    Connection pool saturates Waiting transactions occupy connections
    Service becomes unresponsive Failure spreads beyond the transaction
    Callers time out and retry Additional load on a degraded system
    One Transaction Can Take Down a Service A single in-doubt transaction holding a contended lock can cascade into a full outage, because every subsequent transaction touching that data blocks behind it.

    Durable State Requirements

    Recovery depends entirely on what was written to stable storage before each failure. Both sides must log at specific points.

    CREATE TABLE coordinator_log (
        transaction_id  VARCHAR(80) PRIMARY KEY,
        decision        VARCHAR(10) NOT NULL,
        participants    JSONB       NOT NULL,
        decided_at      TIMESTAMP   NOT NULL,
        completed_at    TIMESTAMP   NULL
    );
    
    CREATE TABLE participant_log (
        transaction_id  VARCHAR(80) PRIMARY KEY,
        state           VARCHAR(20) NOT NULL,
        coordinator_id  VARCHAR(80) NOT NULL,
        redo_payload    JSONB       NOT NULL,
        prepared_at     TIMESTAMP   NOT NULL,
        resolved_at     TIMESTAMP   NULL
    );
    
    -- Recovery: find transactions stuck in doubt
    SELECT transaction_id,
           coordinator_id,
           prepared_at,
           CURRENT_TIMESTAMP - prepared_at AS in_doubt_duration
    FROM participant_log
    WHERE state = 'prepared'
      AND resolved_at IS NULL
    ORDER BY prepared_at;
    Party Must Log Before Purpose
    Participant Replying yes Restart must honour the promise
    Coordinator Sending the decision Restart must repeat the same decision
    Participant Acknowledging completion Outcome survives a later failure
    Coordinator Forgetting the transaction Only after all acknowledgements
    Every phase requires a durable write: This is why 2PC is slow. A distributed transaction costs multiple synchronous disk flushes plus several network round trips before anyone sees a result.

    Recovery

    async function recoverParticipant(participant) {
        const inDoubt = await participant.log.findPrepared();
    
        for (const tx of inDoubt) {
            const coordinator = await participant.locate(tx.coordinatorId);
    
            if (!coordinator.reachable) {
                participant.metrics.recordBlocked(tx.transactionId);
                continue;
            }
    
            const outcome = await coordinator.queryOutcome(tx.transactionId);
    
            switch (outcome.decision) {
                case "commit":
                    await participant.commit(tx.transactionId);
                    break;
    
                case "abort":
                    await participant.abort(tx.transactionId);
                    break;
    
                case "unknown":
                    // Coordinator has no record: it failed before deciding
                    await participant.abort(tx.transactionId);
                    break;
    
                default:
                    participant.metrics.recordBlocked(tx.transactionId);
            }
        }
    }
    Presumed Abort If the coordinator has no record of a transaction, it cannot have logged commit, so abort is safe. This optimization also means aborted transactions require no log entry at all.
    Optimization Assumption Saving
    Presumed abort No record means abort No logging for aborts
    Read-only participant No changes to commit Excluded from phase two
    Last participant Combine its vote with the decision One round trip removed
    Single participant Degenerates to local commit Protocol skipped entirely

    Heuristic Decisions

    When blocking becomes intolerable, an operator may force an outcome manually. This restores availability by abandoning the atomicity guarantee.

    Heuristic Damage If an operator forces abort while another participant already committed, the transaction has partially applied. The system is now inconsistent and no automatic process will detect it.
    Requirement Reason
    Record every heuristic decision Reconciliation depends on knowing it happened
    Alert immediately Inconsistency needs prompt investigation
    Compare against the true outcome later Detect whether damage actually occurred
    Define a compensation path Repair requires knowing what to undo

    Reducing the Blocking Window

    Approach Mechanism Effect
    Replicated coordinator Decision log behind consensus Coordinator failure survivable
    Coordinator failover Standby reads the log and resumes Recovery in seconds, not hours
    Three-phase commit Extra pre-commit round Non-blocking under limited assumptions
    Fewer participants Narrow the transaction boundary Less to block, faster decisions
    Shorter prepare phase Do work before preparing Locks held for less time
    Cohort timeout to abort Abort if the decision never arrives Unsafe, may diverge

    Consensus-Backed Coordinators

    REPLICATED DECISION LOG
    Votes Collected Decision to Consensus Log Any Replica Can Complete
    The Practical Modern Answer Placing the decision log behind a consensus group means no single coordinator failure blocks participants. This is how distributed databases make 2PC operationally viable.

    Three-Phase Commit

    Phase Purpose
    Can-commit Solicit votes without committing resources
    Pre-commit Inform everyone that all voted yes
    Do-commit Apply the change
    Why 3PC Is Rarely Used Its non-blocking property assumes bounded message delay and accurate failure detection. Under a network partition those assumptions break, and it can produce inconsistent outcomes, which is worse than blocking.

    The Cost

    LATENCY LOWER BOUND
    \[ L_{2PC} \geq 2 \times \max_i(RTT_i) + (n+1) \times T_{\text{fsync}} \]
    Cost Consequence
    Two round trips minimum Latency bounded by the slowest participant
    Durable write per participant Disk sync cost multiplied across systems
    Locks held across the network Contention scales with transaction duration
    Availability is multiplicative Any participant down blocks the transaction
    Coordinator is a dependency Its failure blocks all in-flight work
    AVAILABILITY COMPOUNDS
    \[ A_{\text{transaction}} = \prod_{i=1}^{n} A_i \]
    Every participant reduces availability: Five participants each available ninety-nine point nine percent of the time yield a transaction available roughly ninety-nine point five percent of the time, before accounting for the coordinator.

    Alternatives

    Most systems that reach for 2PC would be better served by an approach that avoids distributed atomicity altogether.

    Alternative Trade Suitable When
    Saga with compensation Atomicity for availability Steps can be semantically undone
    Transactional outbox Immediate consistency for eventual One database plus messaging
    Redesigned boundaries Service purity for locality Data can live together
    Idempotent retry Atomicity for repeatability Operations are naturally repeatable
    Reservation pattern Immediate for two-step confirmation Resources can be provisionally held
    Accept and reconcile Prevention for detection Divergence is detectable and rare
    The Boundary Question If two pieces of data must change atomically, that is strong evidence they belong in the same transactional boundary. A distributed transaction is often a symptom of a decomposition drawn in the wrong place.

    When 2PC Is Appropriate

    Reasonable Use

    • Atomicity is a hard requirement
    • Compensation is impossible or unsafe
    • Participants are few and co-located
    • Coordinator is consensus-backed
    • Transaction volume is modest
    • Within one distributed database

    Poor Fit

    • Participants span regions
    • High transaction throughput
    • External systems involved
    • Single coordinator instance
    • Availability outweighs atomicity
    • Steps are naturally compensatable

    Monitoring

    Signals Worth Tracking

    • In-doubt transaction count and age
    • Time from prepare to resolution
    • Abort rate by cause
    • Participants failing to vote
    • Locks held by prepared transactions
    • Coordinator log write latency
    • Recovery attempts and outcomes
    • Heuristic decisions taken
    • Transactions blocked past threshold
    • End-to-end commit latency percentiles
    Alert on in-doubt age, not merely on count. One transaction stuck for an hour is more serious than a hundred resolving within milliseconds.

    Verification

    describe("two-phase commit", function () {
        it("aborts when any participant votes no", async function () {
            const result = await coordinator.execute(transaction, {
                participants: [okParticipant, refusingParticipant]
            });
    
            expect(result.decision).toBe("abort");
            expect(await okParticipant.applied(result.txId)).toBe(false);
        });
    
        it("holds locks while prepared", async function () {
            await participant.prepare(txId, work);
    
            const conflicting = participant.execute(conflictingWork);
    
            await expect(
                Promise.race([conflicting, timeout(1000)])
            ).resolves.toBe("timeout");
        });
    
        it("resumes after coordinator restart", async function () {
            await coordinator.prepareAll(txId, participants);
            await coordinator.logDecision(txId, "commit");
    
            await faultInjector.crashCoordinator();
            await coordinator.restart();
            await coordinator.runRecovery();
    
            for (const p of participants) {
                expect(await p.state(txId)).toBe("committed");
            }
        });
    
        it("blocks rather than diverging when coordinator is lost", async function () {
            await coordinator.prepareAll(txId, participants);
            await faultInjector.destroyCoordinator();
    
            for (const p of participants) {
                expect(await p.state(txId)).toBe("prepared");
            }
        });
    });

    Common Design Mistakes

    Weak Design

    • Single non-replicated coordinator
    • Participants timing out to abort
    • Not logging before voting yes
    • Including many participants unnecessarily
    • Performing slow work inside prepare
    • Spanning regions with one transaction
    • No monitoring of in-doubt state
    • Heuristic decisions left unrecorded
    • Never testing coordinator failure

    Strong Design

    • Consensus-backed decision log
    • Participants block rather than diverge
    • Durable logging at every phase boundary
    • Minimal participant count
    • Work completed before prepare
    • Participants kept close together
    • In-doubt age monitored and alerted
    • Heuristic decisions logged and reconciled
    • Coordinator failure drilled regularly

    System Design Interview Discussion

    Question What Your Answer Should Cover
    Why is 2PC blocking? Prepared participants cannot decide alone
    What is the commit point? The coordinator's durable decision write
    Why can a participant not abort? Others may already have committed
    How do you reduce blocking? Replicated coordinator and fast failover
    Why is 3PC rarely used? Partitions break its safety assumptions
    What does it cost? Round trips, syncs, compounding availability
    What would you use instead? Sagas, outbox, or redrawn boundaries
    When is it justified? Hard atomicity with no compensation path

    Design Checklist

    Production Checklist

    • Confirm atomicity is genuinely required
    • Evaluate sagas and outbox first
    • Back the coordinator with consensus
    • Log durably before voting and before deciding
    • Minimize the participant count
    • Keep participants within one region
    • Complete slow work before prepare
    • Never allow unilateral abort after voting yes
    • Implement automatic recovery on restart
    • Apply presumed abort to reduce logging
    • Bound in-doubt duration with alerting
    • Record heuristic decisions explicitly
    • Reconcile after any heuristic decision
    • Monitor lock contention from prepared transactions
    • Test coordinator crash, restart, and total loss

    Knowledge Check

    1

    What is the commit point in 2PC?

    The coordinator's durable write of the decision. After it, the outcome is fixed and phase two merely delivers it.

    2

    Why can a prepared participant not abort?

    It cannot tell whether the coordinator decided commit before failing, and other participants may already have committed.

    3

    What is presumed abort?

    Treating the absence of a coordinator record as abort, which is safe because no record means commit was never logged.

    4

    Why does availability compound?

    Every participant must be reachable, so the transaction's availability is the product of each participant's availability.

    5

    Why is three-phase commit uncommon?

    It assumes bounded delay and reliable failure detection. Partitions violate both, and it may then produce inconsistency rather than blocking.

    Summary

    Two-phase commit achieves atomic commitment across independent systems by separating agreement from execution. Participants vote in the prepare phase, the coordinator records a decision durably, and that decision is then delivered for application.

    The property that makes it work is that a yes vote is irrevocable. A prepared participant surrenders autonomy, holding locks and refusing conflicting work until told the outcome. This is also what makes the protocol blocking: if the coordinator fails after votes are collected, prepared participants cannot safely decide alone.

    Blocking is structural rather than a defect. A prepared participant that timed out and aborted might diverge from a peer that already committed. The practical remedy is a consensus-backed decision log, so no single coordinator failure leaves anyone in doubt.

    The costs are substantial: two round trips, a durable write per participant, locks held across the network, and availability that compounds downward with each participant added. Most systems should first consider sagas, the transactional outbox, or redrawing service boundaries so the atomic change lives in one place.

    Key Takeaway

    2PC trades availability for atomicity, and the trade is steeper than it appears. Every participant compounds the failure surface, locks are held across network round trips, and a coordinator failure blocks everyone. If you genuinely need it, replicate the decision log behind consensus, keep participants few and close, and monitor in-doubt age as a first-class signal.