High Level Design

Distributed Consensus with Paxos

Two-Phase Commit, how Paxos achieves consensus through majority agreement, and the multi-server write coordination protocol with log-line locking.

August 10, 2026

Prerequisites: Consistency levels, Quorum, Transaction isolation levels, CAP theorem, ACID properties.


Two-Phase Commit (2PC)#

PhaseDescriptionExample
Phase 1: PrepareCoordinator asks all participants if they can commit (vote yes/no)Coordinating DB changes across services
Phase 2: Commit/AbortIf all vote yes → commit; if any vote no → abort allDistributed transactions in banking

Paxos Basics#

Consensus#

Paxos achieves consensus through majority agreement — more than 50% of nodes must agree on a value.

Distributed Nature#

Operates in environments with potential failures and unreliable message delivery. Handles node crashes and network issues.

Use Cases#

  • Apache Zookeeper and Google Chubby
  • Distributed locking
  • Maintaining order in distributed logs

Why It Matters#

Distributed consensus is complex due to failures and timeouts. Paxos offers a fault-tolerant way to reach agreement even in such scenarios.


Distributed Consensus in a Multi-Server Setup#

Consider a system with:

  • 3 Application Servers: S1, S2, S3
  • 3 Database Servers: D1, D2, D3

The Protocol#

1. Consensus Requirement A write is complete and durable only if at least 2 of 3 databases accept it (majority rule).

2. Broadcasting When any app server (e.g., S1) receives a write request, it broadcasts to all three databases in parallel.

3. Log Line Agreement Before writing, servers must agree on the exact log line (position) where data will be written — prevents accidental overwrites.

4. Synchronization Servers query the current log line from all three DBs. If 2 or more agree on a log line, it's selected.

5. Locking To prevent concurrent writes, a server locks the selected log line by becoming its owner.

6. Write Execution The owning server writes. The write is accepted only if ownership is recognized by a majority.

Two Phases#

PhaseAction
Phase 1: LockingServers request and lock the log line by claiming ownership
Phase 2: WritingThe server that locked the line performs the write

Both log line ownership and the write must be acknowledged by a majority (2/3) of database servers. This ensures durability and consistency even under failures and network partitions.


Timestamp-Based Improvements#

Enhance the original algorithm by:

  • Prioritizing recency of operations
  • Preventing stale writes
  • Supporting retry mechanisms and flexible ownership
  • Maintaining consistency without tight coupling between lock ownership and writes