System Design Interview

Course Content

System Design Interview

31 sections · 71 lessons

Building blocks: partitioning, caching, load balancing and queues


Two orthogonal ideas get confused constantly. Partitioning splits data so it does not all live on one machine. Replication copies data so its loss or unavailability does not end the system. Most real systems do both, and the two decisions are independent.

This lesson covers the building blocks that move data and traffic around a system. It starts with partitioning and replication, including the quorum arithmetic behind tunable consistency. Then caching strategies, which decide who fills the cache and when. Then load balancing, where the questions are which layer, which algorithm and how a dead server is detected. It ends with queues and logs, where the most consequential distinction is between the two.

Partitioning strategies

Range partitioning. Keys are split into contiguous ranges: A–F on shard 1, G–M on shard 2. Range scans stay efficient because neighbouring keys live together. The failure is hot spots: partition by date and today's shard receives every write while the rest idle.

Hash partitioning. shard = hash(key) mod N. Distribution is even and range scans become impossible — neighbouring keys land on unrelated shards. The mod N form has a catastrophic resizing property: change N and almost every key moves. Section 7 (Design Consistent Hashing) is entirely about the fix.

Directory-based partitioning. A lookup service stores an explicit map from key or key range to shard. Maximum flexibility — you can move one hot key without touching anything else — at the cost of an extra hop on every request and a service that must not go down. It is usually replicated and cached everywhere.

Consistent hashing. A hash-partitioning scheme where adding or removing a node moves only about 1/N of the keys instead of nearly all of them. Section 7 covers it in full.

Whichever you pick, the partition key choice is what actually decides the design: it determines which queries are single-partition (fast) and which fan out to every shard (slow), and it determines whether load is even or concentrated.

Replication topologies

Single-leader (leader-follower). All writes to one node, which streams changes to followers; reads from any. Simple, no write conflicts ever, and it is what most relational deployments do. Limits: write throughput is one machine's, and there is a failover gap when the leader dies.

Multi-leader. Several nodes accept writes and replicate to each other. Used for multi-region deployments where local write latency matters, and for offline-capable clients. Its cost is real and unavoidable: two leaders can accept conflicting writes to the same key, so you need a conflict resolution rule (Conflict resolution, Section 8).

Leaderless (quorum). Any node accepts a write; the client or a coordinator writes to several nodes and reads from several. No failover, because there is nothing to fail over. The cost is that consistency becomes a tunable rather than a guarantee, and repair mechanisms (read repair, anti-entropy) are needed to converge.

Synchronous versus asynchronous

Independent of topology. Synchronous replication means the write is not acknowledged until a replica has it: no data loss on leader failure, but every write pays the slowest replica's round trip — 0.5 ms in one datacentre, 100 ms+ across regions. Asynchronous means the leader acknowledges immediately: fast, and writes that were acknowledged but not yet shipped are lost if the leader dies. Semi-synchronous — wait for one replica, not all — is the common compromise.

Quorums, and why R + W > N gives strong reads

In a leaderless system with N replicas, a write is acknowledged after W of them confirm, and a read collects responses from R of them.

If R + W > N, the read set and the write set must overlap by at least one node, so at least one node in every read has the newest value. Add a rule for picking the newest among the responses (a version number or timestamp) and the read is strongly consistent.

Worked, with N = 3:

ConfigurationOverlap?Behaviour
W = 2, R = 22 + 2 > 3 yesStrong reads; tolerates one node down for both reads and writes
W = 3, R = 1yesFast reads, but any node down blocks all writes
W = 1, R = 11 + 1 > 3 noFastest and most available; eventually consistent

Latency is set by the W-th fastest response, not the average. With replica latencies of 2 ms, 5 ms, and 40 ms, W = 2 costs 5 ms and W = 3 costs 40 ms. That single fact explains why W = 3 is rare.

