System Design Problem

Design a Distributed Consensus System (Raft / Paxos)

Commonly Asked By:GoogleCoreOSHashiCorpMicrosoft

Interview Setup

Interview Prompt

Design a distributed consensus system (like Raft or Paxos) that allows a cluster of nodes to agree on a sequence of values despite failures.

Clarifying Questions (ask before designing)

QuestionWhy it matters
How many nodes can fail (f)?Quorum = 2f+1. 5 nodes tolerate 2 failures. Drives cluster sizing.
Are we building consensus or using it (e.g., for a KV store)?
  • Building Raft is the interview
  • using etcd is the production answer.
What values are we agreeing on: config changes, client writes, or both?Membership changes require joint consensus, the hardest part of Raft.
Synchronous or asynchronous network model?
  • Raft assumes eventually-delivered messages
  • safety holds regardless of timing.

Scope

In scope

  • Leader election
  • Log replication
  • Safety (election safety, log matching, leader completeness)
  • Fault tolerance (crash failures, not Byzantine)

Out of scope (state explicitly)

  • Byzantine fault tolerance (PBFT)
  • Full production KV store on top (mention as application layer)
  • Performance optimization (batching, pipeline), mention briefly

Functional Requirements

Start by confirming cluster size and whether you should explain Raft or Paxos. Ask about leader election, log replication, membership changes, and snapshot compaction scope.

  • Implement a replicated state machine across N nodes (typically 3, 5, or 7)
  • Leader election: automatically elect a leader; re-elect on failure
  • Log replication: leader replicates log entries to followers in order
  • Linearizable reads and writes: clients see the most recent committed value
  • Membership changes: add/remove nodes without downtime (joint consensus)
  • Snapshot support: compact log by snapshotting state machine
  • Client request forwarding: followers redirect to leader

Non-Functional Requirements

Your interviewer will care most about safety vs liveness, tying directly to CAP Theorem & Consistency Models. Safety is non-negotiable: no split-brain commits. Liveness requires a majority quorum; articulate both before drawing nodes.

  • Safety: Never return incorrect results, even during partitions (no split-brain)
  • Liveness: System makes progress as long as majority of nodes are alive (N/2 + 1)
  • Latency: Write committed in 1 round-trip (leader → majority followers)
  • Durability: Committed entries never lost (persisted to stable storage before ACK)
  • Availability: Tolerate (N-1)/2 node failures (e.g., 2 of 5)
  • Deterministic: Same log → same state machine state on all nodes

Capacity Estimations

Run this math before you pick cluster size. Write throughput and log entry size tell you replication bandwidth; snapshot frequency bounds log replay time on cold starts.

MetricCalculationValue
Cluster sizeGiven (assumption documented in value)3 (dev), 5 (prod), 7 (highly critical)
Writes / secFrom Writes / day ÷ 86400 (+ peak factor in value)10K to 100K
Read / secFrom Read / day ÷ 86400 (+ peak factor in value)100K+ (with read-only followers)
Log entry sizeGiven~100 to 500 bytes
Log growth100K entries/sec x 200 bytes20 MB/sec
Snapshot intervalGiven (assumption documented in value)Every 10K entries or 100 MB
Leader election timeGiven (assumption documented in value)150 to 300ms (Raft election timeout)
Heartbeat intervalGiven (assumption documented in value)50 to 100ms

Architecture Diagram

In the room: default to Raft in the interview; leader election plus log replication is easier to whiteboard than Paxos phases.

Walk your interviewer through leader election and log replication as the two core loops, foundational to Replication, Failover & Leader Election. Consensus algorithms elect a single leader per term, replicate a totally ordered log to a majority of nodes, and apply committed entries to a state machine. Raft is the interview default (easier to explain than Paxos), but staff probes expect you to articulate safety (no split-brain commits) and liveness (leader election on failure). I draw one leader, N followers, and the majority commit quorum explicitly.

Loading...

Raft Leader Election

Normal operation:
  - Leader sends heartbeats (empty AppendEntries) every 100ms
  - Followers reset election timer on heartbeat receipt

