On this page
Tracks

System Design — Building Blocks

Last reviewed 7 Sept 2026

The main guide lists these as one-liners so you can rattle them off. This page is the depth behind each one — enough to survive the “okay, go deeper on that” follow-up.

Numbers to memorise

Estimation is impossible without a few anchors. Know these cold.

ThingRough number
Seconds in a day~86,400 (≈ 10⁵)
L1 cache reference~1 ns
Main memory reference~100 ns
Read 1 MB sequentially from memory~10 µs
SSD random read~100 µs
Round trip within a data centre~0.5 ms
Read 1 MB sequentially from SSD~1 ms
Disk seek (spinning)~10 ms
Round trip US ↔ Europe~150 ms
Single commodity server: QPS ceiling~10K–100K simple req/s
Single Redis instance~100K ops/s, sub-ms
Single Postgres primary, well-tunedlow tens of thousands of writes/s
One row of typical structured data~1 KB order of magnitude

Derived moves you will make constantly: QPS ≈ requests/day ÷ 86,400, peak ≈ 2–10× average, storage/year = write QPS × row size × 3.15×10⁷ s.

Load balancing

Spreads traffic across identical instances and removes dead ones from rotation.

  • L4 (transport) — routes on IP/port, no idea what’s in the packet. Cheap, fast, protocol-agnostic. Can’t do path-based routing or TLS termination.
  • L7 (application) — parses HTTP; routes on path/header/cookie, terminates TLS, retries idempotent requests, does sticky sessions. More CPU per request.
  • Algorithms — round robin (default), least connections (uneven request cost), consistent hashing (want the same client on the same node for cache locality), weighted (mixed instance sizes).
  • Health checks — active (LB probes /healthz) and passive (eject after N consecutive 5xx/timeouts). A shallow check tests the process is up; a deep check tests its dependencies — deep checks can cascade a DB blip into a full outage, so keep them shallow and alert separately.

Follow-ups: How does the LB itself not become a SPOF? (Pair in active-passive with a floating IP, or DNS round-robin across several, or an anycast VIP.) Where do you terminate TLS? (At the L7 LB, usually.) Sticky sessions vs stateless? (Prefer stateless — put session state in Redis so any instance can serve any request.)

Caching

The first lever for read-heavy load. The hard parts are invalidation and stampede.

  • Cache-aside (lazy) — app checks cache, on miss reads the DB and populates the cache. Simple, resilient to cache loss, but every miss pays full latency and stale data lives until TTL.
  • Read-through / write-through — the cache library owns DB access. Write-through keeps cache and DB in lockstep at the cost of write latency.
  • Write-back — write to cache, flush to DB asynchronously. Fast writes, but a cache crash loses data — only for tolerant workloads (counters, metrics).
  • Eviction — LRU is the default; LFU when a small hot set dominates; TTL always, as a correctness backstop even when you also invalidate explicitly.

Cache stampede — a hot key expires and thousands of concurrent requests all miss and hit the DB at once. Fixes, usually combined:

  • Single-flight lock — first miss takes a lock and recomputes; the rest wait for the result.
  • Early / probabilistic recompute — refresh the value slightly before it expires so it’s never actually cold.
  • TTL jitter — ttl ± random, so keys populated together don’t expire together.
  • Serve-stale-while-revalidate — return the expired value immediately and refresh in the background.
flowchart TD
R[Request] --> C{In cache?}
C -- hit --> RET[Return value]
C -- miss --> L{Get recompute lock?}
L -- yes --> DB[(Database)] --> SET[Write to cache] --> RET
L -- no --> W[Wait for lock holder] --> RET
Cache-aside with single-flight on miss

Follow-ups: What’s your hit rate and how did you estimate it? What happens when Redis dies — does the DB fall over? (Cap concurrency to the DB, add a small in-process cache, degrade gracefully.) How do you invalidate on write? (Delete the key on write, or publish an invalidation event to all nodes; never try to update the cached value in place across a fleet.)

Replication