Partitioned, not replicatedcapacity scales; one node down loses its shardA–HI–PQ–Z3× capacity · 0 redundancyReplicated, not partitionedevery node holds everythingFullFullFull1× capacity · survives 2 failuresPartitioned and replicatedwhat production actually runsA–HA–HI–PI–PQ–ZQ–Z3× capacity · survives 1 failure per shardThe two questions every partitioning scheme has to answerHow is the key chosen?A bad partition key concentrates traffic on one shard."user_id" is usually safe; "country" is usually not,because one country is 40% of traffic.What happens when a node joins?Modulo hashing remaps almost every key. Consistenthashing (Lesson 7) moves roughly 1/n of them, which isthe difference between a rolling change and an outage.Who accepts writes?One leader per shard is simple and caps write throughputper shard. Multi-leader removes the cap and buys youconflict resolution.Partitioning buys capacity; replication buys availability. Almost every real system needs both, and they are chosen independently.
The third topology is the only one that both scales and survives a failure — and it is the one that makes rebalancing hard.

Caching strategies

Partitioning and replication decide where the data lives. A cache is a fast store holding copies of data that is expensive to produce, and the strategy is about who populates it and when. There are five, each right in different circumstances.

Who fills the cache, and whenCache strategyCache-aside: the appRead-through: cacheWrite-through: syncWrite-back: asyncRefresh-ahead: guess
Write-back is the only one that can lose an already-acknowledged write, and that is the whole trade.

Cache-aside (lazy loading). The application checks the cache; on a miss it reads the database, writes the result into the cache, and returns it. The cache does not know about the database.

  • Right when: reads dominate and you can tolerate the first request for each key being slow.
  • Cost: every key is cold once. After a cache restart, every key is cold at once (the stampede from the caching lesson in Section 2).
  • This is the default. Roughly 80% of interview systems should use it.

Read-through. The same behaviour, but the cache library does the database read itself. The application only ever talks to the cache.

  • Right when: you want the logic in one place and your cache supports a loader.
  • Cost: the cache must know how to fetch, which couples it to your data source.

Write-through. Every write goes to the cache and the database synchronously.

  • Right when: data is read soon after it is written, and staleness is unacceptable.
  • Cost: every write pays both latencies, and you cache data that may never be read. Combine with a TTL so unread entries expire.

Write-back (write-behind). Write to the cache, acknowledge immediately, flush to the database later in batches.

  • Right when: write volume is very high and small losses are tolerable — view counters, "last seen" timestamps.
  • Cost: an unflushed cache node that dies loses acknowledged writes. Do not use this for anything a user would notice losing.

Write-around. Writes go straight to the database, bypassing the cache entirely.

  • Right when: written data is rarely read back soon — log ingestion, bulk imports.
  • Cost: the first read after a write always misses.

Choosing a caching strategy

WorkloadStrategy
Read-heavy, tolerant of a cold first readCache-aside
Read-heavy, want the fetch logic centralisedRead-through
Written then immediately read, staleness unacceptableWrite-through
Extremely write-heavy, small loss tolerableWrite-back
Write-heavy, rarely read backWrite-around

Mixing is normal: cache-aside for reads plus write-around for a bulk import path is a common and sensible pairing.

Eviction policies

When the cache is full, something must go.

  • LRU (least recently used): discard the item untouched for longest. The sensible default; it matches how access is usually distributed.
  • LFU (least frequently used): discard the least-accessed item. Better for a stable hot set, worse when popularity shifts — it clings to yesterday's favourites.
  • FIFO: discard the oldest inserted. Cheap, and usually worse than LRU.
  • TTL-based: every item expires at a fixed age regardless of use. Not really eviction; it is a freshness policy, and it is usually combined with one of the above.
  • Random: surprisingly competitive, and very cheap. Some production caches use approximate LRU precisely because exact LRU costs bookkeeping.

The invalidation problem, stated plainly

