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.
Prerequisites: Consistency levels, Quorum, Transaction isolation levels, CAP theorem, ACID properties.
Two-Phase Commit (2PC)#
| Phase | Description | Example |
|---|---|---|
| Phase 1: Prepare | Coordinator asks all participants if they can commit (vote yes/no) | Coordinating DB changes across services |
| Phase 2: Commit/Abort | If all vote yes → commit; if any vote no → abort all | Distributed 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#
| Phase | Action |
|---|---|
| Phase 1: Locking | Servers request and lock the log line by claiming ownership |
| Phase 2: Writing | The 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