Leader failure:
  1. Follower's election timer expires (random: 150-300ms)
  2. Follower increments term, transitions to CANDIDATE
  3. Votes for self, sends RequestVote to all peers
  4. Includes: candidate's term, last log index, last log term
  5. Other nodes vote YES if: candidate's term >= voter's term, voter hasn't voted,
     candidate's log is at least as up-to-date
  6. If candidate gets majority → becomes LEADER
  7. If another leader discovered (higher term) → revert to FOLLOWER

Log Replication

1. Client sends write request to leader
2. Leader appends entry to local log: {term, index, command}
3. Leader sends AppendEntries RPC to all followers
4. Follower checks log consistency at prev_log_index + prev_log_term
5. If mismatch → follower rejects; leader backtracks until logs match
6. When majority reply success → leader commits entry, applies to state machine
7. Commit index propagated to followers in next heartbeat

Component Deep Dives

Next we walk through each box on the diagram. I start with the Raft state machine because follower → candidate → leader transitions are what your interviewer will whiteboard first.

Raft State Machine (Follower → Candidate → Leader)

Walk through the three states and the safety invariant: committed entries must appear in all future leader logs.

FOLLOWER: Accept AppendEntries from leader; reset election timer on heartbeat.
CANDIDATE: On election timeout → increment term, vote for self, RequestVote RPC.
LEADER: Send heartbeats; replicate client writes via AppendEntries; commit when majority ack.

Safety invariant: if entry committed in term T, it appears in logs of all future leaders.

Log Matching & Consistency Check

AppendEntries carries prev_log_index + prev_log_term.
Follower rejects if its log doesn't match at that index, forcing the leader to decrement until logs align.
This overwrites conflicting uncommitted entries on minority partitions after heal.

Snapshotting & Log Compaction

Snapshotting prevents unbounded log growth on disk and accelerates restart recovery, linking to Write-Ahead Logging (WAL) and Data Durability.

Without snapshots: log grows unbounded → replay on restart takes minutes.
InstallSnapshot RPC: leader sends compressed state machine snapshot + last included index/term.
Follower discards log prefix before snapshot index. Trade-off: snapshot frequency vs replay time.

Kafka KRaft (Raft in Production)

Kafka KRaft Controller Log (Raft in Production: NOT a generic event bus)

Kafka metadata quorum (KRaft mode, since Kafka 3.3):
  - 3-5 controller nodes run Raft consensus on metadata log
  - Leader replicates: topic configs, partition assignments, broker registrations
  - Commit rule: metadata entry committed when majority of controllers acknowledge
  - Replaces ZooKeeper (which used ZAB ≈ Paxos)

Partition leader election (per topic partition):
  - Each partition has one leader broker + ISR (in-sync replicas)
  - Leader replicates log entries to followers via Fetch requests
  - Unclean leader election disabled for durability (min.insync.replicas=2)

Compared to application-level Raft (etcd, Consul):
  - KRaft: optimized for high-throughput append-only logs (millions msg/sec)
  - etcd Raft: optimized for small KV config stores (thousands ops/sec)
  - Both: single leader per term, majority quorum, linearizable writes

This problem designs consensus primitives; Kafka is an example consumer of Raft, not a user-facing event bus here

API Design

Inter-Node RPCs (Raft Protocol)

PROTOBUF
service RaftNode {
  rpc AppendEntries(AppendEntriesRequest) returns (AppendEntriesResponse);
  rpc RequestVote(RequestVoteRequest) returns (RequestVoteResponse);
  rpc InstallSnapshot(InstallSnapshotRequest) returns (InstallSnapshotResponse);
}

message AppendEntriesRequest {
  uint64 term = 1; string leader_id = 2;
  uint64 prev_log_index = 3; uint64 prev_log_term = 4;
  repeated LogEntry entries = 5; uint64 leader_commit = 6;
}

message LogEntry {
  uint64 term = 1; uint64 index = 2;
  bytes command = 3; EntryType type = 4;
}

Client API