The data changed. The cache does not know. Three ways to handle it:

  1. TTL only. Accept that data is stale for up to the TTL. Trivially simple, and the right answer far more often than candidates expect. State the window: "up to 60 seconds stale."
  2. Delete on write. After writing the database, delete the cache key so the next read repopulates. Correct in the common case, and it has a real race: a concurrent read that missed before your write can write the old value into the cache after your delete. The window is small but non-zero. Mitigations include a short TTL as a backstop, or versioned keys.
  3. Versioned keys. Include a version in the key — route:8213:v7. Nothing is ever invalidated because a new version is a new key; old entries fall out by eviction. Clean, and it needs a reliable version source.

The often-quoted remark that cache invalidation is one of the two hard problems in computer science is usually attributed to Phil Karlton; the attribution is folklore, but the point stands.

Load balancing: layer 4 versus layer 7

A load balancer distributes incoming requests across a pool of servers. The interesting questions are which layer it works at, how it picks a server, and how it decides a server is dead.

Layer 4 against layer 7Layer 4• Routes on IP and port alone• Cheap, fast, protocol-agnostic• Cannot route by path or by headerLayer 7• Reads the HTTP request itself• Routes by path, host or cookie• Terminates TLS, so it costs CPU
Health-check interval times failure threshold is your outage length, and nobody computes it until it matters.

Layer 4 (transport layer) balancers route based on internet protocol (IP) address and port. They forward Transmission Control Protocol (TCP) packets without reading the request.

  • Very fast and cheap — a single node handles millions of concurrent connections.
  • Cannot route by path, header, or cookie, because it never reads them.
  • Cannot retry a failed request, because it does not know what the request was.
  • Right for raw TCP protocols, extreme throughput, and the outermost layer of a large system.

