Consensus & CAP Theorem Interview Questions | JiQuest

add

#

Consensus & CAP Theorem

Advanced module · distributed systems theory

Consensus & CAP Theorem: the theory chapter underneath every distributed system you've already built.

CAP and PACELC, consistency models from linearizable to eventual, Raft's leader election and log replication mechanics (and how Paxos differs), quorum reads and writes via W+R>N, split-brain and the fencing tokens that prevent it, and vector clocks vs. wall-clock time — the vocabulary underneath Kafka's partition leaders, idempotent transfers, and every leader-based coordination pattern on this site.

2Consensus algorithms
4Consistency models
22Interview Q&A
Clientwrite(x=9) Leaderterm 4appends to its log Follower Areplicates entry Follower Breplicates entry majority ack (2 of 3) = committed only then does the leader tell the client it succeeded

Why this module exists, and what it assumes you already built

Every mini project on this site that talks about "the payment ledger's ordered locking," "Kafka's leader-per-partition model," or "idempotency keys" is quietly leaning on the theory in this module without ever naming it. This page names it. It won't teach you to prove Paxos correct from first principles — it will give you the mental models and the exact vocabulary an interviewer expects when a system-design conversation turns from "draw the boxes" to "what happens when a node dies mid-write."

What this module answers preciselyWhat CAP actually trades off (and the misreading almost everyone repeats), how a Raft cluster elects a leader and replicates a log, why W+R>N guarantees a read sees the latest write, and how split-brain is prevented in production.
Where you've already used this, unnamedAn idempotency key sidesteps the same ambiguity that leader failover creates. Ordered lock acquisition on a payment ledger is a single-node stand-in for linearizability. A Kafka partition's leader broker is a live instance of the leader-election problem.

Jump to a section

CAP theorem: what it actually claims

CAP gets reduced to "pick two of three" so often that the reduction is now more widely known than the theorem itself — and it's a misleading way to remember it. Consistency, Availability, and Partition tolerance aren't three independent knobs you set once at design time; Partition tolerance isn't really optional the moment you have more than one node talking over a real network, because real networks drop and delay packets. The only knob you actually control is what happens during an active partition: do you keep answering requests and risk returning stale or conflicting data (Availability), or do you refuse to answer until you can guarantee correctness (Consistency)? Outside of an active partition — which, for most systems, is the overwhelming majority of the time — there's no forced trade-off at all.

C A P Consistency Availability Partition tolerance CAsingle-node RDBMS CPZooKeeper, etcd, HBase APCassandra, DynamoDB refuses to answer rather than answer wrong not partition-tolerant at all

The "CA" corner is the honest asterisk: a single-node database can genuinely offer both consistency and availability because there's no second node to partition from. The instant you replicate across machines for durability or scale, CA quietly stops being available as an option, and the real design decision becomes CP vs. AP — refuse-when-uncertain vs. answer-when-uncertain.

PACELC: the more useful extension

CAP only describes behavior during an active partition (P), which for most systems is rare. PACELC adds the observation that matters the other 99% of the time: Else (E, no partition happening), a system still trades Latency (L) against Consistency (C). A system might choose Availability over Consistency during a partition, and independently choose to trade extra latency for stronger consistency during normal operation (or the reverse) — that second choice is the one you're actually making on nearly every request, since true network partitions are the uncommon case, not the common one.

SystemDuring a partition (CAP)Normal operation (PACELC)
etcd / ZooKeeper (Raft / ZAB)CP — unavailable if it can't reach a majorityPC — accepts higher write latency for linearizable reads
Cassandra (tunable)AP by default — always answers, possibly stalePL by default — low latency, weaker consistency, tunable per query
DynamoDB (on-demand)AP for eventually-consistent reads; CP-like for strongly-consistent readsCaller picks per request: PL (eventual) or PC (strong)

Consistency models: how strong a guarantee do you actually need?

"Consistency" isn't binary. Each model below is a different, precisely defined promise about what order operations appear to happen in, and each one costs more latency and availability than the one below it.