PUT    /api/kv/{key}           → Write key-value (linearizable)
GET    /api/kv/{key}           → Read key-value (linearizable or stale)
DELETE /api/kv/{key}           → Delete key
GET    /api/cluster/status     → Cluster health, leader info
POST   /api/cluster/add_node   → Add node (membership change)
POST   /api/cluster/remove_node → Remove node

Common Error Responses

400 Bad Request: invalid input, missing required fields, or malformed JSON payload
401 Unauthorized: missing or invalid authentication token or API key
403 Forbidden: authenticated caller lacks required permissions for this resource
404 Not Found: requested resource ID does not exist
409 Conflict: duplicate write or version conflict, retry with a unique idempotency key
422 Unprocessable Entity: syntactically valid request failed semantic business validation
429 Too Many Requests: rate limit quota exceeded, client should honor Retry-After header
500 Internal Error: unexpected server failure, retry safely with an idempotency key
503 Service Unavailable: downstream dependency is unavailable or overloaded, retry with exponential backoff

Data Model

WAL (Write-Ahead Log): On-Disk Format

File: wal-segment-00042.log
Each entry (binary, protobuf-encoded):
┌──────────┬──────────┬──────────┬──────────────┬──────────┐
│ Length(4B)│ CRC32(4B)│ Term(8B) │ Index(8B)    │ Data(var)│
└──────────┴──────────┴──────────┴──────────────┴──────────┘

- CRC32 for corruption detection
- fsync after every append (safety) or batch fsync (performance)
- Segment rotation: new file every 64 MB

Persisted State

JSON
{
  "current_term": 42,
  "voted_for": "node-1",
  "log": [...],
  "snapshot": {
    "last_included_index": 10000,
    "last_included_term": 38,
    "state_machine_data": <binary>
  }
}

Fault Tolerance

Split Brain Prevention

Scenario: Network partition splits 5-node cluster into [A, B, C] and [D, E]

Partition [A, B, C] (3 nodes = majority):
  - Continues operating normally
  - All writes succeed

Partition [D, E] (2 nodes = minority):
  - Cannot elect leader (need 3 votes, only have 2)
  - No writes accepted

When partition heals:
  - D and E discover higher term from majority
  - Revert to followers, sync log from new leader

KEY GUARANTEE: At most ONE leader at any time
  - Because you need majority to win election
  - Two majorities always overlap by at least one node

Read Scalability

1. Follower Reads (stale, simplest): Read from any follower
2. Read Index (linearizable, no disk I/O): Leader records commit index → confirms leadership
3. Lease-Based Reads (linearizable, lowest latency): Leader uses time-limited lease
4. Follower Reads with Read Index: Follower asks leader for commit index, waits to catch up

Membership Changes (Joint Consensus)

Raft solution: Joint Consensus (two-phase)
  Phase 1: Leader logs C_old,new configuration entry
    - Both old AND new configs must agree (majority of old AND majority of new)
  Phase 2: Leader logs C_new configuration entry
    - Now only new config rules apply

Single-server changes (used in practice by etcd, CockroachDB):
  Only add or remove ONE node at a time (3→4→5 or 5→4→3)

Additional Considerations

Performance Optimizations

1. Batching: Buffer multiple client requests, replicate as one AppendEntries
2. Pipeline: Send AppendEntries for index N+1 before N is acknowledged
3. Parallel AppendEntries: Send to all followers simultaneously
4. Pre-vote: Before starting election, candidate asks "would you vote for me?"
5. Read optimization: Lease-based reads (no log append needed for reads)

Where Raft/Paxos Is Used

  • etcd (Kubernetes config store): Raft
  • ZooKeeper (Kafka, HBase): ZAB ≈ Paxos
  • CockroachDB, TiKV: Multi-Raft
  • Google Spanner: Multi-Paxos
  • Kafka (KRaft mode): Raft
  • Consul: Raft
  • MongoDB replication: Raft-like

Multi-Raft

Single Raft group: all data on 3-5 nodes → limited by single node's disk/CPU

