System Design · Lesson 2 of 18

Every Term, Explained

Scaling, sharding, caching, replication, quorums, idempotency and 40 more terms — one by one.

  • Intermediate
  • 30 min read
  • 3 objectives

Before this lessonLesson 1: System Design Fundamentals

What you will learn

  • Define every term interviewers use
  • Know which problem each one solves
  • Pick the right tool per bottleneck

Your Progress

0 of 18 lessons 0%

  • Lessons0 / 18
  • Completed0
  • Est. time left~ 16 hours

Create a free account to keep your progress on every device.

Tip: pressing Next marks this lesson complete automatically.

Interviewers throw around forty or fifty words as if everyone agrees on them: shard, partition, quorum, idempotent, backpressure. This lesson defines them one at a time. For each term: what it means, the problem it solves, and the cost you pay. Read it once end to end, then come back to it as a dictionary while you work through the case studies.

Everything below fits one mental model. A request travels client → DNS → CDN → load balancer → app server → cache → database, and every term names either a stage on that path, a way of multiplying a stage, or a way of surviving a stage failing.

flowchart LR C[Client] --> DNS[DNS / GSLB] DNS --> CDN[CDN edge] CDN --> LB[Load balancer] LB --> GW[API gateway] GW --> APP[App servers
stateless] APP --> CACHE[(Cache)] APP --> DB[(Database
primary)] DB --> REP[(Replicas)] DB --> SH[(Shards)] APP --> Q[(Queue)] Q --> W[Workers] APP --> OBJ[(Object storage)]

Every term in this lesson is a box on this path, a way to duplicate a box, or a way to keep working when a box dies.

1. Scaling

Vertical scaling (scale up)

Buy a bigger machine: more CPU, more RAM, faster NVMe. Solves: load, with zero code change — this is genuinely the right first answer for most systems. Costs: a hard ceiling (there is a largest instance), price grows faster than capacity, and one machine is one failure domain. You still need a second machine for availability.

Horizontal scaling (scale out)

Run many identical machines behind a load balancer. Solves: unbounded capacity and redundancy. Costs: you must now handle coordination — where does session state live, how do you find a peer, how do you avoid two workers doing the same job. Horizontal scaling is easy for stateless servers and hard for databases; that asymmetry drives most of system design.

Stateless vs stateful

A stateless service keeps nothing important in local memory or disk: any server can serve any request, so you can add, kill or restart machines freely. A stateful service owns data (a database, a cache node, a WebSocket connection) so requests care which instance they reach. Make app servers stateless by pushing session data into Redis or a signed token (JWT) and files into object storage.

Sticky sessions (session affinity)

The load balancer pins a client to one server (by cookie or IP hash). Solves: in-memory session or an open WebSocket. Costs: uneven load, and a deploy or crash drops that user's state. Prefer statelessness; use stickiness only for long-lived connections.

Bottleneck

The single resource that saturates first — CPU, memory, disk IOPS, network, database connections, or a lock. Scaling anything else changes nothing. Interviewers want you to name the bottleneck before proposing a fix.

2. Getting the request to a server

DNS and GSLB

DNS maps a name to IPs and is your cheapest load balancer: return several A records, or return a different IP per region (GSLB / geo-DNS / anycast) so users hit the nearest datacenter. Costs: TTL caching means DNS failover takes minutes, not seconds, so never rely on DNS alone for fast failover.

Load balancer

Spreads requests over healthy servers and hides individual failures. Layer 4 balances TCP/UDP by IP and port — fast, protocol-blind. Layer 7 reads HTTP, so it can route by path or header, terminate TLS, retry, and compress — flexible, slightly slower. Algorithms: round robin (simple), least connections (better with uneven request cost), least response time, random with two choices (nearly as good as least-connections without global state), and hash on a key (sends the same user or key to the same server).

Health check

The balancer polls /health and removes failing instances. Make it check real dependencies shallowly — a health check that pings the database can take your whole fleet out of rotation when the database hiccups. Separate liveness (am I alive — restart me if not) from readiness (should I get traffic right now).

