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.
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.
Explain what a CRDT (Conflict-free Replicated Data Type) is and why state-based and operation-based CRDTs let replicas converge to the same value without any coordination between them. Walk through two concrete examples: a grow-only counter (G-Counter) and an observed-remove set, and describe what property of the underlying merge operation makes convergence guaranteed.
Sample Answer
A CRDT (Conflict-free Replicated Data Type) is a data structure whose update and merge operations are mathematically guaranteed to make every replica converge to the same value, with no locking, coordination, or central authority, as long as every update eventually reaches every replica. State-based CRDTs ship the whole replica state and merge it with a commutative, associative, idempotent join function; operation-based CRDTs ship individual operations that must themselves be commutative and be delivered with causal ordering. A grow-only counter (G-Counter) and an observed-remove set (OR-Set) are the two simplest concrete examples of this guarantee in action.
Why convergence is guaranteed
Convergence works because the merge operation is commutative (order doesn't matter), associative (grouping doesn't matter), and idempotent (merging a state with itself changes nothing), so applying merges in any order, any number of times, produces the same final state. Formally, this kind of merge is called a join, and a replica's state is modeled as an element of a join-semilattice: a partially ordered set where the join always computes the least upper bound of two states. That is the actual property behind convergence without coordination: it is not that conflicts never happen, it is that the merge function is defined so a conflict has exactly one well-defined resolution no matter how or when it gets computed.
G-Counter (grow-only counter)
- State: a vector with one non-negative integer slot per replica, c[i].
- Local update: a replica only ever increments its own slot.
- Merge: element-wise maximum across the two vectors.
c′[i]=max(c1[i],c2[i])
- Read: sum across all slots.
total=∑ic[i]
- Because each slot only ever grows for its own replica, taking the max per slot can never lose an increment either side already recorded.
OR-Set (observed-remove set)
- State: a set of (element, unique tag) pairs, split into an add-set and a remove-set of tags.
- Add(e): mint a fresh tag, insert (e, tag) into the add-set.
- Remove(e): copy every tag currently observed for e into the remove-set; it removes only tags this replica has actually seen, never tags added elsewhere that haven't arrived yet.
- Merge: union the add-sets, union the remove-sets.
- An element counts as present if it has at least one tag in the add-set that is not in the remove-set.
Naming the comparators explicitly
| Strategy | How a conflict is resolved | What it guarantees | Where it fails |
|---|---|---|---|
| Last-write-wins (LWW) | Keep the value with the later timestamp, discard the other | Deterministic if clocks are totally ordered | Silently discards a concurrent write; a clock-skewed node can win even though its update happened earlier in real time |
| Vector clocks | Compare vectors to detect that two writes are concurrent | Tells you a conflict exists | Detection only. It does not resolve the conflict; an application or a person still has to pick a winner |
| CRDTs (this answer) | The merge function is commutative, associative, and idempotent by construction | Automatic, coordination-free convergence | Only works for data types whose semantics fit that mold; does not generalize to arbitrary business logic |
| Application-specific merge | Domain code decides, for example union two shopping carts, or keep the higher of two account balances | Correctness tailored to the domain | Bespoke code per data type; nothing about it is automatic or reusable |
Worked example: G-Counter convergence
Three replicas A, B, C start at (0,0,0):
- Replica A processes 2 local increments: its state becomes (2,0,0).
- Replica B processes 3 local increments, concurrently, before hearing from A: (0,3,0).
- Replica C stays idle: (0,0,0).
A and B exchange state and merge (element-wise max): merge((2,0,0),(0,3,0)) = (2,3,0). Read = 2+3+0 = 5. C later merges with that result: merge((0,0,0),(2,3,0)) = (2,3,0). Read = 5. Whichever order the three replicas merge in, the final vector is (2,3,0) and the read is 5, exactly matching the 2+3=5 real increments actually performed. No increment is lost and none is double-counted.
Worked example: OR-Set add and remove race
Replicas R1 and R2 have already converged on a set containing 'milk' with tag t1. The two replicas are then partitioned from each other:
- R1's user removes 'milk': remove-set gains {t1}, the only tag R1 has ever observed for 'milk'.
- R2's user, unaware of the removal, re-adds 'milk': add-set gains a brand-new tag {t2}, so the add-set is now {t1, t2}.
On merge: add-set = {t1, t2} (union), remove-set = {t1} (union). 'milk' is present because t2 is in the add-set and not in the remove-set. This is the correct outcome: R2's re-add introduced a tag the remover never saw, so it survives, exactly the observed-remove semantics the name describes.
Trade-offs & pitfalls
- Storage and bandwidth: every element needs extra metadata (a vector slot per replica for counters, a unique tag per add for sets), and removed elements don't disappear until a garbage-collection pass establishes causal stability across replicas.
- The edge case that catches teams out: CRDTs don't compose across non-commutative operations. A G-Counter or OR-Set is safe because the operations that define it (increment, tagged add and remove) are commutative by construction. But if you build an append-only log CRDT and then bolt on an application-level 'delete the last 3 entries' operation defined by position, that composition is not well-defined under concurrency: 'last 3' means something different on each replica depending on how many entries have been concurrently appended there at the time the delete runs, so two replicas can end up deleting different entries even though each individually applied a correct-looking CRDT merge. The fix is the principle OR-Set already uses: target deletions by a stable element identifier, never by position or count.
- When to avoid: anywhere a global invariant spans multiple keys (uniqueness, a balance that must never go negative), or the business logic genuinely isn't commutative. CRDTs solve convergence, not arbitrary correctness.
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.
Explain how checkpointing works in a stateful stream-processing framework: how a barrier or snapshot marker flowing through the pipeline lets the system capture a consistent point-in-time state across many parallel operators, and how the system uses that checkpoint to restore and resume with exactly-once semantics after a failure.
Sample Answer
Direct answer
A checkpoint barrier is a special marker the coordinator injects into every source stream at a chosen moment. As it flows downstream mixed in with real data, each operator uses its arrival to mark a cut: everything on that input channel before the barrier belongs to checkpoint N, everything after belongs to checkpoint N+1. Once an operator has seen the barrier on all of its input channels, it takes a local snapshot of its own state and forwards the barrier onward. Because every operator's snapshot is cut at the same logical point in the data rather than the same wall-clock instant, the union of all local snapshots plus the recorded source read-positions forms one consistent global snapshot the whole job can be rewound to after a failure.
Barrier injection and alignment
The checkpoint coordinator periodically assigns an increasing checkpoint id and injects a barrier carrying that id into every source partition. Barriers travel with the data on each channel, in order, never overtaking a record. When an operator with multiple input channels receives the barrier on one channel before the others, it stops consuming further records on that channel and buffers them, while continuing to process the channels where the barrier hasn't arrived yet. This buffering-until-all-channels-caught-up step is called alignment; it guarantees the operator's eventual snapshot reflects exactly the same cut point on every input.
Local snapshot and coordinator commit
Once aligned, the operator snapshots its local state (for a keyed aggregation, the current value per key) to durable storage and forwards the barrier to its downstream operators. The coordinator marks a checkpoint complete only once every operator, all the way to the sinks, has acknowledged it, and only then persists the checkpoint's metadata as the new restore point.
Unaligned checkpoints
Alignment can add latency under backpressure or skew, since a fast channel has to wait on a slow one before the operator can snapshot. Unaligned checkpoints avoid this by not waiting at all: the framework snapshots the buffered in-flight records themselves as part of the checkpoint, alongside the operator's own state, trading more storage and I/O for lower checkpoint latency.
Worked example: barrier flow through a keyed aggregation
Job topology: Source -> Map -> KeyedSum -> Sink, with KeyedSum running as two parallel instances, A and B.
sequenceDiagram
participant Coord as Coordinator
participant A as KeyedSum A
participant B as KeyedSum B
participant Sink
Coord->>A: barrier(42) on P0
Coord->>A: barrier(42) on P1
A->>A: snapshot state
A->>Sink: forward barrier(42)
Coord->>B: barrier(42) on P0
Coord->>B: barrier(42) on P1
B->>B: snapshot state
B->>Sink: forward barrier(42)
Sink->>Coord: ack checkpoint 42
- Coordinator starts checkpoint 42 and injects
barrier(42)into the source's two partitions, P0 at read-offset 1000 and P1 at read-offset 850. - Instance A receives input from both P0 and P1 after the shuffle (the keyed redistribution that routes every record for a given key consistently to the same parallel instance).
barrier(42)arrives on A's P0 channel first: A stops consuming new P0 records and buffers them (alignment), while continuing to process P1 records normally, since P1's barrier hasn't arrived yet. barrier(42)arrives on A's P1 channel. A has now seen the barrier on every input channel, so it snapshots its local running sums, say{key=X: 17, key=Y: 42}, to durable storage, forwardsbarrier(42)downstream to the Sink, and unblocks the buffered P0 records, which now belong to checkpoint 43.- Instance B does the same independently for its own keys.
- The Sink receives
barrier(42)from both A and B, snapshots (or, if it's a transactional sink, pre-commits) its own pending output, and acknowledges checkpoint 42 to the coordinator. - Once the coordinator has acknowledgements from the source (offsets P0=1000, P1=850), A, B, and the Sink, checkpoint 42 is marked complete and persisted.
- If the job then crashes after checkpoint 42 completed but before checkpoint 43 finished, restart loads checkpoint 42: sources reset to offsets P0=1000/P1=850, A and B restore their snapshotted key-sums, and the Sink either commits its pre-committed checkpoint-42 output or relies on idempotent writes if it isn't transactional. Processing then resumes from exactly those offsets, so no record before the barrier is reprocessed and no record after it is lost.
End-to-end exactly-once needs the sink to participate
Internal exactly-once state is only as strong as what happens at the external sink. A transactional (two-phase commit, 2PC) sink pattern treats each checkpoint as a transaction boundary: during the checkpoint, the sink prepares and flushes its output but does not commit; once the coordinator marks the checkpoint complete, it tells the sink to commit. On failure, any uncommitted transaction is aborted, so restored state and committed output stay consistent. If a sink can't participate in a transaction this way (for example, one that only supports single-item idempotent writes rather than a cross-partition transaction), the fallback is idempotent writes keyed by (checkpoint_id, record_id), so the same output re-emitted after restoring from checkpoint 42 doesn't double-apply. The same barrier/snapshot mechanism underlies this regardless of what the job is computing: a change-data-capture (CDC) ingestion pipeline, a feature-computation job for machine learning, or an ordinary aggregation are all just other kinds of stateful stream jobs from the checkpoint coordinator's point of view.
Trade-offs & pitfalls
| Aligned checkpoints | Unaligned checkpoints | |
|---|---|---|
| Latency under backpressure/skew | Can stall waiting for the slowest channel | Avoids the wait entirely |
| Storage/I/O overhead | Lower (only operator state is stored) | Higher (in-flight buffered records are stored too) |
| Reasoning simplicity | Simple, single well-defined cut point | More moving parts to restore correctly |
Rescaling (changing parallelism) needs a savepoint, a user-triggered durable checkpoint, plus a defined remapping of keyed state across the new number of instances. The most common pitfall is assuming "exactly-once" automatically covers the whole pipeline: if the sink isn't transactional or idempotent, restoring from a checkpoint after a crash can re-emit output that was already delivered before the crash, turning exactly-once internal state into at-least-once external effects.
Explain quorum-based reads and writes using the N/R/W notation (N replicas, W write quorum, R read quorum). Using a concrete example with N=5, show why W + R > N is required to guarantee that every read sees the most recent write, and discuss how shifting R and W trades off latency, availability, and durability when nodes fail.
Sample Answer
Direct Answer
In a system with N replicas, a write is only considered committed once W of those replicas have acknowledged it, and a read is only considered complete once R replicas have been queried and the freshest value among their answers is returned. If you pick W and R so that
W+R>Nthen every possible set of W replicas and every possible set of R replicas are guaranteed to overlap in at least one replica, which means any read is guaranteed to touch at least one replica that has the most recent write.
Why the Overlap Guarantee Holds
This falls out of a simple counting fact: if you pick two subsets of a set of N items, and the sizes of those two subsets add up to more than N, they cannot be disjoint. If a write-set of size W and a read-set of size R were completely disjoint, sharing no replica at all, together they would use W + R distinct replicas out of only N available, which is impossible once W + R > N. So the two sets must share at least one replica, and since the write-set includes every replica the write reached, the shared replica is guaranteed to have seen the latest write.
∣A∣+∣B∣>N⟹A∩B=∅for A,B⊆{1,…,N}Worked Example, N = 5
Take five replicas, labeled 1 through 5. Choose W = 3 and R = 3 (3 + 3 = 6 > 5, so the guarantee holds).
A write commits to replicas {1, 2, 3}, the write quorum. A later read queries replicas {3, 4, 5}, the read quorum. The overlap between {1, 2, 3} and {3, 4, 5} is {3}, so replica 3 is guaranteed to be in both sets, and since replica 3 has the latest write, the read correctly returns the fresh value even though replicas 4 and 5 are still stale.
Now see what happens if you drop below the threshold: keep W = 3 but use R = 2 (3 + 2 = 5, not greater than N = 5, so the guarantee no longer holds). A read that happens to query {4, 5} shares no replica at all with the write quorum {1, 2, 3} and would return the stale value those two replicas still hold, with no way for the client to know it missed the latest write.
Trading Off Latency, Availability, and Durability
- Lowering W speeds up writes, since fewer replicas have to acknowledge, and lets writes succeed even if more replicas are down, but it weakens durability (fewer copies exist right after the write) and forces R to be larger to keep W + R > N, which slows reads down instead.
- Lowering R speeds up reads the same way, at the cost of needing a larger W.
- To keep serving at a chosen W or R while tolerating f replica failures, you need enough surviving replicas to still form that quorum, so majority quorums, such as W = R = 3 for N = 5 (the smallest quorum size bigger than half of 5), are a common default: they satisfy W + R > N for any N, and they keep working as long as a majority of replicas are reachable.
Leaderless Quorums vs. a Leader-Based Design
Quorum systems like this are naturally leaderless: any client can attempt a write or a read against any W or R replicas without funneling through one elected coordinator, unlike a Raft-based design where every write has to go through the single current leader. That gives quorum systems more availability during a partition, since any reachable set of W or R replicas can keep working, at the cost of needing real conflict handling: two writes that each reach a different, overlapping-but-not-identical set of replicas can produce concurrent versions that a read has to reconcile, by comparing versions and taking the latest or surfacing both to the application, which a single-leader system avoids by construction since all writes are already serialized through the leader.
Choosing Sane Defaults in a Client Library
A client library that exposes N, R, and W as tunable knobs should default to a majority quorum on both sides, W = R = the smallest integer greater than N / 2, rather than exposing the raw numbers with no guidance, because majority-on-both-sides is the smallest configuration that always satisfies W + R > N regardless of N, and it gives a reasonable latency-versus-safety balance without requiring the caller to re-derive the inequality themselves. The library should still let advanced callers override it, such as R = 1 for the fastest possible read when the caller is prepared to handle occasional staleness itself, or W = N for maximum durability when the caller can tolerate slower writes, with the safety trade-off documented at each override.
Trade-offs and Pitfalls
- Quorum overlap guarantees that a read touches at least one replica with the latest write; it does not by itself guarantee the read correctly identifies which of the R responses is the latest one. Without comparing versions or timestamps correctly across the R responses, you can still return a stale value even though the fresh one was right there in the response set.
- Concurrent writes are a real gap: if two writes race and land on different, only-partially-overlapping write quorums, you can end up with genuinely concurrent versions that need reconciliation, not just staleness that time will fix.
- Picking W = 1 to maximize write availability forces R = N to keep the safety guarantee, which makes every read fragile to a single unavailable replica; it's rarely a good default outside very read-light, write-heavy workloads that can tolerate that risk.
Unlock Full Question Bank
Get access to all 32 Distributed Systems Fundamentals interview questions and detailed answers.
Sign in to ContinueJoin thousands of developers preparing for their dream job.