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.
You're operating a service on DynamoDB. A downstream job writes an item, then immediately reads it back with a default GetItem call and sometimes gets a stale or missing result. Walk me through why, and what you'd change.
Sample Answer
Direct answer
By default, DynamoDB's GetItem and Query calls perform eventually consistent reads, which can be served from a replica that hasn't yet applied the most recent write, so a read immediately after a write can come back stale or missing. The fix is either to request a strongly consistent read on that specific call, or to redesign the flow so the downstream step doesn't need to re-read a value it just wrote.
Structured elaboration
Why it happens: DynamoDB replicates every write across multiple storage nodes in the region before acknowledging the write as successful, but a default GetItem/Query can be routed to a replica that hasn't received that write yet. This is a deliberate cost and latency trade-off, not a bug: eventually consistent reads use half the read capacity and typically have lower latency than strongly consistent reads.
Your options as the operator:
- Pass ConsistentRead: true on the specific read that needs the fresh value. It costs twice the read capacity units of an eventually consistent read and only works within the same region, but it guarantees you see every write that was acknowledged before the read started.
- Avoid the read-after-write pattern entirely: have the writer hand the value it just wrote directly to the downstream step instead of making it re-read from the table.
- If a separate process genuinely has to re-read, add a short retry with backoff, since same-region replication lag is typically single-digit milliseconds.
Where the ConsistentRead flag has limits: DynamoDB Global Tables (cross-region replicas) replicate asynchronously, so a strongly consistent read in one region's replica still only guarantees you see every write already acknowledged in that region, never a write still in flight from another region; there is no cross-region strong-consistency option. DynamoDB Accelerator (DAX, an in-memory cache in front of DynamoDB) is a different case: requesting ConsistentRead: true through DAX does work, DAX simply passes that request straight to DynamoDB without serving or populating it from cache, so you get a genuinely fresh read at the cost of losing DAX's cache acceleration for that one call.
Worked example
A checkout service writes an order row, then a fulfillment worker in the same request path reads it back to grab the shipping address. For a sub-4KB item, an eventually consistent GetItem consumes 0.5 RCU while a strongly consistent one consumes 1 RCU, so setting ConsistentRead: true on that one call costs an extra 0.5 RCU and removes the race entirely, versus a blind retry loop that adds latency and still isn't guaranteed to succeed on the first attempt.
Trade-offs and pitfalls
- Turning on strongly consistent reads everywhere "to be safe" roughly doubles read capacity cost and latency across the service; it should be applied surgically to the one call with the race, not the whole read path.
- Strongly consistent reads do nothing for a Global Table's cross-region replica (there is no cross-region equivalent), and while they do work through DAX, they lose all cache acceleration when they do, a gap teams often discover only after a multi-region or DAX rollout.
- The most robust fix is usually architectural (pass the value forward instead of re-reading it), since it removes the timing dependency entirely and also removes the extra read-capacity cost.
What the interviewer probes next
Expect a follow-up comparing this to S3, which has provided strong read-after-write consistency for every operation (new objects, overwrites, and deletes) since late 2020, unlike DynamoDB's opt-in ConsistentRead, and how you'd catch this class of race in production before a customer reports it.
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.
Say you're placing a 5-node quorum-based cluster. Compare spreading those 5 nodes across 3 Availability Zones in one AWS region versus splitting them across two separate regions. How does quorum placement change, and how do you avoid split-brain in each topology?
Sample Answer
Direct answer
Within a single region, 3 AZs give you low-latency, redundant links (typically sub-2ms), so placing 5 nodes as 2-2-1 across the AZs lets you lose one whole AZ and still have a majority (3 of 5) reachable to keep accepting writes. Across two regions you cannot split 5 nodes evenly, so one region ends up holding the majority (3-2), and the region with only 2 nodes can never reach quorum on its own. Split-brain avoidance is the same rule in both cases (only the side that can prove it holds a majority may accept writes), but the two-region case carries a real risk that someone fails the minority side over anyway during a long partition, which is how split-brain actually happens in practice.
Structured elaboration
Single-region, 3-AZ placement (e.g., 2-2-1):
- AZ-to-AZ links are low-latency and rarely fully partition from each other (same region, redundant fiber paths).
- Losing one AZ still leaves 3 of 5 nodes reachable, so majority quorum holds and writes keep flowing.
- This mainly defends against a single AZ outage (power, networking gear), not against a true network split.
Two-region placement (e.g., 3-2):
- The region holding 3 nodes always has majority quorum by itself; the region with 2 never does.
- If the inter-region link drops, the 2-node region correctly refuses writes (it can't reach quorum) while the 3-node region keeps serving. That is the safe outcome, but it means the minority region's healthy nodes go write-unavailable.
- The dangerous failure mode is a human or an automation script promoting the minority region to "keep serving" during the partition. That creates two sides independently accepting writes, i.e. split-brain, and the divergence has to be reconciled or discarded once the link heals.
Why this is mostly an operational risk, not a protocol risk: consensus protocols like Raft already refuse to commit without a majority, so the protocol itself prevents split-brain as long as nobody forces an override. The real risk is a health check or runbook that misreads "can't reach the majority" as "the majority must be down" and promotes the wrong side.
Worked example
With N=5, majority is ceil((5+1)/2) = 3 nodes. In the 2-2-1 AZ layout, losing any single AZ still leaves at least 3 nodes across the remaining two AZs, so quorum holds. In the 3-2 region layout, if the WAN link between regions fails, the 3-node region has quorum (3/5, can elect a leader and accept writes) and the 2-node region does not (2/5, must reject writes and serve stale reads at best) until the link recovers.
Trade-offs and pitfalls
- Multi-region protects against a whole-region outage that multi-AZ cannot, but it adds real commit-path latency (cross-region round trips) and makes the minority-region-unavailable outcome unavoidable with an odd node count split unevenly.
- A common mistake is assuming a 3-2 region split protects both regions equally. It does not: only the 3-node region can survive a partition alone.
- An even split (e.g., 4 nodes as 2-2 across two regions) is worse, not safer: neither side can reach majority alone, so both stop accepting writes, or someone bolts on a tie-breaker node that becomes a new single point of failure.
- Overly aggressive health-check timeouts can misread a slow but healthy cross-region link as a partition and trigger an unnecessary failover, which is why managed cross-region services (e.g., Aurora Global Database) use deliberately conservative promotion procedures rather than fast automatic failover.
What the interviewer probes next
Expect a follow-up on what happens to writes that were in flight on the minority side when the partition started, and whether you'd ever choose an even node count.
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.
Walk me through the CAP theorem: what do consistency, availability, and partition tolerance each guarantee, and why can a distributed system only provide two of the three once a network partition actually occurs? Give one example of a system design that would lean toward consistency (CP) and one that would lean toward availability (AP), and state precisely what each choice gives up. Also clarify how this notion of 'consistency' differs from the one used in ACID transactions.
Sample Answer
Direct Answer
The CAP theorem says a distributed system that can be split by a network partition can only guarantee two of three properties at once: Consistency, Availability, and Partition tolerance. Because real networks do partition (links fail, messages get delayed or dropped), partition tolerance isn't really an optional design choice, so the actual trade-off every replicated system makes, and only makes while a partition is actually happening, is between Consistency and Availability.
What Each Property Guarantees
- Consistency (C): every read returns the result of the most recent completed write, as if there were only one copy of the data (this is the strong, linearizable notion of consistency).
- Availability (A): every request that reaches a non-failed node gets a response, without a guarantee that the response reflects the latest write.
- Partition tolerance (P): the system keeps operating even when the network drops or delays messages between nodes, splitting them into groups that can't talk to each other.
Why You Only Get Two, and Only During a Partition
When there is no partition, a well-built system can offer both C and A: every node can talk to every other node, so it can confirm it has the latest data before answering. The theorem only bites once a partition actually separates the cluster into two or more groups. At that point, a node in the minority (or either side, in a symmetric split) that receives a request has exactly two choices:
- Answer immediately with whatever data it has locally. That satisfies Availability, but the data might be stale relative to a write that landed on the other side of the partition, so it does not satisfy strong Consistency.
- Refuse to answer (return an error or block) until it can confirm it isn't giving out stale data, typically by waiting for the partition to heal or for enough of the cluster to be reachable. That satisfies Consistency, but it fails Availability for that request.
There is no third option that gives both while the partition is open. That is the entire content of the theorem: it's about behavior during the partition window, not a permanent label on a system.
CP and AP Examples
- A CP-leaning example: a consensus-backed coordination store, such as etcd (a distributed key-value store built on the Raft consensus protocol). If a partition isolates a minority of nodes from the quorum, that minority stops serving both reads and writes rather than risk returning stale or conflicting data. It gives up availability on the minority side to preserve strong consistency everywhere it does respond.
- An AP-leaning example: a Dynamo-style, eventually-consistent key-value store. During a partition, every reachable node keeps accepting reads and writes on both sides, so the system stays available, but the two sides can accumulate divergent writes that must be reconciled once the partition heals (via version vectors, last-write-wins, or application-level merge logic). It gives up guaranteed-fresh reads to preserve availability.
CAP's "Consistency" vs. ACID's "Consistency"
These are two different axes, and conflating them is a common interview trap. ACID (atomicity, consistency, isolation, durability) describes properties of a single transaction, typically on one database: its "C" means a transaction only ever moves the database from one state that satisfies its own defined invariants (foreign keys, uniqueness constraints, application-level rules) to another such state. It says nothing about how fresh a read on a different replica is.
CAP's "C" is about replication: whether a read anywhere in the system reflects the most recent completed write, regardless of which physical replica served it. A system can be perfectly ACID-consistent (every transaction respects its constraints) on every individual replica while still being CAP-inconsistent overall, because a stale replica can return an old value that was, at the time it was written, a perfectly valid state.
Trade-offs and Common Pitfalls
- Treating CAP as a fixed label for an entire system is a common misreading. The choice is scoped to a partition and can even be scoped per operation: a single system can serve some requests (say, checkout) with a CP posture and others (say, product-view counts) with an AP posture.
- Don't assume "P" is a design choice you can decline. Every distributed system that spans more than one process over a real network needs to survive partial network failure, so the honest framing is which of C or A you give up when partitioned, not whether to support partition tolerance.
- A frequent good follow-up is PACELC, which asks what you trade off between latency and consistency even when there is no partition happening, since CAP alone is silent about that normal-operation case.
Unlock Full Question Bank
Get access to all 10 Distributed Systems Fundamentals interview questions and detailed answers.
Sign in to ContinueJoin thousands of developers preparing for their dream job.