2PC
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.
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.
| 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.
The Two Phases
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.
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.
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 };
}
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 |
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.
| 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 |
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 |
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 |
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);
}
}
}
| 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.
| 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
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 |
The Cost
| 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 |
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 |
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
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
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.
Why can a prepared participant not abort?
It cannot tell whether the coordinator decided commit before failing, and other participants may already have committed.
What is presumed abort?
Treating the absence of a coordinator record as abort, which is safe because no record means commit was never logged.
Why does availability compound?
Every participant must be reachable, so the transaction's availability is the product of each participant's availability.
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.