Distributed & scalable systems
What changes once a system spans more than one machine: replication, consensus, partial failure, time, and the scalability practice built on top of them. It is the deepest area of most senior interviews because the intuitions from single-process code stop holding.
on this pageshowhide
guide
overview
~2 minDistributed systems is where an interview goes once the problem stops fitting on one machine. What interviewers probe is whether your single-process intuitions have been replaced: a remote call can fail without telling you whether it did anything, two machines disagree about the time, a copy of the data can lag behind, and part of the system can be unreachable while the rest keeps serving. At junior levels the questions check vocabulary: what replication buys, what eventual consistency means, why a load balancer exists. At senior and principal levels they check whether you can name the guarantee a design actually gives, the failure that breaks it, and the price paid to keep it. The hub splits into three sections. [Core distributed theory](/topics/found-distributed-systems-core-theory) is the foundation: replication, partitioning, consensus and leader election, physical and logical time, failure detection, distributed transactions and delivery semantics. A good share of it is about limits, the things no protocol can do, as much as about how to build what is possible. [Consistency models](/topics/found-eventual-consistency) takes one thread of that theory to depth: the spectrum from linearizable down to eventual, quorums, session guarantees, conflict resolution and CRDTs. [Scalability and system design](/topics/found-system-design) is the largest and most applied section: capacity estimation, load balancing, caching, messaging, rate limiting, observability and the data tier, then the systems design rounds keep returning to, such as feeds, search, proximity, job schedulers and file storage. Learn the theory first, even if your target is the design round. Design answers are judged by the trade-offs they name, and those trade-offs come from the theory: a cache is a replica with a staleness budget, a shard key is a partitioning decision, a queue is a promise about delivery. Start with failure models and replication, then partitioning and consistency, then consensus, and move to the design sections once you can say what each building block guarantees and what it gives up. Several ideas, quorums and read-your-writes in particular, return in more than one section at different depths, so expect to revisit them as your level rises.
primer
A few ideas carry the whole subject. Each section below applies them to a different problem, and most hard questions are one of them in disguise. - **Partial failure is the normal case.** Inside one process a function call either returns or the process dies. Across a network a request has a third outcome: no answer, which cannot tell a crashed node from a slow one, or a lost request from a lost reply. Retries, timeouts, failure detectors and consensus protocols all exist to act sensibly under that ambiguity, and interviewers reward candidates who say which outcome their design assumes. - **There is no shared now.** Physical clocks drift and are corrected, occasionally backwards, so timestamps from different machines are weak evidence of order. What can be tracked reliably is causality: which event could have influenced which. Logical and vector clocks capture that relation. Any design that orders events by wall-clock time should state what happens when two clocks disagree. - **Copies buy availability and create the consistency problem.** Replication exists to survive machine loss and to spread reads, and the moment a second copy exists you must decide which one answers and how far behind it may be. Whether a write is acknowledged before or after other copies have it decides whether an acknowledged write can be lost. - **Consistency is a set of named promises, not a yes-or-no.** Linearizable, sequential, causal, the per-client session guarantees, eventual: each promises something specific about what a read may return. Stronger promises need coordination, which costs latency on the requests that need it and, during a partition, availability. Precision about which promise you mean is what separates a senior answer from a vague one. - **Agreement rests on overlapping majorities.** Consensus, leader election and quorum-replicated stores share one mechanism: any two groups large enough to act must share a member, so conflicting decisions cannot both win. The cost is that a minority cannot make progress, and a deposed leader has to be stopped from acting on stale authority. - **Retries create duplicates; idempotency makes them harmless.** Over an unreliable network you choose between possibly losing a message and possibly delivering it twice. What systems advertise as exactly-once is an exactly-once *effect*: repeated delivery combined with processing that recognises a repeat. - **Scale comes from splitting work and splitting data.** Stateless tiers grow by adding copies behind a balancer; state grows by partitioning, where the key decides which queries stay cheap and where hot spots form. A capacity estimate tells you whether any of this is needed, and at what size.
- Partial failure
- A state where some components of a system have failed while others keep running, and the survivors cannot always tell which is which.
- Network partition
- A fault that splits nodes into groups that cannot reach each other, while each group may still be running and serving its own clients.
- Failure model
- The assumption a design makes about how components break: stopping, stopping and later restarting, losing messages, or behaving arbitrarily. Stronger faults need costlier protection.
- Replication lag
- The delay between a write landing on one copy and appearing on another; the source of stale reads in asynchronously replicated systems.
- Quorum
- The number of replicas that must take part before a read or write counts as done; sized so that groups overlap where the guarantee requires it.
- Linearizability
- A single-object guarantee that every operation appears to take effect at one instant between its start and its finish, consistent with real-time order.
- Serializability
- A transaction guarantee that concurrent transactions produce a result equal to some serial order; on its own it makes no promise about real-time order.
- Eventual consistency
- The guarantee that replicas converge to the same value once writes stop arriving, with no promise about what a read returns before that.
- Session guarantees
- Per-client promises layered over a weak store: read-your-writes, monotonic reads, monotonic writes and writes-follow-reads.
- Happens-before
- The partial order of events defined by causality: an earlier event on the same process, or a send before its receive, and anything reachable through those links.
- Vector clock
- A per-node counter array attached to events or versions; comparing two of them shows whether one causally precedes the other or they are concurrent.
- Consensus
- Getting a group of nodes to agree on one value, or one sequence of values, despite some members failing; the core of replicated logs and leader election.
- Term
- A monotonically increasing number identifying a leadership period, used by nodes to reject messages from a leader that has since been replaced. Some protocols call it an epoch.
- Split brain
- Two nodes both believing they are the leader, typically after a partition or a false failure detection, and both accepting writes.
- Fencing token
- A number that increases with each grant of a lock or lease; the protected resource refuses requests carrying a number older than one it has already seen.
- Two-phase commit
- An atomic commit protocol in which a coordinator collects votes from every participant before announcing a single decision, and blocks if the coordinator fails at the wrong moment.
- Idempotency
- The property that performing an operation more than once has the same effect as performing it once, which makes retries safe.
- Consistent hashing
- A way of mapping keys to nodes so that adding or removing a node moves only a small share of keys instead of nearly all of them.
- CRDT
- A replicated data type whose merge operation lets replicas accept concurrent updates independently and still converge to the same state.
The three sections are not parallel chapters. Theory supplies the guarantees, the consistency section makes them precise, and the design section spends them. ### Inside the theory Replication and partitioning are the two axes along which data is distributed: copies of the same data, and different data on different nodes. Most real systems do both, and most theory questions sit on one axis or where they cross. Consensus is how replicas agree on a leader and on the order of writes; it depends on failure detection to notice that a leader is gone, and on terms to stop the old one. Time feeds conflict handling: last-write-wins leans on physical timestamps, while [logical and vector clocks](/topics/found-distributed-systems-vector-logical-clocks) detect when two writes are truly concurrent. [Distributed transactions](/topics/found-distributed-systems-distributed-transactions) and delivery semantics are the same problem, acting once across machines that fail independently, seen from the database side and from the messaging side. ### From theory to guarantees The [consistency models](/topics/found-eventual-consistency) section names what a replicated system promises to a reader. Quorum sizing, read repair and background reconciliation are the mechanisms; linearizability, causal order and the session guarantees are the promises; CRDTs and application-level merges are what you use when you refuse to coordinate at all. CAP and its practical reading are the bridge between that section and the core theory. ### From guarantees to designs Most components in [system design](/topics/found-system-design) are theory decisions wearing practical names. A cache is a replica you have agreed may be stale. The data tier is replication and partitioning with a concrete key. A message queue is a delivery-semantics choice plus back-pressure. A distributed rate limiter is a shared counter that trades accuracy for coordination. Observability is how you reconstruct one request after it has crossed many machines. Design rounds reward candidates who can make that translation out loud.
- Fault Tolerance & Reliability →
The subject starts with what can break: failure models and partial failure frame every protocol and trade-off that follows.
- Replication →
Copies are the first tool for surviving failure, and replication lag introduces the stale-read problem the rest of the hub keeps returning to.
- Partitioning & Sharding →
The second axis of distribution: how data is split across nodes, how keys are placed, and where hot spots come from.
- Consistency Model Spectrum →
Gives names to the guarantees a replicated store can offer, so later answers can state which one they provide.
- Consensus Protocols →
Majority quorums, replicated logs and leaders: the machinery behind every strongly consistent component a design relies on.
- Capacity Estimation →
The entry point to design rounds: sizing the load first tells you which building blocks the design actually needs.
Treating a timeout as proof that a remote operation failed; it may have succeeded, so a blind retry of a non-idempotent call can apply it twice.
Saying "eventually consistent" without naming what a reader may observe meanwhile, or which session guarantee the product needs on top of it.
Reading CAP as a permanent choice of two out of three; the choice applies while a partition lasts, and latency is the everyday cost of stronger consistency.
Ordering events across machines by wall-clock timestamps, then using those timestamps to discard writes under last-write-wins without admitting that concurrent updates get lost.
Claiming exactly-once delivery over a lossy network instead of at-least-once delivery with idempotent or deduplicated processing.
Sizing a quorum so reads and writes need not overlap, then expecting a read to see the latest acknowledged write.
Relying on a lock or lease for safety with no fencing: a paused process can wake up after its lease expired and still write.
Choosing a partition key by even spread alone and ignoring access skew, leaving one celebrity key or one time range to overload a single node.
Sizing capacity from average traffic instead of peak plus headroom, or forgetting that replication multiplies the storage figure.
Proposing a distributed transaction across services without saying what happens when the coordinator or a participant fails mid-protocol.
The same handful of choices appears in nearly every section. Naming the one you are making, and what would change your mind, is usually worth more than the diagram. - **Consistency versus latency and availability.** Stronger guarantees need coordination: typically a round trip to a leader or a majority, and refusal to answer when that majority is unreachable. Pay it where a wrong answer is expensive, such as balances, locks and uniqueness, and avoid it where staleness is cheap. - **Synchronous versus asynchronous replication.** Waiting for copies before acknowledging protects acknowledged writes and slows every write; acknowledging first is fast and can lose the most recent writes on failover. - **Detection speed versus false positives.** A short failure timeout reacts quickly and wrongly condemns slow nodes, which can trigger needless failovers; a long one is calm and slow to notice real crashes. - **Coordination versus convergence.** You can prevent conflicts by routing every write through one decision point, or accept them and resolve later with timestamps, version vectors or mergeable types. The first costs throughput and availability, the second costs lost updates or application complexity. - **Even spread versus locality.** Hash placement balances load and scatters ranges; range placement keeps related keys together and invites hot spots. - **Freshness versus load.** Caches, replicas and edge copies remove load from the source in exchange for a staleness window and an invalidation problem. - **Synchronous calls versus queues.** A queue absorbs bursts and decouples failure, and in return adds delay, ordering questions and duplicate handling.
Certain shapes recur across the hub under different names. Recognising them is how you place an unfamiliar question quickly. - **Overlapping majorities.** Consensus commits, leader votes and quorum reads and writes all rely on two sufficiently large groups sharing a member. When a question asks why a number must exceed half, this is the pattern. - **One writer, many readers.** Single-leader replication, a primary with read replicas, and a cache in front of a database all funnel writes to one place and spread reads, then deal with the lag between them. - **Retry plus deduplication.** Idempotency keys on requests, sequence numbers on messages and deduplication tables on consumers turn at-least-once delivery into a single effect. - **Monotonic numbers as authority.** Terms, epochs, fencing tokens, version numbers and logical clocks all use a number that only grows to reject stale actors or order events without trusting wall time. - **Detect divergence, then repair.** Read repair, hinted handoff, background anti-entropy and gossip let replicas drift and then pull them back together. - **Split by key.** Sharding, consistent hashing, per-user rate limits and partitioned queues all route by a key, and all inherit the same hot-key problem. - **Absorb bursts with a buffer.** Token buckets, queues and back-pressure let a system accept short spikes while enforcing a sustained rate. - **Estimate before you build.** Throughput, storage and bandwidth estimates decide whether a design needs caching, sharding or fan-out at all; the [capacity estimation](/topics/found-system-design-capacity-estimation) section trains the arithmetic.
explore
- Core Distributed Theory64 questions
- Replication6 questions
- Consensus Protocols6 questions
- CAP & Consistency Models6 questions
- Partitioning & Sharding6 questions
- Time & Physical Clocks6 questions
- Fault Tolerance & Reliability5 questions
- Logical & Vector Clocks6 questions
- Leader Election6 questions
- Distributed Transactions (2PC/3PC)6 questions
- Idempotency & Delivery Semantics6 questions
- Failure Detection & Gossip5 questions
- Scalability & System Design419 questions
- Capacity Estimation5 questions
- Load Balancing6 questions
- Caching at Scale16 questions
- Data Tier18 questions
- Consistency & Availability5 questions
- Async Messaging & Queues6 questions
- Rate Limiting & Throttling6 questions
- Content Delivery12 questions
- Distributed Observability6 questions
- Search & Indexing22 questions
- Real-Time Fanout & Feeds24 questions
- Geospatial & Proximity11 questions
- Job Scheduling & Delayed Execution16 questions
- File & Object Storage21 questions
- ML System Design245 questions
- Consistency Models45 questions
- Read/Write Quorums5 questions
- Reconciliation & Session Guarantees6 questions
- Consistency Model Spectrum5 questions
- Convergence & Propagation6 questions
- Conflict Resolution6 questions
- CRDTs6 questions
- Linearizability vs Serializability6 questions
- Causal Consistency5 questions
- API Designskillanchors this topic
- Data Engineerroleanchors this topic
- DevOps / SRE Engineerroleanchors this topic
- PostgreSQL DBAroleanchors this topic
- Software Design & Architectureskillanchors this topic
- AI & Data Scientistrole
- AI Engineerrole
- Android Developerrole
- Backend Developerrole
- Blockchain Developerrole
- Computer Scienceskill
- Forward Deployed Engineerrole
- Frontend Developerrole
- Full Stack Developerrole
- Game Developerrole
- Java Backend Developerrole
- Kotlin Backend Developerrole
- MLOps Engineerrole
- Machine Learning Engineerrole
- MongoDBskill
- Redisskill
- SQLskill
- Server-Side Game Developerrole
- Software Architectrole
- System Designskill
- iOS Developerrole
questions
528 · 3 sectionsIn plain terms, what does the CAP theorem say a distributed database must sacrifice when a network partition occurs, and why can't a system just have all three of consistency, availability, and partition tolerance?
basics
~20 sWhen two halves of a system can't talk to each other, you must pick: either every request works but might return old/wrong data (available), or requests wait/fail until the data is provably correct (consistent). You can't have perfect answers and instant answers at once during that outage.
Two servers in different data centers each timestamp an event using their local system clock (wall-clock/time-of-day). Why can't you safely assume that whichever timestamp is numerically smaller happened first in real time?
basics
~20 sComputer clocks aren't perfectly synced - they drift apart between corrections and can even jump backward during a correction, so timestamps from two different machines can't be trusted to show which event really happened first.
In a Raft or Paxos cluster of N nodes, a write is considered committed once it has been acknowledged by a quorum. Why does that quorum have to be a strict majority (more than N/2) rather than some fixed count like 2 nodes, regardless of cluster size?
basics
~10 sA majority is needed so any two quorums always share at least one node - that overlap is what stops two different decisions from being made at once.
A distributed transaction coordinator needs to make sure a payment write to Database A and an inventory write to Database B either both happen or neither happens. Walk through how the two-phase commit (2PC) protocol achieves this, phase by phase.
basics
~10 s2PC asks everyone 'can you commit?' first (prepare phase). Only if all say yes does it tell them to actually commit (commit phase). If anyone says no, everyone rolls back.
A cluster node pings its peers every second and marks a peer 'dead' if it misses 5 heartbeats in a row (5 seconds of silence). What is the basic trade-off this fixed-timeout heartbeat approach faces when choosing the timeout value, and why can't a single value be 'correct' for all conditions?
basics
~10 sA short timeout catches crashes fast but often wrongly declares a slow-but-alive node dead. A long timeout avoids false alarms but takes longer to notice a real crash.
A scoring fleet takes 12,000 sensor readings per second at peak, each instance sustains 200, target utilisation 60% - how many instances?
basics
~10 sOne hundred instances. Dividing the 12,000-per-second peak by 200 gives 60 instances running flat out; dividing again by the 0.6 utilisation target gives 100, whose 20,000-per-second ceiling leaves peak sitting at 60%.
A bulk catalogue re-embedding job runs 400 accelerator-hours a month at $3 an hour for 20 million embeddings - what is the cost per embedding?
basics
~10 sSix hundredths of a cent. 400 hours at $3 is $1,200 of compute, divided by 20 million embeddings gives $0.00006 each, or $0.06 per thousand. That figure covers marginal compute only.
In a video-moderation service, which half of the spend grows with upload volume: the quarterly training run or per-clip scoring?
basics
~20 sPer-clip scoring grows with upload volume; the quarterly training run does not. Training is paid once per model version, while scoring is paid again for every clip that arrives, so only the serving half is multiplied by traffic.
What must a training run's record hold before two quarterly claim-severity runs can be compared at all?
basics
~20 sA run record must pin every input and the yardstick: the dataset snapshot identifier, the code commit, the full parameter set, a content digest of the model artifact, and the evaluation set identifier with its exact metric definition. Two runs compare only when that evaluation set and metric are identical.
In a content-moderation labeling operation, why is each reported post judged by three annotators instead of one?
basics
~20 sOne verdict on a moderated post carries no error signal of its own. Three verdicts make disagreement visible, route the contested post into an adjudication path, and turn agreement into a running health check on the written guideline.
In a distributed key-value store, what does 'causal consistency' guarantee about the order in which different replicas observe writes, and how does that differ from plain eventual consistency?
basics
~20 sCausal consistency guarantees that writes which are causally related (one happened because of, or after seeing, another) are seen by every replica in that same order. Unrelated writes may appear in different orders on different replicas. Plain eventual consistency promises no ordering at all, just eventual agreement.
In an eventually consistent data store, two replicas each accept a write to the same key while a network partition separates them. When the partition heals and the store uses last-write-wins (LWW) to reconcile, what happens to the value, and what has to happen for that to work reliably?
basics
~20 sLast-write-wins keeps whichever write has the newer timestamp and throws away the other one. For this to work, every replica needs a timestamp on each write, and those timestamps must be trustworthy enough to compare across machines.
In a distributed database, what does 'eventual consistency' mean, and how is it different from a system that guarantees you always read the latest write?
basics
~10 sEventual consistency means all copies of data will match up eventually, but right after a write some copies might still be stale. Strong consistency means everyone sees the newest value immediately.
In an eventually consistent distributed database, what does it mean for replicas to 'converge' after a write, and what is one basic mechanism that helps them get there?
basics
~20 sConvergence means all copies of the data eventually end up matching each other, even if some copies see the update late. One simple way this happens is copies periodically comparing notes and fixing whichever one has stale data.
What problem does a CRDT (Conflict-free Replicated Data Type) solve, and how does a G-Counter's merge rule guarantee replicas converge without any coordination between nodes?
basics
~20 sA CRDT lets many computers update the same piece of data at the same time without talking to each other, and their copies can always be combined back into one correct answer later. A G-Counter only grows, so merging two copies means keeping the bigger count each machine reported.