Layer 7 (application layer) balancers read the Hypertext Transfer Protocol (HTTP) request.

  • Route /api/* to one pool and /static/* to another.
  • Terminate Transport Layer Security (TLS) so backends do not each need certificates.
  • Retry a failed idempotent request on another server, and shed load intelligently.
  • Compress, rewrite, and add headers.
  • Costs more CPU per request, and adds a small amount of latency.

Most systems use both: layer 4 at the edge for volume, layer 7 behind it for routing.

Routing algorithms

AlgorithmHow it picksUse when
Round robinNext server in the listServers identical, requests uniform
Weighted round robinProportional to a configured weightMixed instance sizes
Least connectionsFewest in-flight requestsRequest durations vary a lot
Least response timeLowest recent latencyBackend performance is uneven
IP hash / consistent hashHash of client IP or a keyYou want the same key on the same server (cache affinity)
Power of two choicesPick two at random, take the less loadedLarge pools — nearly as good as least-connections at a fraction of the coordination

Round robin is the correct default. Least connections is the correct answer when request cost varies widely — with round robin, a server that draws three slow requests keeps receiving new ones at the same rate as an idle peer.

Health checks, and the arithmetic of failure detection

Active checks: the balancer requests a known path (say /healthz) on an interval, and removes a server after a threshold of consecutive failures.

The detection window is the product of the two settings. A 5-second interval with a 3-failure threshold means up to 15 seconds during which a dead server keeps receiving traffic. At 1,700 requests per second across ten servers, that is 170 requests per second failing for 15 seconds — roughly 2,500 failed requests per incident. Tightening to a 2-second interval and 2 failures reduces that to about 680, at the cost of more false removals during a transient blip.

Passive checks (outlier ejection): the balancer watches real traffic and ejects a server that returns errors or times out. Detection is immediate because it uses requests that were happening anyway. Best used alongside active checks, not instead of them.

Distinguish two check types on the server side: liveness ("is the process alive?", restart if not) and readiness ("can it serve right now?", remove from rotation if not). Conflating them causes a warming server to be killed instead of temporarily removed.

Reverse proxies, API gateways, and sidecars

Reverse proxy: a server that receives client requests and forwards them to backends. Every layer 7 load balancer is a reverse proxy; the term emphasises TLS termination, compression, static file serving, and response caching.

API gateway: a reverse proxy with product concerns attached — authentication, rate limiting (Section 6), request validation, usage metering, versioning, and per-route routing to different services. In a microservice system it gives clients one entry point instead of twenty. Its risk is becoming a bottleneck and a place where business logic accumulates.

Service mesh sidecar: a proxy deployed next to every service instance, handling service-to-service traffic — retries, timeouts, mutual TLS, and traces — without the application implementing them. It pushes the reliability patterns into infrastructure. Its cost is an extra network hop each way and a large operational surface.

Global load balancing sits above all of this: DNS-based routing or anycast IP addresses send a user to the nearest healthy region. DNS-based routing is limited by record caching — clients hold a stale answer for the TTL, so failover is not instant.

Message queues and event streaming: queue versus log

A load balancer spreads synchronous requests. A queue decouples a producer from a consumer, so a slow or absent consumer does not slow the producer down. The most consequential distinction here is between a queue and a log.

A log is read, not consumedoff 0off 1off 2off 3off 4nullgroup Agroup BReading does not delete; each consumer group owns its own offset.
A queue removes on read and cannot replay; a log keeps the bytes, so a second consumer can start again at zero.

Queue. Messages are added, delivered to one consumer from a pool of competing consumers, and deleted once acknowledged. The queue is a buffer for work.

  • Natural fit for task distribution: resize this image, send this email.
  • Consumers scale by adding more of them; each message is handled once.
  • Once consumed, a message is gone. No replay, no second consumer.
  • Common implementations: RabbitMQ, Amazon SQS, and similar brokers.

Log. Messages are appended to an ordered, immutable sequence and retained for a configured period regardless of consumption. Each consumer group tracks its own position (offset) in the log.

  • Multiple independent consumers read the same stream: one writes to the search index, another updates analytics, a third feeds a cache. None affects the others.
  • A consumer can rewind and reprocess — after a bug, or to populate a new service.
  • Ordering is guaranteed within a partition (see below).
  • Storage grows with retention rather than with backlog.
  • The most common implementation is Apache Kafka; Apache Pulsar and cloud equivalents follow the same model.

The log model won for event-driven architectures, because replay and independent consumers turn out to be worth more than automatic cleanup. Section 21 (Design a Distributed Message Queue) builds one from scratch.

Partitions and ordering

A log topic is split into partitions for throughput; each partition is an independent ordered sequence handled by one consumer in a group at a time.

The consequence is precise and frequently examined: ordering is guaranteed within a partition and nowhere else. If you need all events for one user processed in order, the partition key must be the user identifier. Choose the key badly and either ordering breaks or one partition becomes hot while the rest idle.

Delivery semantics

At-most-once. Acknowledge before processing. A crash mid-processing loses the message. Use only where loss is acceptable — some telemetry.

At-least-once. Acknowledge after processing. A crash after processing but before acknowledging causes redelivery, so the message is handled twice. This is what almost every real system uses.

Exactly-once. Every message has exactly one effect.

Consumer groups and backpressure

A consumer group is a set of consumers sharing the work of a topic, with each partition assigned to exactly one member. Adding a consumer triggers a rebalance, which reassigns partitions and briefly pauses processing. A topic with 12 partitions supports at most 12 useful consumers in a group — partition count is your maximum parallelism, and it is awkward to change later, so pick it generously.

Backpressure is what happens when consumers cannot keep up. Producers keep appending, and consumer lag — the gap between the newest offset and the consumed offset — grows. The retention window is your buffer: if it is 7 days and lag is 6 days, you have one day to fix it before data is lost permanently.

Size it. At 10,000 messages per second and 1 KB each, that is 10 MB/s, or 864 GB a day, so a 7-day retention needs about 6 TB before replication — 18 TB at replication factor 3. That arithmetic is what turns "we'll retain a week" into a decision.

Responses to lag: add consumers (up to the partition count), make processing cheaper, shed low-priority messages, or push backpressure to producers by rejecting writes. Doing nothing means silent data loss when retention expires.