Copy the data to more than one node for read scaling and failover.

  • Primary–replica — all writes to the primary, which streams its log to replicas. Reads can go to replicas.
  • Replication lag — replicas are behind by milliseconds to seconds. “Read your own writes” breaks: user posts a comment, next page load hits a lagging replica, comment is gone. Fix by routing that user’s reads to the primary for a few seconds, or reading from the primary for session-critical paths.
  • Sync vs async — synchronous replication (wait for one replica to ack) protects against data loss on primary failure but adds write latency and stalls if the replica is slow. Async is the common default; semi-sync (ack from any one of N) is the middle ground.
  • Failover — promoting a replica needs leader election (or a human), fencing the old primary to prevent split-brain, and clients rediscovering the new primary. Automated failover can lose the last few async-replicated writes.

Follow-ups: What does replication not solve? (Write throughput and total dataset size — one primary still takes every write and holds the whole dataset. That’s what sharding is for.) Multi-primary? (Only with conflict resolution — last-write-wins, CRDTs, or app-level merge; avoid unless you truly need multi-region writes.)

Sharding / partitioning

Split writes and storage across independent primaries when one can’t hold the load or the data.

  • By key range — users A–M / N–Z. Range scans stay local; prone to hot ranges (everything recent lands in the last shard).
  • By hash of key — even spread, but range queries now fan out to every shard.
  • By entity / tenant — all of one customer’s data on one shard. Clean blast radius, natural for B2B; noisy-neighbour risk and big-tenant skew.
  • Directory / lookup table — a service maps key → shard, so you can rebalance by moving ranges and updating the map. Flexible; the directory is now a critical dependency.

Costs you must name: cross-shard queries and joins (fan-out + merge in the app), cross-shard transactions (need 2PC or saga), rebalancing when you add capacity, and hot keys — one celebrity / one viral item can swamp a single shard regardless of scheme. Mitigate hot keys by splitting that key into sub-keys or caching it hard.

flowchart TD
A[App] --> RT["Shard router: hash(key) mod N"]
RT --> S0[(Shard 0)]
RT --> S1[(Shard 1)]
RT --> S2[(Shard 2)]
A -. cross-shard query .-> SC["Scatter to all, gather, merge"]
Hash sharding with an app-side router

Consistent hashing

Standard hashing (hash(key) % N) remaps almost every key when N changes — catastrophic for a cache or a shard set that grows.

Consistent hashing places nodes and keys on a hash ring; a key belongs to the next node clockwise. Add or remove a node and only the keys between it and its predecessor move — about 1/N of keys, not all of them.

  • Virtual nodes — each physical node gets many points on the ring, so load evens out and a departing node’s keys spread across all survivors rather than dumping onto one neighbour.
  • Where it shows up — distributed caches (memcached clients), Cassandra/DynamoDB partitioning, shard assignment, L7 LBs doing session affinity.

Follow-ups: How many virtual nodes? (100–200 per physical node is typical.) How do you handle a node that’s a hot spot anyway? (More vnodes for weak nodes / fewer for hot data, or bound-load consistent hashing that caps any node at (1+ε)×average.)

Message queues and event logs

Decouple producers from consumers; absorb spikes; retry failures without blocking the caller.

  • Queue (SQS, RabbitMQ) — a message is delivered to one consumer and deleted on ack. Good for task distribution — send email, resize image, charge card. Has a visibility timeout and a dead-letter queue for messages that fail N times.
  • Log (Kafka, Kinesis) — an ordered, partitioned, retained stream. Many independent consumer groups each read at their own offset; you can replay from any point. Good for event sourcing, analytics pipelines, fan-out to multiple systems.
  • Ordering — only within a partition / message group. If you need per-user ordering, partition by user id.
  • Delivery — at-least-once is the norm, so consumers must be idempotent (next section). “Exactly-once” in Kafka is really atomic read-process-write within Kafka, not end-to-end.
  • Backpressure — queue depth is your early-warning metric. Alarm on it. Decide the policy when consumers can’t keep up: scale consumers, shed load, or let the queue buffer (and for how long).