ModelGuaranteeWhen it's the right choice
LinearizabilityEvery operation appears to take effect instantaneously at one point between its start and end, in an order every client agrees matches real time.A distributed lock's "who holds it now," a bank balance check immediately before a transfer, leader election state itself.
Sequential consistencyAll processes see operations in the same order as each other, but that shared order need not match real, wall-clock time.A replicated cache where clients just need a single coherent view, not one tied to real-time recency.
Causal consistencyOperations that are causally related (a reply to a comment) are seen in that order by everyone; unrelated, concurrent operations can be seen in different orders by different clients.A comment thread or chat app — replies must follow what they reply to, but two unrelated top-level comments can render in a different order per viewer without anything looking broken.
Eventual consistencyNo ordering guarantee beyond: if writes stop, every replica eventually converges to the same value.A view counter, a "likes" count, DNS propagation — briefly stale is cheap, and every replica staying available is worth more than instant agreement.
The trade-off in one line Each step down this list buys you lower latency and higher availability, and costs you a guarantee about what "now" means. Picking the weakest model that your actual correctness requirements allow — not the strongest one available — is the skill; defaulting to linearizability everywhere is usually paying a latency tax the application never needed.

Raft: leader election and log replication, mechanically

Every Raft node is always in exactly one of three states — Follower, Candidate, or Leader — and every message carries a term, a monotonically increasing integer acting as the cluster's logical clock. At most one leader can exist per term. That single fact is what makes the whole algorithm mechanically checkable rather than something you have to trust.

Leader election

Followers expect periodic heartbeats (empty AppendEntries RPCs) from the current leader. If a randomized election timeout elapses with no heartbeat, a follower assumes the leader is gone, becomes a candidate, increments its term, votes for itself, and sends RequestVote to every other node. A node grants its vote at most once per term, and only if the candidate's log is at least as up-to-date as its own — that log-completeness check is what prevents a node with a stale, incomplete log from ever becoming leader. Win a majority of votes and you become leader for that term; split the vote (common with several simultaneous timeouts) and the term simply expires unresolved, and a new election starts with a fresh, randomized timeout, which is precisely what makes the randomization necessary — it keeps repeated split votes statistically rare instead of guaranteed.

Node A (Follower) Node B (Follower) Node C (Follower) election timeout: term 1→2, becomes Candidate RequestVote(term=2) RequestVote(term=2) VoteGranted(term=2) VoteGranted(term=2) Leader, term 2 2 of 3 votes = majority AppendEntries (heartbeat, term=2) AppendEntries (heartbeat, term=2) heartbeats reset B and C's election timeouts

Log replication and commitment

Once elected, the leader is the only node that accepts client writes. Each write is appended to the leader's local log as an uncommitted entry, then sent to followers via AppendEntries. An entry becomes committed — safe to apply to the state machine and safe to acknowledge to the client — only once a majority of nodes (including the leader) have persisted it. This is the same majority-overlap guarantee the quorum section below generalizes: because any two majorities of the cluster must share at least one node, whichever node wins the next election is guaranteed to have seen every entry that was ever committed, so committed history can never be silently lost by a leader change.

// Simplified Raft RPC shapes
RequestVote(term, candidateId, lastLogIndex, lastLogTerm) -> (term, voteGranted)
  // grant vote only if: term >= currentTerm, haven't voted this term yet,
  // and candidate's log is at least as up-to-date as the voter's own log

AppendEntries(term, leaderId, prevLogIndex, prevLogTerm, entries[], leaderCommit) -> (term, success)
  // success=false if prevLogIndex/prevLogTerm don't match -> forces the
  // leader to walk backward and repair a divergent follower log
  // entry becomes committed once leaderCommit reflects majority replication
The scenario that actually gets asked in interviews A 3-node cluster's leader partitions away from the other two. The isolated leader can no longer reach a majority, so a correct implementation stops acknowledging new writes rather than trusting its own clock. The majority side's election timeout fires, a new leader wins term 3, and keeps serving writes. When the partition heals, the old leader receives a message carrying term 3, recognizes its own term 2 is stale, steps down to follower, and truncates any of its uncommitted entries that conflict with the new leader's log — entries that were never acknowledged to a client are the only ones that can be safely discarded this way.