Reverse proxy and API gateway

A reverse proxy (nginx, Envoy) sits in front of app servers to terminate TLS, cache, compress and buffer slow clients. An API gateway adds product concerns: authentication, rate limiting, quotas, request validation, routing to many microservices, API versioning and usage metering. One place to enforce cross-cutting rules; also one more hop and one more thing that can fall over.

CDN (content delivery network)

Hundreds of edge locations cache your static assets — images, JS, video segments — near users. Solves: latency (100 ms across an ocean becomes 10 ms), and bandwidth cost on your origin. Costs: invalidation is now a problem (fix it with content-hashed filenames like app.7f3a91.js and long Cache-Control), and per-GB egress is a real line item.

Object storage (blob store)

S3, GCS, Azure Blob: cheap, effectively infinite, durable storage for immutable files, addressed by key, served through a CDN. Never store photos or videos in your relational database — store the bytes in object storage and the key plus metadata in the DB. Uploads go straight from the client to the store with a pre-signed URL, so the bytes never touch your app servers.

3. Caching

Cache

A small, fast copy of data you will read again soon. Solves: read latency and database load — in a 100:1 read-heavy product, a 90% hit rate removes 90% of database reads. Costs: staleness, memory cost, one more failure mode, and a whole class of consistency bugs.

The cache layers

  • Browser / client cache: free, closest, hardest to invalidate.
  • CDN / edge cache: static and public content.
  • Application in-process cache: nanoseconds, but each server has its own copy, so N servers means N versions of the truth.
  • Distributed cache (Redis, Memcached): sub-millisecond, shared by all servers — the workhorse.
  • Database buffer pool / query cache: the DB caching its own hot pages.

Hit rate, hot key, working set

Hit rate is the fraction of reads served by cache. The working set is the data actively being read — size your cache for that, not the whole dataset (Pareto: 20% of items get 80% of traffic). A hot key is one item so popular it saturates a single cache node (a celebrity profile, a viral link); fix it by replicating that key to several nodes or caching it locally in every app server for a few seconds.

Eviction policies

LRU (least recently used) — the sane default. LFU (least frequently used) — better when popularity is stable, resists a scan wiping your cache. FIFO — cheap, worse. TTL — every entry expires after N seconds, which bounds staleness regardless of policy.

Read patterns

Cache-aside (lazy loading): app checks the cache, on a miss reads the DB and writes the cache. Only requested data is cached; every miss pays double latency. Read-through: the cache library does it for you. See lesson 3 for code.

Write patterns

  • Write-through: write cache and DB together. Cache is always warm and fresh; writes are slower.
  • Write-around: write only the DB and invalidate the key. Best when written data is rarely read straight back.
  • Write-back (write-behind): write the cache, flush to the DB asynchronously. Fastest writes, and you lose data if the cache node dies before the flush.

Cache stampede, thundering herd, penetration

A popular key expires and a thousand concurrent requests all miss and all hit the database at once — a stampede. Fixes: add random jitter to TTLs so keys do not expire together; take a short per-key lock so one request refills while others wait or serve stale; refresh hot keys in the background before expiry. Cache penetration is repeated lookups for keys that do not exist — cache the negative result, or put a bloom filter in front.

4. Databases: replication

Replication

Keep copies of the same data on several machines. Solves: read capacity, durability, and failover. Costs: copies can disagree.

Leader–follower (primary–replica)

All writes go to one leader, which streams its change log to followers that serve reads. Synchronous replication waits for a follower to acknowledge — no data loss on failover, higher write latency, and a stalled follower stalls writes. Asynchronous is fast but a leader crash loses the last few writes. Common middle ground: semi-synchronous — one synchronous follower, the rest async.

Replication lag and read-your-writes

A follower is milliseconds to seconds behind. The classic bug: a user posts a comment, the next read hits a lagging replica, and their own comment is missing. Fixes: read from the leader for that user for a few seconds, route a session's reads to one replica (monotonic reads), or track a write timestamp / log position per session and only use a replica that has caught up.