Follow-ups: Queue vs log for this case? What’s your retry policy and DLQ handling? What’s the max acceptable end-to-end lag? How do you handle a poison message?

Idempotency and delivery semantics

Any consumer, retry, or externally-triggered write must tolerate running twice. There is no true exactly-once delivery over a network — you engineer effectively-once processing.

  • Idempotency key — client generates a unique id per logical operation; server records “this id → this result” and on a repeat returns the stored result instead of re-executing. Essential for payments, order creation, any POST that a client might retry.
  • Conditional writes — UPDATE ... WHERE status = 'pending', compare-and-set, INSERT ... ON CONFLICT DO NOTHING. The second attempt is a no-op.
  • Dedupe window — store processed message ids (with a TTL) and skip duplicates. Sized to cover your max redelivery delay.
  • Natural idempotency — “set balance to X” is idempotent; “add 10 to balance” is not. Prefer absolute operations where you can.

Outbox / change data capture

Problem: you need to update your DB and publish an event, and you can’t do both atomically — a crash between them loses the event or emits a phantom.

Transactional outbox — in the same DB transaction as the domain change, insert a row into an outbox table. A separate relay polls the table (or tails the DB log) and publishes each row to the broker, marking it sent. The event is now guaranteed if and only if the transaction committed. Consumers still dedupe, because the relay is at-least-once.

flowchart TD
H[Handler] --> TX[BEGIN]
TX --> D[Write domain row]
TX --> O[Write outbox row]
D --> CM[COMMIT]
O --> CM
CM --> RL[Relay polls / tails log]
RL --> B[(Message broker)]
RL --> MK[Mark outbox row sent]
Transactional outbox

CDC (Debezium etc.) is the same idea without an explicit table — the relay reads the database’s replication log directly and turns row changes into events.

CDN and blob storage

Static assets, images, video, and user uploads should never touch your application servers.

  • Blob store (S3, GCS) — cheap, durable (11 nines), effectively infinite. Clients upload and download directly using pre-signed URLs so bytes bypass your backend. Store only the object key in your DB.
  • CDN — caches those objects (and cacheable API responses) at edge PoPs near users. Cuts latency and origin load. Control it with Cache-Control, use content-hashed filenames for immutable assets, and purge by URL or surrogate key on change.
  • Big uploads — multipart upload direct to the blob store; for video, drop a message on a queue for a transcoding pipeline.

API gateway

The single front door in front of many services.

Handles cross-cutting concerns so services don’t each reimplement them: TLS termination, authn/authz, rate limiting, request routing, API-key management, request/response shaping, and aggregation (one client call → several service calls). Keep it thin — no business logic. It’s a critical path component, so it must be horizontally scaled and must fail fast rather than adding latency.

Consistency models (practical CAP)

Partitions happen, so during one you are choosing consistency or availability — and even without a partition (PACELC) you trade latency vs consistency.

ModelGuaranteeTypical use
Strong / linearizableEvery read sees the latest writeBalances, inventory, locks, uniqueness
Read-your-writesYou see your own updates; others may lagProfile edits, settings
Monotonic readsYou never see time go backwardsFeeds, comment threads
EventualReplicas converge “soon”Like counts, view counts, DNS, caches

Default posture for web systems: available + eventually consistent + idempotent writes, with strong consistency carved out only for the money/inventory critical path (single-row transactions, conditional updates, or a dedicated strongly-consistent store).

Observability

You can’t operate what you can’t see.

  • Metrics — cheap, aggregated, for dashboards and alerts. Track the RED set per service (Rate, Errors, Duration) and USE for resources (Utilisation, Saturation, Errors). Alert on p99 latency and error rate — the symptoms users feel — not on CPU alone.
  • Logs — structured (JSON), with a request id on every line. Sample the high-volume happy path; keep all errors.
  • Traces — one trace id propagated across every hop, so you can see where a slow request spent its time. Essential once you have more than a couple of services.
  • SLI / SLO — pick a few Service Level Indicators (e.g. “p99 < 300 ms”, “success rate > 99.9%”), set the objective, and spend the resulting error budget deliberately.