How Raft differs from Paxos, conceptually

Paxos came first and is provably correct, but it's notorious for being hard to reason about and even harder to turn into a working replicated log without inventing a substantial amount of unspecified engineering on top — which is precisely why Raft exists. Raft's own paper is literally titled "In Search of an Understandable Consensus Algorithm."

PaxosRaft
Core unit of agreementA single value, agreed on by Proposers and Acceptors via Prepare/Promise then Accept/Accepted rounds.A replicated log entry, agreed on by an explicit Leader that all writes flow through.
LeadershipNo leader required by the base protocol — any proposer can propose, which is part of why it's hard to reason about under contention.Exactly one leader per term, elected explicitly and mechanically before any writes happen.
Building a logRequires Multi-Paxos, a widely-implemented but loosely-specified extension, to agree on a sequence of values efficiently.Log replication is the protocol's native design, not an extension bolted on afterward.
Staleness detectionHandled through proposal numbers, compared implicitly across rounds.Handled explicitly through the term number carried on every single message.
Real-world adoptionChubby, Spanner's internals, older Kafka-adjacent systems in spirit.etcd, Consul, CockroachDB, TiKV, Kafka's KRaft metadata quorum.

You don't need to reproduce Paxos's proof in an interview. What you do need is the shape of the comparison above: Paxos is a more general, harder-to-implement-correctly primitive for agreeing on one value at a time; Raft trades some of that generality for a leader-centric design that maps directly onto "replicate a log," which is what almost every real system actually needs.

Quorum reads and writes: W + R > N

With N total replicas, a write is acknowledged once W of them confirm it, and a read is served once R of them respond (the client picks the freshest value among those R responses, typically by comparing version numbers). The guarantee that makes this useful: if W + R > N, every possible write quorum and every possible read quorum, drawn from the same N nodes, are mathematically guaranteed to overlap in at least one node — and that shared node is guaranteed to hold the latest write.

Write quorum W=3 — acknowledged v7 Read quorum R=3 — queried for v N1 N2 N3 N4 N5 N3 sits in both quorums — it must return v7 N=5, W=3, R=3, W+R=6>5 — every write quorum intersects every read quorum

The latency / availability trade-off in choosing W and R

