Distributed Systems Fundamentals Questions
Core theory that underpins any multi-node system: the CAP and PACELC theorems, consistency models (strong, causal, eventual), partitioning, replication, and the fundamental tradeoffs between latency, availability, and consistency. Covers how network partitions, clock skew, and partial failure change the reasoning compared to single-node systems. This is the vocabulary layer every distributed design question builds on.
Explain Lamport clocks and vector clocks: how each captures a happens-before relationship between events, and what information a vector clock encodes that a Lamport clock does not (distinguishing genuine causality from mere concurrency). Walk through why two events can be 'concurrent' under this model even though one clearly happened at an earlier wall-clock time.
Sample Answer
Lamport clocks and vector clocks both order events in a distributed system without relying on wall-clock time, which cannot be trusted to stay synchronized across machines. A Lamport clock is a single integer per process that increases on every local event and every message received, guaranteeing that if event A happened-before event B, A's counter is smaller than B's, but not the reverse: two events can tie or land on comparable counter values without one having actually caused the other. A vector clock is a full vector, one counter per process, that lets you tell exactly whether two events are causally related or genuinely concurrent, which is the extra information a single Lamport counter throws away.
Lamport clocks
- Each process keeps one integer counter, starting at 0.
- Local event: increment own counter.
- Send: increment, then attach the counter to the message.
- Receive: set counter = max(local counter, counter in message) + 1.
- Guarantee: if A happened-before B, then LC(A) < LC(B). The converse does not hold: LC(A) < LC(B) does not imply A happened-before B.
Vector clocks
- Each process keeps a vector with one slot per process, all starting at 0.
- Local event: increment own slot.
- Send: increment own slot, attach the whole vector.
- Receive: take the element-wise maximum of the local vector and the incoming vector, then increment own slot.
- Comparison rule:
V(A)≤V(B)⟺∀i, V(A)i≤V(B)i and ∃j, V(A)j<V(B)j
- If neither V(A) <= V(B) nor V(B) <= V(A) holds, the vectors are incomparable, and the events are genuinely concurrent: no message path connects them in either direction, regardless of what wall-clock time either happened at.
Worked example: a two-person chat, printed event trace
Two people, on process P1 and process P2, are chatting. Message ordering here needs to respect causality: a reply should never appear to precede the message it replies to, which is exactly what vector clocks are for.
- e1 (P1, local event, user starts typing): Lamport clock 1, vector clock [1,0].
- e2 (P1, sends message m1 to P2): Lamport clock 2, vector clock [2,0], attached to m1.
- e3 (P2, local event, user independently opens the chat window before receiving anything from P1): Lamport clock 1, vector clock [0,1]. In real wall-clock terms, say this happens several seconds before e1 even occurs on P1's machine, since the two users' actions are completely independent at this point.
- e4 (P2, receives m1): Lamport clock = max(1, 2) + 1 = 3. Vector clock = elementwise max([0,1], [2,0]) = [2,1], then increment P2's own slot: [2,2].
Now compare e1 and e3: Lamport clocks are LC(e1)=1 and LC(e3)=1, a tie. A Lamport clock alone gives no way to tell whether these are causally related from the numbers themselves; forcing a total order would need an arbitrary tie-break, like comparing process identifiers, and that tie-break tells you nothing true about causality. The vector clocks settle it precisely: V(e1)=[1,0] and V(e3)=[0,1] are incomparable, since 1 > 0 in the first slot but 0 < 1 in the second, so e1 and e3 are concurrent by definition, even though e3 happened earlier in real wall-clock time in this scenario. Concurrency here is about the absence of a causal path, not about which one occurred first on a wall clock.
Now compare e3 and e4: V(e3)=[0,1], V(e4)=[2,2]. Every slot of V(e3) is less than or equal to the corresponding slot of V(e4), and the first slot is strictly less (0<2), so V(e3) <= V(e4), and e3 happened-before e4, correctly, since e3 and e4 both occurred on P2 in that program order.
Trade-offs & pitfalls
- Vector clocks only detect concurrency; they do not resolve it. When V(A) and V(B) are incomparable and both represent a write to the same piece of data, the vector clock correctly tells you there is a genuine conflict, but not which write should win. An application still needs a policy on top, last-write-wins by some tie-break, a CRDT merge, or surfacing both versions for a user or client to reconcile; the vector clock's job stops at detection.
- Storage cost: a vector clock needs one slot per participating process, so it grows with the number of writers, unlike a Lamport clock's single integer. Systems with many writers usually prune or cap this, for example with dotted version vectors or per-shard writer sets, rather than keep an ever-growing vector per object.
- Common wrong turn: assuming a Lamport clock's total order reflects real causality. It gives a valid total order consistent with happened-before, so if A really did happen before B, Lamport respects that, but not every pair the Lamport order ranks is actually causally related, so Lamport clock values alone cannot answer whether A caused B.
What is PACELC, and how does it extend the CAP theorem? Walk through an example decision where PACELC's latency-versus-consistency trade-off matters even when there is no active network partition.
Sample Answer
Direct answer
PACELC, short for "if Partition, Availability vs. Consistency; Else, Latency vs. Consistency", says that CAP's dilemma, choose Consistency or Availability when a network Partition is happening, is only half the story. Even when there is no partition at all, a system still has to choose between Latency and Consistency for every write it replicates, because making a write durable on every replica before acknowledging it takes longer than acknowledging it once it's durable on a single node. PACELC packages this as: if Partition occurs, trade off Availability against Consistency (exactly what CAP already says); Else, meaning no partition, trade off Latency against Consistency.
Restating CAP precisely first
CAP says that during an actual network partition, a distributed system can guarantee only one of Consistency (every read sees the latest completed write) or Availability (every request gets a non-error response) for the nodes on either side of the split, not both. A common misreading treats CAP as "pick two of three, always"; it isn't. CAP's teeth are specifically about behavior during a partition. Most systems are both consistent and available almost all of the time, precisely because a true network partition is a rare event relative to total uptime, not something happening continuously.
flowchart TD
Start[Write occurs] --> P{Partition active?}
P -->|Yes| AC[Choose Availability or Consistency]
P -->|No| LC[Choose Latency or Consistency]
What PACELC adds
PACELC names the trade-off CAP is silent about: during normal operation, with no partition, you still choose between Latency (L) and Consistency (C), because synchronous replication that waits for a majority of replicas costs a round trip before it can acknowledge a write, while asynchronous or single-node-acknowledged replication returns faster but risks a reader seeing stale data, or the acknowledged write being lost outright if that one node fails before it propagates. Systems are commonly labeled by both branches together, for example PA/EL (favor Availability under partition, favor Latency otherwise, the Cassandra/Dynamo-style default) or PC/EC (favor Consistency in both cases, the HBase-style default).
Worked example: a decision with no partition occurring
A write to a piece of user data must be replicated to three nodes: R1 in the local region, and R2, R3 in two remote regions. All three are reachable; no partition is happening anywhere in this example.
- Favor consistency (the "C" side of the Else branch): the write path waits for acknowledgment from a majority, at least two of the three replicas, say R1 and R2, before returning success to the caller. Any subsequent read from a majority quorum is now guaranteed to see this write. Cost: the caller's write waits on the round trip to R2, a remote replica, even though R1, the local one, already has it durably.
- Favor latency (the "L" side of the Else branch): the write path acknowledges as soon as R1 has it durably, and replicates to R2 and R3 asynchronously in the background. Cost: the caller gets a fast, local acknowledgment, but a read served from R2 immediately afterward, before the async replication catches up, will not see the write yet. If R1 crashes before that background replication completes, the already-acknowledged write can be lost entirely, with zero partition ever occurring.
This decision, wait for two of three versus acknowledge on one, is made on every single write regardless of whether any partition is happening, which is exactly the trade-off PACELC's Else branch names and CAP alone has nothing to say about, since CAP only speaks to a system that is not fully connected.
Trade-offs & pitfalls
A common misconception is treating a database's PACELC label as a fixed law of the software rather than a description of its typical default: most systems let you tune the replication wait per request (via quorum size), so "Cassandra is PA/EL" describes its usual configuration, not something it's incapable of changing. It's also easy to blur this Else-branch trade-off with an availability discussion; in the worked example above, no node was ever unreachable, so the trade being made is purely about how long the write path waits before acknowledging, not about surviving an outage, which is a separate concern belonging to the partition branch of the theorem.
Compare Raft and Paxos as consensus protocols: how does each actually reach agreement, and why is Raft generally considered easier to reason about and implement? Give a situation where a team might still reach for Paxos (or a Paxos variant) over Raft, and one where you'd rather rely on an external coordination service (etcd, ZooKeeper, Consul) than embed a consensus implementation yourself.
Sample Answer
Direct Answer
Paxos and Raft are both protocols that let a cluster of nodes agree on a value, or, in the log-replication case most real systems use, an ordered sequence of values, despite crashes and message delays, using a majority quorum so that any two decisions are guaranteed to have at least one node in common. Raft reaches the same safety guarantee as Paxos but organizes the protocol into named, sequential subproblems, mainly a single strong leader that serializes all writes during its term, which most engineers find much easier to implement correctly than Paxos's more general and symmetric design.
How Each Actually Reaches Agreement
Paxos (Multi-Paxos in practice). A proposer picks a proposal number and sends a Prepare message to the acceptors; each acceptor promises not to accept any proposal numbered lower and reports back the highest-numbered proposal it has already accepted, if any. Once the proposer hears back from a majority, it sends an Accept message carrying the value from the highest-numbered already-accepted proposal it was told about, not necessarily its own original value, and the value is chosen once a majority of acceptors accept it. Multi-Paxos elects a stable leader so that steady-state operation can skip repeating the Prepare phase for every new value.
Raft. Raft splits the same problem into leader election, where nodes agree on a single leader for a numbered term using randomized timeouts and majority votes, log replication, where the leader appends client commands to its own log and replicates them to followers, treating an entry as committed once a majority of nodes have stored it, and a safety rule that a candidate can only win an election if its log is at least as up to date as a majority of the cluster, which prevents a new leader from ever overwriting an already-committed entry.
Why Raft Is Easier to Implement and Reason About
Raft's decomposition gives each subproblem, who's the leader, how entries get replicated, how membership changes safely, its own explicit invariant, so an implementer can reason about one piece at a time. Paxos's proposer and acceptor roles are more general and symmetric, any node can propose at any time, which is elegant but produces more possible interleavings of concurrent proposals to reason about, especially once you move from the single-value textbook description to a real, steady-state, multi-value system, which is where most of the genuinely tricky Paxos engineering, such as stable leader election, log compaction, and membership changes, actually lives, and where the original paper says relatively little.
Comparison Table
| Paxos (Multi-Paxos) | Raft | |
|---|---|---|
| Roles | Proposer, acceptor, learner; any node can propose | Single leader, followers, candidates; the leader serializes all writes during its term |
| Phases | Prepare/Promise then Accept/Accepted, repeated per value (a steady leader skips Prepare) | Leader election once per term, then log replication per entry, plus a separate membership-change protocol |
| Core safety argument | Quorum intersection across proposal numbers that might be concurrent | A candidate can only win with a log at least as up to date as a majority, so a new leader can never miss committed entries |
| Common production use | Google's Chubby, an internal Paxos-based lock and coordination service, plus various in-house tuned variants | etcd, Consul, and CockroachDB, which runs one Raft group per data range |
A Worked Trace: Why a Competing Proposer Can't Just Overwrite the Value
Three acceptors, A1, A2, A3. Proposer P1 sends Prepare(1) to all three; none has accepted anything yet, so all three promise and report nothing. P1 gets a majority, 3 of 3, so it sends Accept(1, X). A1 and A2 accept proposal (1, X) before a second proposer, P2, starts a competing round. P2 sends Prepare(2) to A2 and A3; it doesn't reach A1. A2 has already accepted (1, X), so it promises not to accept below 2 and reports that it already accepted (1, X). A3 has accepted nothing, so it promises and reports nothing. P2 now has a majority of promises, A2 and A3, but because A2 reported an already-accepted value, the protocol requires P2 to propose that same value, X, rather than whatever value P2 originally intended. P2 sends Accept(2, X), not Accept(2, Y). Even though P2 won the second round, the value that gets chosen is still X. This is exactly the mechanism that keeps Paxos safe under concurrent proposers: a later round can change who proposes, but it cannot change what gets chosen once a value has reached a majority.
When to Still Reach for Paxos
Reach for Paxos, or a variant, instead of Raft when you're extending or must interoperate with an existing Paxos-based system where a rewrite isn't worth the risk, or when you need the extra flexibility Paxos's more general, symmetric design supports, such as non-majority or weighted quorum arrangements tuned for a specific wide-area latency profile, or when your team already has deep, tested Paxos expertise and infrastructure that a switch to Raft wouldn't meaningfully improve on.
When to Rely on an External Coordination Service Instead of Embedding Consensus Yourself
Most application teams don't actually need to choose between Raft and Paxos at all; they need a small set of coordination primitives, such as leader election for their own service, a shared lock, a bit of shared configuration, or service discovery, and implementing a correct Raft or Paxos group from scratch means owning a lot of subtle correctness surface, including log compaction, snapshotting, membership changes, and safe leadership transfer, for something an existing, battle-tested coordination service already does well. In that case, point your application at etcd, ZooKeeper, or Consul rather than embedding a consensus implementation. The case for building your own is when the consensus group needs to sit directly in your own data path for latency or throughput reasons, for example when you're building a replicated database yourself and every write needs to go through your own consensus group rather than round-tripping to an external service.
Trade-offs and Pitfalls
- It's a common misreading to treat Paxos as worse than Raft; it's a general, provably minimal algorithm. What actually makes it hard is that the original description covers a single value, and turning that into a real, steady-state, multi-value system requires additional engineering, such as a stable leader, log compaction, and membership changes, that Raft specifies as part of its core design instead of leaving as an exercise.
- Don't conflate which consensus algorithm to use with the more common real decision, which is whether to implement any consensus algorithm yourself at all; for most teams, depending on an existing coordination service is the right default, and only teams actually building infrastructure-level replicated systems typically end up choosing between Raft and Paxos directly.
What problems does clock skew between machines create in a distributed system? Give at least three concrete examples (event ordering across services, a lease that expires early or late, a TLS certificate that appears valid or invalid depending on which node's clock you ask) and describe, at a high level, why this makes naive wall-clock-based ordering unsafe.
Sample Answer
Direct Answer
Clock skew is the difference between what two machines' clocks read at the same real instant. It matters because any decision that compares timestamps from different machines to decide what happened first, whether a lease is still valid, or whether a certificate is still in its valid window, quietly assumes those clocks agree, and in a real network they don't. A numerically later timestamp on one machine's clock does not reliably mean later in real time once you're comparing across machines.
Three Concrete Problems
1. Event ordering across services. If service A stamps an event with its own local clock and service B stamps a related event with its own local clock, and A's clock runs even slightly ahead of B's, an event that actually happened after (in real time) on B can end up with a numerically smaller timestamp than an earlier event on A. Anything that reconstructs what happened in what order by sorting on raw timestamps can get the sequence backwards. The same failure shows up in machine learning feature pipelines: if a feature-store write on one node happens slightly after a model-serving read on another node consumed the old value, but the writer's clock runs a bit fast, the write can carry an earlier or overlapping timestamp than the read, which corrupts any point-in-time audit of which feature value was actually used for a given prediction, if that audit trusts the raw timestamps.
2. A lease that expires early or late. Distributed locks are commonly granted as leases valid until wall-clock time T. If the holder's clock runs slow relative to the granting service's clock, the holder can believe it still owns the lease past the point the granting service has already reassigned it, expiring late from the holder's point of view and risking two nodes both acting as if they hold the resource. If the holder's clock runs fast instead, it can abandon a still-valid lease early and stop acting well before the granting service considers it expired, causing unnecessary churn.
3. A TLS (Transport Layer Security) certificate that looks valid on one node and invalid on another. Certificate validity is a wall-clock range check, evaluated locally by whichever machine happens to be doing the handshake, against the certificate's not-before and not-after dates. If one node's clock has drifted backward past the not-before date, or forward past the not-after date, that single node rejects a certificate every correctly-clocked node accepts, or the reverse, producing a confusing, node-specific TLS failure that looks like a certificate problem but is actually a clock problem.
A Concrete Trace of Why Naive Ordering Is Unsafe
Node A's clock reads 100 and Node B's clock reads 96 at the same real instant, a 4-unit skew. Event E1 happens on Node A at that instant and is stamped 100. Two time units later in real time, event E2 happens on Node B; by then B's clock reads 98 (96 plus 2), so E2 is stamped 98. Comparing the raw timestamps, 100 is greater than 98, so E1 looks like it happened after E2. But in real time, E2 actually happened after E1. Any process that orders events purely by comparing these timestamps gets the sequence exactly backwards, even though the comparison itself, 100 greater than 98, is arithmetically correct.
Why This Isn't Just a Sync-the-Clocks-Better Problem
Synchronizing physical clocks, with NTP (the Network Time Protocol) or a hardware-disciplined protocol like PTP (Precision Time Protocol), reduces the size of the skew, but it does not make comparing two independently-running clocks perfectly safe; it only shrinks the window in which the trace above can happen. The standard engineering answer for ordering that has to be correct, not just human-readable, is to stop relying on raw wall-clock comparison for that purpose and use a logical clock instead: a counter that each node increments on its own events and carries along with outgoing messages, which correctly captures which events could have influenced which others regardless of clock drift. Production distributed databases often go a step further and use a hybrid logical clock (HLC), which combines a physical-time component with a logical counter, to get correct ordering without giving up a timestamp that's still roughly readable as wall-clock time.
Trade-offs and Pitfalls
- Don't confuse clock skew, which is that two clocks disagree right now, with clock drift, the rate at which they diverge over time; skew is the instantaneous symptom you observe, drift is the ongoing cause, and a monitoring setup that only alerts on one of them will miss the other.
- A common wrong turn is assuming that running NTP means this is handled. NTP typically keeps clocks within milliseconds of each other, which is fine for human-readable log timestamps, but any nonzero skew is still a real correctness risk for anything that depends on strict ordering, so it reduces the problem rather than eliminating it.
- The lease-expiry problem specifically is best addressed by combining a conservative time-to-live with fencing at the protected resource, rejecting stale operations based on a monotonically increasing token rather than on wall-clock time at all, instead of trying to shrink clock skew to zero, which isn't achievable.
A system serves linearizable reads from a single leader to guarantee strong consistency, but read latency from remote regions is high. Propose at least three ways to reduce that latency for reads that don't strictly need the freshest possible value, while preserving strong guarantees for the reads that do.
Sample Answer
Direct answer
Keep exactly one path for reads that must be linearizable, meaning the read is guaranteed to see every write that completed before the read began, and give everything else a cheaper path that trades a bounded amount of staleness for a shorter round trip. The three techniques below differ only in how they let something other than a full leader round trip answer safely, without breaking the guarantee for the reads that actually ask for it.
Technique 1: bounded-staleness follower reads
Each replica tracks the highest write position (a log index or sequence number) it has applied. A client that can tolerate some staleness issues a "read as of no more than X positions old" request to its nearest replica; if that replica's applied position is within X of the leader's latest, it answers locally. A strict "give me the current value" request is never eligible for this path and always goes to the leader (or through one of the other two techniques below).
Technique 2: leader lease reads
The leader holds a time-bound lease, renewed through its normal heartbeat or replication round trip with followers, that certifies no other node could have become leader before the lease expires. While the lease is valid, the leader can answer a read from its own local state without running a fresh round of consensus for that specific read, because the lease itself is the proof no newer leader could have committed a write elsewhere in the meantime. This only removes the coordination round trip, not the trip to the leader itself, so it mainly helps latency when the leader happens to be close to the requester; it does not by itself let a remote follower answer.
Technique 3: read-index protocol (two-tier API)
Before answering a strong read, the leader (or a follower proxying to it) records the currently committed log index, the "read index", and confirms with a lightweight quorum check-in, not a full new log entry, that it is still the leader. It can then let any replica serve the read once that replica's own applied index has reached the read index, including a nearby follower, without paying for a full consensus round trip per read. Exposing this as two API surfaces, a ReadStrong() that always uses the read-index path and a ReadFast() that skips straight to the nearest replica's current state, lets the client declare which guarantee it actually needs.
Worked example: comparing index positions
A product's price is replicated with a monotonically increasing log index. The leader's last committed index is 582. A follower in a remote region has an applied index of 579, three entries behind due to ordinary asynchronous replication lag.
- Strict read request: the system requires the answering replica to show an applied index of at least 582 before answering. The remote follower, at 579, does not qualify, so this read is served either by the leader directly, or the follower must first catch up to 582 (via read-index or a direct proxy).
- Bounded-staleness read request, tolerance = 10 positions: the remote follower's applied index (579) is only 3 behind the leader's 582, well within the stated tolerance of 10, so it answers locally with no round trip to the leader at all.
This is the actual mechanism behind "reduce latency for reads that don't need the freshest value": comparing an explicit applied-index number against an explicit, stated tolerance, not a vague notion of "probably fresh enough."
Trade-offs & pitfalls
| Technique | Where latency actually drops | Failure mode if misused |
|---|---|---|
| Bounded-staleness follower reads | Any read routed to a nearby follower within tolerance | Silently returns a stale value if an application defaults every read to this path, including ones that needed read-your-writes |
| Leader lease reads | Only reads served by the leader itself | Correctness depends on bounded clock drift; too long a lease widens the window in which a partitioned old leader could still believe itself current |
| Read-index protocol | Any caught-up replica, without a full consensus write per read | A replica behind the read index has to wait or catch up, which reintroduces latency proportional to replication lag in the worst case |
The most common pitfall across all three is exposing a fast path and a strong path as separate API calls and then trusting the client to pick correctly every time; a safer design escalates automatically, for example retrying as a strong read when a fast read's own returned version looks suspiciously far behind what was expected, rather than relying solely on the caller's judgment.
Unlock Full Question Bank
Get access to all 31 Distributed Systems Fundamentals interview questions and detailed answers.
Sign in to ContinueJoin thousands of developers preparing for their dream job.