Multi-Raft (CockroachDB, TiKV):
  - Data split into ranges/regions (e.g., key range [a-f] → one Raft group)
  - Each range has its own Raft group (3-5 replicas)
  - Hundreds of Raft groups per node
  
Benefits: Parallel writes across ranges, better load distribution
Challenge: Cross-range transactions need 2PC on top of Raft

Interview Walkthrough

  • 25-minute cut

    Skip arch50/arch75 depth unless staff.

    • Goal: totally ordered log across unreliable nodes (8 min)
    • Raft leader election and log replication loops (9 min)
    • Know AppendEntries prev_log_index alignment rules (8 min)
  • State the goal first: agree on a totally ordered log across unreliable nodes with one leader, a majority quorum, and linearizable writes.
  • Walk through Raft leader election: heartbeat timeout → increment term → RequestVote → majority wins; at most one leader per term.
  • Explain log replication via AppendEntries with prev_log_index + prev_log_term: follower rejects mismatches, leader backtracks until logs align.
  • Split-brain prevention: minority partition (2 of 5 nodes) cannot elect a leader, because majority overlap guarantees safety.
  • Cover snapshots for log compaction and joint consensus for membership changes: never change quorum size in a single step.
  • Mention real deployments: etcd (Kubernetes), Consul, CockroachDB Multi-Raft; interviewers want concrete examples, not just theory, often integrated into Service Discovery.
  • Common pitfall: starting with Paxos notation; Raft's explicit leader model is easier to explain and what most production systems actually implement.

Engineering Trade-offs

Raft vs Paxos vs ZAB Comparison

Consensus trades write throughput against fault tolerance, introducing odd-sized clusters and a leader bottleneck.

AspectRaftMulti-PaxosZAB (ZooKeeper)
UnderstandabilityDesigned to be simpleNotoriously complexMedium
LeaderStrong leaderProposer (weak leader)Leader
LogContiguous, no gapsCan have gaps, fill laterContiguous
Membership changeJoint consensusComplexAtomic broadcast
Used byetcd, CockroachDB, TiKV, Kafka (KRaft)Cassandra (LWT), SpannerZooKeeper, HBase
Rounds for commit1 (AppendEntries)2 (prepare + accept)1 (proposal)

Leader Election Trade-offs

Election timeout tuning:
  Too short (e.g., 50ms):
    ✓ Fast re-election → high availability
    ✗ Network hiccup causes false elections
  
  Too long (e.g., 5 seconds):
    ✓ Stable leadership, no false elections
    ✗ Leader fails → 5s unavailability for writes

  Typical sweet spot: 150-300ms election timeout, 50-100ms heartbeat

The "unavailability window" on leader failure:
  Worst case: ~500ms total write unavailability
  Accepted trade-off vs. allowing multiple leaders → data corruption

When Consensus is NOT Needed

Use cases that DO need consensus:
  - Distributed lock, Leader election, Distributed counter
  - Configuration store, Distributed transactions

Use cases that DON'T need consensus:
  - Read-heavy data: simple primary+replica replication
  - Best-effort counters: Redis INCR, no consensus
  - Event logs (Kafka): partition leader with ISR
  - Caches: Redis cluster with async replication

Rule: use consensus for metadata/coordination; use primary replication for data.

Linearizability vs Sequential vs Eventual

Linearizability (Raft): Every op appears instantaneously between invocation and response.
  Cost: all reads through leader (or verify with majority)

Sequential consistency: Ops appear in a sequence consistent with program order.
  Does not map to real time.

Eventual consistency: All nodes converge eventually. No timing guarantees.

For Raft: Linearizable reads via leader.
  Optimization: Lease-based reads (time-limited lease → local reads without consensus).
  Risk: Clock skew → stale reads. Use monotonic clocks + conservative lease timing.

💬Review

Help Us Improve

How helpful was this walkthrough?

Click a star to rate. We actively use this feedback to refine and update our system design content.

Placeholder
Optional but highly appreciated!

Discussion

Share your thoughts, ask questions, or help others.

Loading comments...