Every unit you add to W makes writes slower and less available (more nodes must be reachable and healthy to succeed) but makes reads cheaper. Every unit you add to R does the reverse. The strict inequality also has a cost independent of which side you weight: it means neither reads nor writes can ever tolerate more than N−W (or N−R) nodes being unreachable, even if the remaining nodes are perfectly healthy — a strict quorum sacrifices some availability for the overlap guarantee. Dynamo-style systems sometimes relax this deliberately with a sloppy quorum (accepting writes on whichever nodes are reachable, even if they're not the "correct" N, and reconciling later via hinted handoff), trading the strict overlap guarantee for still higher availability during a partition — a real-world instance of the AP choice from the CAP section above, made concrete.

Choice (N=5)Write behaviorRead behavior
W=5, R=1Slow, fragile — any one node down blocks all writes.Fast, cheap — any single node has the latest data.
W=1, R=5Fast, cheap — returns after the first ack.Slow, fragile — must contact and merge every node.
W=3, R=3 (majority)Tolerates 2 nodes down.Tolerates 2 nodes down. The standard default.
The connection to Raft Raft's majority-commit rule is exactly a W=majority quorum write. A linearizable read in Raft (via a read-index or a lease-based read on the leader) is effectively R=1, but against a node whose data is guaranteed current precisely because W was already a majority — the same overlap math, expressed through a leader instead of an explicit per-request quorum vote.

Split-brain: how it happens, and how fencing tokens and leases prevent it

Split-brain is what happens when two nodes simultaneously believe they are the leader (or the lock holder) and both act on that belief. It's rarely caused by a bug in the election logic itself — it's almost always caused by a leader that hasn't yet noticed it lost leadership: a long GC pause, a slow disk, or a transient network blip can isolate a leader just long enough for the rest of the cluster to elect a replacement, while the original leader, unaware, wakes up and keeps acting exactly as before.

Old Leader (paused)resumes, still thinks it leads New Leaderelected after quorum lost old one Shared storagelast-accepted token = 6rejects anything < 6 write(token=5) write(token=6) rejected: 5 < 6 accepted: 6 ≥ 6

Leases: bounded-time leadership

A lease grants leadership for a fixed window that must be actively renewed before it expires. If a leader can't renew — because it's partitioned, paused, or just slow — the coordination service is free to hand leadership elsewhere once the lease lapses. The critical discipline this requires: a leader must verify it can still reach quorum before acting, not just trust that its local clock says the lease hasn't expired yet. A leader that's been paused by a stop-the-world GC pause has no way of knowing how much wall-clock time actually passed while it was frozen.

Fencing tokens: closing the gap a lease alone leaves open

A lease reduces the window for split-brain but doesn't close it completely — a leader can resume right at the boundary, believe its lease is still valid, and send a write before anything times out on the writing side. A fencing token closes that gap by moving the check to the resource being protected instead of trusting the leader's own judgment: every time leadership changes, the token increments, every write carries it, and the protected resource rejects any write whose token is lower than the highest it has already accepted — regardless of what the sender believes about its own status.

// storage layer enforcing a fencing token, independent of who "thinks" they're leader
if (incomingToken < storage.lastAcceptedToken) {
    reject("stale leader: token " + incomingToken + " < " + storage.lastAcceptedToken);
} else {
    storage.lastAcceptedToken = incomingToken;
    storage.apply(write);
}
The generalization worth naming in an interview A fencing token is a monotonically increasing version number enforced at the point of writing — the same core idea as optimistic locking on a database row, or the ordered, versioned writes a ledger uses to reject a stale concurrent update. Consensus produces the token; the resource enforces it.

Vector clocks and Lamport timestamps: ordering without a shared clock

Wall clocks on different machines drift, even with NTP running continuously — skew of tens to hundreds of milliseconds is routine, and a bad NTP correction or a paused VM can produce much larger jumps. Two events that are genuinely causally related can end up with timestamps in the wrong order simply because one machine's clock was running behind, which silently corrupts any "highest timestamp wins" conflict-resolution strategy.

Lamport timestamps

Each node keeps a single counter. It increments the counter on every local event, and on receiving a message, sets its counter to max(local, received) + 1. This guarantees that if event A happened-before event B (causally), A's Lamport timestamp is smaller than B's. What it doesn't give you is the reverse: two genuinely concurrent, unrelated events can still land on different Lamport values, so from the numbers alone you can't tell whether two events were causally connected or simply coincidentally ordered.

Vector clocks

A vector clock keeps one counter per node, as an array. Comparing two vector clocks tells you, exactly, one of three things: one happened-before the other, one happened-after the other, or the two are truly concurrent (neither entry dominates the other). That third case — detecting genuine concurrency, not just imposing an arbitrary order over it — is what a system like the original Dynamo used to recognize a real write conflict (two clients updating the same shopping cart from different replicas before either write propagated) instead of silently discarding one write in favor of the other based on an unreliable wall-clock comparison.

Concrete example Alice and Bob each update the same record on two different replicas that haven't yet synced. A wall-clock "last write wins" system picks whichever has the later timestamp and silently drops the other's update — a real data loss bug that looks like it worked. A vector-clock-aware system detects that neither update happened-before the other (they're concurrent), surfaces both versions, and lets the application (or the user) reconcile them explicitly, which is slower but doesn't lose data.

Where this shows up in production, and in this site's own projects

This material stops being abstract the moment you connect it to systems you already reach for.

ZooKeeperUses ZAB (ZooKeeper Atomic Broadcast), a leader-based protocol conceptually close to Raft — one leader totally orders writes, committed once a majority of the ensemble persists them, giving linearizable writes and (by default) sequentially consistent reads.
etcdImplements Raft directly, which is why Kubernetes' control plane can treat etcd as its linearizable source of truth — every write to cluster state goes through the same leader-election and majority-commit machinery covered above.
Kafka's partition leadersEach partition has exactly one leader broker; a producer with acks=all waits for the in-sync replica set (ISR) to acknowledge — structurally a quorum write. Older Kafka used ZooKeeper to elect partition leaders; KRaft replaces that with Kafka's own built-in Raft-based metadata quorum.
Idempotency keysLeader failover leaves a client unable to tell whether its in-flight write committed before the old leader died. An idempotency key makes the safe retry a no-op if the write already landed — the request-level answer to the same uncertainty consensus resolves at the cluster level.
The payment ledger's ordered lockingAcquiring account locks in a fixed global order to avoid deadlock is a single-node stand-in for linearizability — scale that ledger across partitions or regions, and "lock ordering" becomes "quorum writes and consensus," the exact material this module covers.
Distributed locks generallyAny "only one worker processes this job" pattern is a lease in disguise — and is exactly where a paused worker resuming after its lease expired can cause the split-brain scenario diagrammed above, unless the downstream resource enforces a fencing token.

Key design decisions and interview talking points

What does CAP theorem actually claim, precisely?

CAP says that when a network partition actually occurs between nodes that hold replicated data, a system can only choose one of two things during that partition: consistency (every node returns the latest acknowledged write, or an error) or availability (every request gets a non-error response, even if some node is unreachable). It is not a claim about normal operation — outside of an active partition, a well-built system can offer both, and most of the time nothing is partitioned at all.

Why is "pick two of three" a misleading way to remember CAP?

It implies Partition tolerance is an optional choice on the same footing as the other two, when in practice any system with more than one node over a real network must tolerate partitions happening — you don't get to opt out of the network dropping packets. The only real, ongoing decision is what to do when a partition happens: favor Consistency or favor Availability. Treating all three as equally optional obscures that P is effectively a given at multi-node scale.

What does PACELC add that CAP leaves out?

CAP only describes behavior during an actual partition, which is the uncommon case. PACELC adds that even during normal operation (Else), a system still trades Latency against Consistency — a choice made on nearly every request, not just during rare network failures. A system's PACELC classification (for example PA/EL for Cassandra, or PC/EC for etcd) describes its everyday behavior far more often than its CAP classification does.

Your service uses a 3-node Raft cluster and the leader partitions away from the other two — walk through what happens.

The isolated leader keeps sending heartbeats into the void and gets no acknowledgments, so it can no longer confirm a majority for any new entry — a correct implementation stops accepting writes once it can't reach quorum, rather than trusting its own clock or local state. Meanwhile the two majority-side nodes stop receiving heartbeats, an election timeout fires, one becomes a candidate, increments the term, and wins a majority vote (2 of 3) to become the new leader. When the partition heals, the old leader receives a message carrying the new, higher term, recognizes it is stale, steps down to follower, and truncates any of its own uncommitted log entries that conflict with the new leader's log.

What is a Raft term, and why does the algorithm need it?

A term is a monotonically increasing integer acting as a logical clock across the cluster — at most one leader can be elected per term, and every message carries the sender's term. It's what lets a node cheaply detect staleness: a message with a term lower than your own came from a leader or candidate that hasn't heard about a more recent election, and you reject it — the exact mechanism that resolves a split-brain leader once a partition heals.

Why does a Raft leader need acknowledgment from a majority before committing a log entry, not from all followers?

Requiring all followers would mean a single slow or down node blocks every write, destroying availability with no safety benefit. A majority is the minimum that still guarantees safety: any two majorities of the same cluster must overlap in at least one node, so whichever node wins the next election is guaranteed to include at least one node that saw every previously committed entry — that overlapping node is what carries committed history forward across a leader change.

How does Raft differ from Paxos at a conceptual level?

Paxos separately proves agreement on a single value through Proposers and Acceptors exchanging Prepare/Promise and Accept/Accepted messages, with no leader mandated by the base protocol — agreeing on a sequence of values (a replicated log) requires the extra, loosely-specified engineering of Multi-Paxos. Raft was designed to be more understandable: it bakes in a strong, exclusive leader from the start, models the whole problem as replicating a log rather than agreeing on abstract values, and makes staleness detection explicit through the term number on every message instead of leaving it to the implementer.

RaftPaxosConsensus

Derive why W + R > N guarantees a read sees the latest write.

With N total replicas, a write quorum touches W of them and a read quorum touches R of them. If W + R were less than or equal to N, you could pick a write quorum and a read quorum that share zero nodes, and the read would consult only replicas that never saw the write. Once W + R exceeds N, any W-sized set and any R-sized set drawn from the same N nodes are guaranteed by the pigeonhole principle to share at least one node, and that shared node holds the latest write — so comparing version numbers across the R responses is guaranteed to surface it.

You have N=5 replicas. What are the trade-offs between W=5,R=1 and W=1,R=5 and W=3,R=3?

W=5,R=1 makes reads maximally fast and cheap but writes wait on every replica — one slow or down node blocks all writes. W=1,R=5 is the mirror image: writes return instantly, but every read has to contact all five nodes and merge results, and a down node blocks reads instead. W=3,R=3 (majority) balances both directions — writes and reads each tolerate up to 2 unreachable nodes while still guaranteeing overlap (3+3=6>5), which is why majority quorums are the default in most production systems rather than either extreme.

What's a "sloppy quorum," and why would a system deliberately weaken the strict quorum guarantee?

A sloppy quorum accepts a write on whichever W nodes are actually reachable during a partition, even if they aren't the "correct," designated N nodes for that key — then reconciles via hinted handoff once connectivity is restored. It trades away the strict overlap guarantee (a read might briefly miss a write accepted on a substitute node) in exchange for the write succeeding at all during a partition, a deliberate, concrete instance of choosing Availability over Consistency, the AP side of CAP made operational.

Give a concrete example of a split-brain scenario in a leader-based system.

A leader holding a distributed lock experiences a long GC pause or network blip that isolates it just long enough for the coordination service to declare its lease expired and hand leadership to another node. The original leader, unaware anything happened, resumes exactly where it left off and writes to a shared resource believing it's still leader — meanwhile the new leader is also writing to the same resource, and without a mechanism to reject the stale leader's writes, both corrupt shared state simultaneously.

How does a lease prevent split-brain, and what's the danger of a leader trusting its own local clock?

A lease grants leadership for a bounded time window the leader must actively renew before expiry; failing to renew means the coordination service is free to grant leadership elsewhere. The danger is a leader assuming it still holds the lease just because its own clock says it hasn't expired — clock drift, a stop-the-world GC pause, or a scheduling delay can all make a leader's local sense of time diverge from the coordination service's, so a leader must verify it can still reach quorum, not just check a timer, before acting as leader.

What is a fencing token, and what specific failure does it prevent that a lease alone doesn't?

A fencing token is a monotonically increasing number handed out every time leadership or a lock is granted, carried on every downstream write. The protected resource rejects any write whose token is lower than the highest it has already accepted. This closes the exact gap a lease leaves open: a leader that paused right before writing, had its lease reassigned, then resumes moments later isn't late enough for a timeout to catch — but its token is stale, so the resource refuses it while the new leader's higher-token writes succeed.

Why can't you just use wall-clock timestamps to order events across machines?

Clocks on different machines drift and skew relative to each other even with NTP running — differences of tens to hundreds of milliseconds are routine, and larger jumps happen during NTP corrections or VM pauses. Two causally-related events (B happened because of A) can end up with B's timestamp earlier than A's simply because B's machine's clock ran slightly behind, silently producing an order that contradicts reality — exactly the bug that causes a "last write wins" system to quietly drop a newer write in favor of an older one that got a later timestamp.

What's the difference between a Lamport timestamp and a vector clock?

A Lamport timestamp is a single counter per node giving a total order consistent with causality — if A happened-before B, A's timestamp is guaranteed smaller — but not the reverse: two concurrent, unrelated events can still get different Lamport values, so the numbers alone can't tell you whether two events were causally related or just coincidentally ordered. A vector clock keeps one counter per node in an array, so comparing two vector clocks tells you exactly whether one happened-before, happened-after, or is truly concurrent with the other — the extra information needed to correctly detect real write conflicts instead of imposing an arbitrary order over them.

Explain linearizability with a concrete example of when weaker consistency would cause a bug.

Linearizability means every operation appears to take effect atomically at some single instant between invocation and completion, in an order every client agrees matches real time — once a write returns success, every subsequent read by anyone must see it. A distributed lock is the textbook case: if "who currently holds the lock" were only eventually consistent, two different clients could each read a stale "unlocked" state and both believe they acquired the lock, defeating the entire purpose of having a lock.

When would you deliberately choose eventual consistency over something stronger?

When staleness is cheap and unavailability is expensive — a product view counter, a "likes" count, or a shopping cart that just needs to converge before checkout are fine being briefly out of date in exchange for every replica always accepting reads and writes, even during a partition. The moment "briefly wrong" has a real cost — money moves, a lock's ownership, an inventory count that can go negative — is the moment eventual consistency stops being the right default.

How does ZooKeeper (or etcd) provide linearizable writes — what's actually happening underneath?

ZooKeeper uses ZAB (ZooKeeper Atomic Broadcast), a leader-based protocol conceptually close to Raft — a single leader orders every write and replicates it to followers, committing once a majority acknowledge. etcd is more direct: it implements Raft itself. In both cases, all writes funnel through one leader per term/epoch and require majority replication before acknowledgment, which is exactly the mechanism that makes "read the latest committed value" a well-defined, linearizable operation rather than a best-effort guess.

How does Kafka's leader-per-partition model relate to consensus, and how did KRaft change it?

Each Kafka partition has exactly one broker acting as leader, with followers replicating from it — a producer using acks=all waits for the partition's in-sync replica set (ISR) to acknowledge, structurally a quorum write. Historically Kafka relied on ZooKeeper to elect and track partition leaders and cluster metadata, a consensus system bolted on from outside; KRaft replaces that with Kafka's own built-in Raft-based metadata quorum, so the cluster now runs its own leader election and log replication for metadata instead of depending on a separate ZooKeeper ensemble.

How do idempotency keys relate to consensus, if they don't solve it directly?

Leader failover creates a specific uncertainty for a client: a write might have committed on the old leader right before it died, or it might not have — there's no way to know, only that the request is now unanswered. Retrying is correct, but retrying blindly risks double-applying a write that actually succeeded. An idempotency key turns that retry into a safe no-op if the original request already landed — the practical, request-level answer to the same ambiguity that leader election and failover introduce at the cluster level.

What's the trade-off a payments system faces when choosing between CP and AP during a partition?

A payments system almost always chooses CP — refusing to process a transfer when it can't confirm the account's true current balance is far safer than accepting the write and risking an overdraft or double-spend that has to be unwound later, possibly after the money has already left the system. An AP choice optimizes for uptime at the cost of exactly the correctness a ledger cannot compromise on, which is why this site's payment and wallet projects lean on strict ordering and locking rather than eventual consistency for balance state.

If an interviewer asks for the single most practically useful idea from this whole module, what do you say?

That every one of these ideas — CAP, quorums, fencing tokens, vector clocks — is really answering one question from a different angle: when multiple nodes might disagree about the true current state, how do you get them to agree on one answer, or safely detect that they haven't? Naming which of these tools a given system leans on (a majority quorum, a fencing token, a vector clock) is what separates an answer that sounds right from one that demonstrates you understand why the mechanism is actually necessary.

Related guides

Update these hrefs to your published Blogger post URLs once each page is live.

No comments
Leave a Comment