Failover, split brain, fencing

Failover promotes a follower when the leader dies — automatically (fast, risky) or manually (safe, slow). Split brain is two nodes both believing they are leader, each accepting writes; you resolve it with a majority-based election and fencing (a monotonically increasing term or epoch number so the old leader's writes get rejected).

Multi-leader and leaderless

Multi-leader: several regions accept writes, so local writes are fast — but concurrent edits to the same row conflict and must be merged (last-write-wins, per-field merge, CRDTs, or ask the user). Leaderless (Dynamo, Cassandra): the client writes to several replicas directly and reads from several, using quorums.

Quorum: N, W, R

With N replicas, a write waits for W acknowledgements and a read collects R responses. If W + R > N, any read overlaps at least one up-to-date replica, giving strong consistency. N=3, W=2, R=2 is the standard balanced choice; W=1 gives fast writes and stale reads; R=1 gives fast reads and possibly stale data. Hinted handoff parks writes for a down replica on a neighbour; anti-entropy (Merkle-tree comparison) repairs divergence later.

5. Databases: partitioning and sharding

Partitioning vs sharding

Partitioning is splitting one logical table into pieces — it can be vertical (some columns to another table or service — move that rarely-read bio TEXT off the hot row) or horizontal (some rows to another place). Sharding is horizontal partitioning across separate machines, each with its own CPU, memory and disk. Loosely: partitioning splits data, sharding splits data and the hardware. Solves: a dataset or write rate one machine cannot hold. Costs: cross-shard queries, cross-shard transactions, rebalancing, and hot shards — this is the most expensive decision in the lesson, so do it last.

Shard key

The column used to route a row (user_id, tenant_id, video_id). Pick one that (a) appears in nearly every query, so you hit one shard, and (b) spreads traffic evenly. Sharding a chat app by user_id means one user's messages are in one place; sharding it by message_id means every conversation read touches every shard.

Sharding strategies

  • Range-based: A–F on shard 1, G–M on shard 2. Range scans stay local; skewed data creates hot shards (everyone named "S", or all of today's rows if you shard by time).
  • Hash-based: shard = hash(key) % num_shards. Even distribution, no range scans, and changing num_shards moves almost all data.
  • Consistent hashing: place shards and keys on a hash ring; a key belongs to the next node clockwise. Adding or removing a node moves only 1/N of keys. Virtual nodes (each physical node holds ~100–300 ring positions) smooth out imbalance and let heterogeneous machines carry different weights. Used by Cassandra, DynamoDB and Memcached clients.
  • Directory-based: a lookup service maps key → shard. Maximum flexibility (move a big tenant to its own shard), plus one more service in the critical path.
  • Geo / entity-based: shard by region or tenant, for latency and data-residency rules.

Hot spot, skew, resharding

A hot spot is one shard or key taking a disproportionate share of traffic — usually a celebrity, a viral item, or a monotonically increasing key that sends every new write to the last shard. Fixes: salt the key (celebrity_id:bucket_0..15), split that key's data, cache it hard, or dedicate a shard. Resharding is splitting shards as you grow: doubling shard count, or pre-creating many logical shards (say 1,024) and mapping several onto each physical machine so growth is a move, not a rehash.

Federation and celebrity/fan-out problems

Federation splits by feature into separate databases (users DB, orders DB, analytics DB) — easy early wins, no cross-DB joins. Fan-out is one write producing many follow-up writes (a tweet delivered to 40 million followers); hybrid fan-out — push for normal users, pull for celebrities — is the standard answer (lesson 11).

6. Databases: modelling and engines

SQL vs NoSQL

Relational databases give you a fixed schema, joins, secondary indexes and real transactions — the right default whenever you have relationships and need correctness. NoSQL is a family: key-value (Redis, DynamoDB), document (MongoDB), wide-column (Cassandra, HBase — huge write volume, queries by known key), graph (Neo4j — relationships as first-class), time-series (InfluxDB, Timescale), search (Elasticsearch — inverted index, full-text, ranking). Choose NoSQL for a stated reason: write volume, flexible shape, horizontal scale, or a specific access pattern — not because it is new.

ACID vs BASE

ACID: Atomicity (all or nothing), Consistency (invariants hold), Isolation (concurrent transactions do not corrupt each other), Durability (committed means on disk). BASE: Basically Available, Soft state, Eventually consistent — the relaxed contract most distributed stores offer.

Isolation levels

From weak to strong: read uncommitted (dirty reads), read committed (the common default), repeatable read / snapshot (no changing rows mid-transaction; write skew still possible), serializable (as if transactions ran one at a time — correct and slowest). Know the anomalies by name: dirty read, non-repeatable read, phantom read, lost update, write skew.

Index

A separate sorted structure mapping column values to rows, turning a full scan into a lookup. Costs: every write must update every index, and indexes take space. A composite index on (user_id, created_at) serves "this user's recent rows" in one range read — column order matters. A covering index contains every column the query needs, so the table itself is never touched.

Normalization and denormalization

Normalized means each fact is stored once — no update anomalies, but reads need joins. Denormalized duplicates data (storing author_name next to each post, or a pre-computed timeline row) so a read is one lookup. Reads get fast; writes get more expensive and you own the job of keeping copies in step. Denormalize deliberately, for a measured read path.

B-tree vs LSM tree

B-tree (Postgres, MySQL/InnoDB): updates pages in place, excellent reads and range scans, random write I/O. LSM tree (Cassandra, RocksDB, LevelDB): appends to an in-memory memtable, flushes sorted SSTables to disk, merges them during compaction — superb write throughput, reads may touch several files, and compaction costs background I/O. Write-heavy workload → LSM; read and range heavy → B-tree.

WAL, bloom filter, checksum

A write-ahead log records every change before applying it, so a crash can be replayed — it is also the stream replication and change-data-capture read from. A bloom filter is a small bit array that answers "is this key definitely absent, or probably present?" with no false negatives — it saves LSM engines and caches from pointless disk reads. A checksum detects silent corruption on disk or over the wire.

MVCC and locking

MVCC (multi-version concurrency control) keeps multiple row versions so readers never block writers. Pessimistic locking takes a lock up front (SELECT ... FOR UPDATE) — correct under contention, risks deadlocks. Optimistic locking reads a version number and fails the write if it changed (UPDATE ... WHERE version = 7) — better when conflicts are rare.

Distributed transactions: 2PC and sagas

Two-phase commit has a coordinator ask every participant to prepare, then commit — atomic across services, but it blocks if the coordinator dies and it scales poorly. A saga replaces one distributed transaction with a sequence of local transactions plus compensating actions (charge the card; if dispatch fails, refund). Sagas are eventually consistent and are what real systems use. The outbox pattern makes this safe: write the row and the event to the same database transaction, and a separate process publishes the event — no lost or phantom messages.

7. Consistency and the CAP trade-off

CAP

During a network partition you must choose: refuse requests to stay consistent (CP — payments, inventory, bookings) or answer with possibly stale data to stay available (AP — feeds, likes, recommendations). "CA" is not an option on a real network. PACELC extends it: else, when there is no partition, you still trade latency against consistency, because stronger guarantees mean more round trips.

Consistency models, strongest first

  • Linearizable / strong: every read sees the latest committed write, as if there were one copy.
  • Sequential / causal: operations that depend on each other are seen in order by everyone (your reply never appears before the message it replies to).
  • Read-your-writes: you always see your own changes; others may lag.
  • Monotonic reads: you never see time go backwards.
  • Eventual: given no new writes, all replicas converge — eventually.

Name the model per feature, not per system: a shopping app can be eventual for the product-view counter and linearizable for the checkout.

Idempotency

An operation is idempotent when doing it twice has the same effect as doing it once. Since every network call can time out and be retried, this is the property that keeps you from double-charging a card. Implement it with a client-supplied idempotency key: store key → result, and on a repeat return the stored result instead of re-executing. GET, PUT and DELETE are naturally idempotent; POST is not, so it needs a key.

Exactly-once, at-least-once, at-most-once

Messaging gives you at-least-once (may duplicate) or at-most-once (may lose). True exactly-once delivery is impossible across a network; exactly-once processing is achievable by combining at-least-once delivery with idempotent consumers or a dedupe table.

8. Asynchronous work

Message queue and log

A queue (SQS, RabbitMQ) hands each message to one consumer and then forgets it. A log (Kafka, Kinesis) is an append-only ordered record that many independent consumer groups read at their own offset, and can replay. Solves: decoupling producers from consumers, absorbing traffic spikes, retries, and keeping slow work off the request path. Costs: eventual consistency, duplicate delivery, ordering caveats, and operational weight.

Partition, consumer group, offset

A Kafka topic is split into partitions; order is guaranteed only within a partition, so choose a partition key (say conversation_id) that puts things needing order together. Each consumer group tracks its own offset, so adding a new consumer never steals another's messages.

Dead letter queue and poison pill

A message that keeps failing (a poison pill) would block or loop forever. After N attempts, move it to a dead letter queue for humans, and keep the pipeline flowing.

Backpressure

When downstream cannot keep up, the system must push back — bounded queues, rejecting with 429, slowing producers, or shedding low-priority load. Without backpressure, queues grow without bound, memory fills, latency explodes and the failure arrives all at once instead of gracefully.

Batching, coalescing, debouncing

Group many small operations into one to amortise overhead: batch 1,000 events into a single insert, coalesce repeated invalidations of the same key, debounce a burst of edits into one notification. Throughput improves, per-item latency gets slightly worse.

CDC and event sourcing

Change data capture tails the database's WAL and publishes row changes, so caches, search indexes and analytics stay current without dual writes. Event sourcing goes further: the log of events is the source of truth, and current state is a projection you can rebuild.

9. Reliability and failure

Timeout, retry, jitter

Always set a timeout — an unbounded wait converts one slow dependency into a fleet-wide outage as threads pile up. Retry only idempotent operations, with a bounded count and exponential backoff plus jitter (random(0, base·2^attempt)), so clients do not resynchronise into a retry storm. A retry budget caps retries as a fraction of total traffic.

Circuit breaker, bulkhead, load shedding

A circuit breaker watches a dependency's error rate; past a threshold it opens and fails fast for a cooldown, then lets a trial request through (half-open) before closing. Bulkheads give each dependency its own connection and thread pool, so one slow service cannot consume every worker. Load shedding drops the least valuable requests at the edge to keep the rest healthy; graceful degradation serves a reduced feature (a cached, unranked feed) instead of an error page.

Single point of failure and failure domains

Any component whose death takes the system with it. Find it by asking "what if exactly this one thing dies?" for every box in your diagram. Spread replicas across availability zones (independent power and network, ~1 ms apart) and, if the product needs it, across regions (hundreds of ms apart).

Availability math

99.9% ≈ 8.8 hours down a year, 99.99% ≈ 53 minutes, 99.999% ≈ 5 minutes. Components in series multiply (0.999 × 0.999 = 0.998), so every added hop costs availability. Redundant components in parallel improve it: two 99% instances give 1 − 0.01² = 99.99% — if their failures are genuinely independent.

SLI, SLO, SLA, error budget

An SLI is a measured indicator (p99 latency, success rate). An SLO is your internal target ("99.9% of reads under 200 ms over 30 days"). An SLA is the contractual promise with penalties — always looser than your SLO. The error budget is the allowed failure left over (0.1% of a month ≈ 43 minutes); spend it on releases, and freeze changes when it runs out.

RPO and RTO

RPO (recovery point objective) is how much data you can afford to lose — it sets backup and replication frequency. RTO (recovery time objective) is how long recovery may take — it sets whether you need a warm standby or can restore from a snapshot. A backup you have never restored is not a backup.

Observability

Metrics (cheap numeric time series — rates, errors, durations), logs (structured events for detail), traces (one request's path and timing across services, via a propagated trace id). Always quote percentiles, not averages: the mean hides the tail, and p99 is what your loudest users experience. Watch out for tail latency amplification — if one request fans out to 100 services, a 1% slow rate means almost every request is slow.

10. Traffic control and delivery

Rate limiting and throttling

Rate limiting rejects requests over a quota (429 Too Many Requests with Retry-After) to protect capacity and stop abuse; throttling slows or queues them instead. Algorithms — token bucket, leaky bucket, fixed window, sliding window — get a full lesson (lesson 9).

Realtime transports

Short polling: ask every few seconds — trivial, wasteful. Long polling: the server holds the request open until there is news — decent fallback. SSE: one long-lived HTTP stream, server → client only, auto-reconnecting — perfect for feeds and token streaming. WebSocket: full-duplex over one TCP connection — the choice for chat and multiplayer, and now you have stateful servers to manage. gRPC / HTTP/2 streams: efficient service-to-service.

Deployment safety

Blue-green runs two full environments and flips traffic (instant rollback, double the cost). Canary sends 1% of traffic to the new version and watches the metrics. Rolling replaces instances a few at a time. Feature flags separate deploying code from enabling behaviour, which makes both rollout and rollback a config change.

Consistent hashing, bloom filters, geohash — where they show up

Three structures pay for themselves repeatedly: consistent hashing for assigning keys to nodes (caches, shards, WebSocket servers), bloom filters for cheaply skipping work on absent keys (LSM reads, crawler URL dedupe, "have I sent this notification?"), and geohash / quadtree / S2 for turning two-dimensional coordinates into one-dimensional prefix-searchable keys (nearby search, lesson 16).

Numbers to reason with

# Latency, in nanoseconds, so ratios are obvious
L1_cache            = 1
main_memory         = 100
ssd_read            = 100_000          # 100 us
same_dc_round_trip  = 500_000          # 0.5 ms
disk_seek           = 10_000_000       # 10 ms
cross_continent_rtt = 150_000_000      # 150 ms

print("memory vs SSD:      ", ssd_read // main_memory, "x")
print("SSD vs disk seek:   ", disk_seek // ssd_read, "x")
print("local vs global hop:", cross_continent_rtt // same_dc_round_trip, "x")

# Why a cache hit rate matters: 4,000 reads/s at a 95% hit rate
reads = 4_000
for hit_rate in (0.0, 0.8, 0.95, 0.99):
    print(f"hit {hit_rate:.0%} -> DB sees {round(reads * (1 - hit_rate))} reads/s")
Output
memory vs SSD:       1000 x
SSD vs disk seek:    100 x
local vs global hop: 300 x
hit 0% -> DB sees 4000 reads/s
hit 80% -> DB sees 800 reads/s
hit 95% -> DB sees 200 reads/s
hit 99% -> DB sees 40 reads/s

That last block is the whole argument for caching in four lines: a 95% hit rate turns a database problem into a non-problem. It is also the argument for measuring the hit rate, because 80% and 99% differ by 20x in load.

The order you should reach for these

Escalate in this order, and only when a number forces the next step. Most interview candidates lose points by starting at step 7.

  1. One server, one database. Measure.
  2. Indexes and query fixes. (Most "we need to shard" is a missing index.)
  3. Vertical scale.
  4. Stateless app servers behind a load balancer; CDN for static assets.
  5. Cache the hot reads.
  6. Read replicas for read load; queues to move slow work off the request path.
  7. Shard, or split services, when one machine genuinely cannot hold the write volume or the data.
// Write your solution here

Finished reading? Mark this lesson complete to track your progress.

Up next · Lesson 3Scaling and CachingLoad balancers, horizontal scaling, caches, CDNs and consistent hashing.