../

System design concepts

The building blocks behind system-design answers, explained well enough to reason with. For each block you get what it is, when to reach for it, what it costs, and one sentence to say in an interview. The production reference with numbers, caching tables, queues, rate limiting and a worked URL shortener is System design; data modeling is in Database design. How to run the interview itself, with nine worked designs, is in System design interviews.

Networking & protocols

BlockWhat it isReach for it whenCosts
DNSname → IP; stub resolver → recursive resolver → root → TLD → authoritative server; answers cached for their TTLgeo routing, weighted and failover records, internal service namesfailover is only as fast as the TTL (and some clients ignore it); a cold lookup costs a round trip
TCPreliable, ordered byte stream; 3-way handshake (1 RTT), congestion controlnearly everything: HTTP/1.1, HTTP/2, database protocolshandshake latency; one lost packet stalls every byte behind it (head-of-line blocking)
UDPbare datagrams: no handshake, order or retransmitvoice, video, games, DNS; QUIC is built on ityou rebuild reliability and ordering yourself; some networks throttle it
TLS 1.3encryption and server identity on top of TCP; 1-RTT handshake, 0-RTT on resumption (RFC 8446 (opens in a new tab))all public traffic; mTLS (both sides show certificates) between servicesa round trip, certificate rotation, a little CPU; 0-RTT data can be replayed, so only for idempotent requests
HTTP/1.1HTTP/2HTTP/3
TransportTCPTCPQUIC over UDP
Concurrencyone request at a time per connection; browsers open about 6 per originmany streams on one connectionmany independent streams
Head-of-line blockingper connectionat the TCP layer: one lost packet stalls all streamsper stream only
Headersplain text, repeatedHPACK compressionQPACK compression
New connectionTCP + TLS 1.3 = 2 RTTsame1 RTT (0-RTT on resume)
Extraserver push, since removed from Chrome and Firefoxconnection migration (Wi-Fi to cellular keeps the connection)
SpecRFC 9112 (opens in a new tab)RFC 9113 (opens in a new tab)RFC 9114 (opens in a new tab)

Server-to-client updates

OptionHow it worksReach for it whenCosts
Short pollingclient asks every N secondsupdates are rare, simplicity winswasted requests; up to N seconds late
Long pollingserver holds the request until there is data or a timeout; the client reconnectsstreaming is blocked by proxiesa request per burst; reconnect storms after deploys
Server-Sent Events (SSE)one HTTP response left open, text/event-stream; EventSource reconnects and resends Last-Event-ID (MDN (opens in a new tab))one-way feeds: notifications, live scores, LLM token streamsserver → client only, text only; on HTTP/1.1 each stream holds one of the ~6 connections per domain, shared across all tabs (HTTP/2 multiplexes them)
WebSocketsHTTP Upgrade to a full-duplex framed connection (RFC 6455 (opens in a new tab))both sides send often: chat, multiplayer, collaborative editingstateful: connection servers, a pub/sub backplane to reach the right one, reconnect and resync logic, draining on deploy
Webhooksyour server POSTs to a URL the customer registeredtelling another company's server that something happened (payments, CI)receiver may be down: retries with backoff, HMAC signatures, receivers must dedupe

Say it like this

  • DNS: "DNS failover is bounded by the TTL, so for fast regional failover I'd put anycast or a global load balancer in front."
  • TCP vs UDP: "It's live audio, so late packets are useless; UDP (or WebRTC) beats TCP retransmits here."
  • HTTP/3: "Mobile clients on lossy networks benefit from HTTP/3, because a lost packet only stalls its own stream."
  • SSE vs WebSockets: "Updates only flow server to client, so SSE over plain HTTP is enough; WebSockets only if the client streams too."
  • Webhooks: "We deliver webhooks at least once with retries and a signature, and document that receivers must dedupe by event ID."

API design

RESTgRPCGraphQL
Shaperesources + HTTP verbs, usually JSONRPC methods, Protobuf over HTTP/2one endpoint; the client picks the fields
Strengthscacheable by browsers and CDNs, universal tooling, easy to debugsmall and fast, typed code generation, streaming both ways, deadlinesno over- or under-fetching, nested data in one round trip, typed schema
Weaknessesover-fetching; many round trips for nested databrowsers need gRPC-Web or a proxy; binary is harder to inspectHTTP caching is hard, N+1 resolvers (batch with DataLoader), must cap query depth and cost
Reach for it whenpublic APIs, CRUDservice-to-service calls inside your networkmany clients needing different shapes (web, mobile), one API over many services
PaginationOffset (?page=3&limit=20)Cursor / keyset (?after=token&limit=20)
QueryORDER BY id LIMIT 20 OFFSET 40WHERE (created_at, id) < (:t, :id) ORDER BY created_at DESC, id DESC LIMIT 20
Deep pagesslow: the database still reads and skips the offset rowsan index seek per page, constant cost
New rows arrivepages shift: duplicates and skipped itemsstable
Jump to page Nyesno, next and previous only
Use foradmin tables, small result setsfeeds, infinite scroll, public APIs at scale
ConcernPractice
Idempotency keysclient sends a unique Idempotency-Key per logical operation; server stores key → status and response under a unique constraint, replays the stored response on retry, rejects the same key with a different body; keys expire (Stripe keeps them 24 h). The header name comes from an expired IETF draft (opens in a new tab) and is a convention, not a standard
Safe retriesGET, PUT, DELETE are idempotent by definition (RFC 9110 (opens in a new tab)); POST and PATCH need a key
Versioningpath (/v1/), header or media type, or date-pinned versions per account (Stripe); make additive changes without a new version; never reuse Protobuf field numbers; deprecate GraphQL fields
Errorsconsistent envelope, machine-readable code, application/problem+json (RFC 9457 (opens in a new tab))
Rate limitsper API key, user and IP; 429 Too Many Requests with Retry-After; RateLimit-Policy / RateLimit headers are an IETF draft (opens in a new tab); algorithms in Reliability & availability
Long operations202 Accepted + a job resource to poll, or a webhook when done
Large uploadspresigned URL straight to object storage, multipart for big files

Say it like this

  • REST vs gRPC: "The public API is REST for cacheability and reach; internal calls are gRPC for typed contracts and deadlines."
  • GraphQL: "Web and mobile need different shapes of the same data, so a GraphQL layer over the services saves round trips; I'd cap query cost."
  • Cursor pagination: "Cursor pagination on (created_at, id) keeps deep pages an index seek and doesn't skip items when new posts arrive."
  • Idempotency: "Clients retry payments on timeouts, so POST /payments takes an idempotency key and replays the stored result."

Load balancing & proxies

BlockWhat it isReach for it whenCosts
L4 load balancerroutes TCP/UDP connections by IP and portraw throughput, non-HTTP protocols, TLS passthroughcan't see paths or headers; a long-lived connection stays pinned to one backend
L7 load balancerparses HTTP; routes by host, path, header; terminates TLSmicroservices, canaries, per-request gRPC balancing, retriesmore CPU; it sees plaintext; one more hop
Reverse proxyserver in front of your app: TLS, compression, caching, buffering slow clients (nginx, Envoy, Caddy)always, even with one app serveranother component to configure and monitor
API gatewayL7 proxy with API concerns: auth, rate limits, quotas, request shaping, aggregationone public API over many services; a backend-for-frontend per clientbusiness logic creeps in; must be replicated like any hot path
Service discoveryregistry of healthy instances: DNS (Kubernetes Services), Consul, etcd; client-side or server-side lookupautoscaled, short-lived instancesthe registry must be highly available; stale entries need health checks and TTLs
Service mesha sidecar proxy per instance for mTLS, retries, telemetry (Istio, Linkerd)many services, uniform security and traffic policyextra latency per hop, operational weight
AlgorithmPicksGood for
Round robin (weighted)next backend in turnuniform, short requests
Least connectionsfewest open connectionslong or uneven requests, WebSockets
Power of two choicesbetter of 2 random backendsbig fleets; avoids every LB herding onto the same "least loaded" node
Consistent hashingkey's position on a ring (see Scaling data)caches and shards: same key, same node, few keys move on resize

Full table with L4/L7 products and health-check details: System design: Scaling & load balancing.

Stateless services keep no user state in process memory: sessions go to Redis or a signed token, files to object storage, caches to a shared cache. Then any instance can serve any request, instances can die or be added at will, and autoscaling works. State lives in a few purpose-built stores that you scale deliberately.

Say it like this

  • L7: "An L7 balancer lets me route /api and /ws to different pools and do canary releases by header."
  • Stateless: "The API tier is stateless, sessions live in Redis, so I can scale it horizontally behind the balancer."
  • Gateway: "Auth and per-key rate limits sit in the gateway so every service doesn't reimplement them."
  • Discovery: "Instances register with the service registry and the balancer only routes to ones passing readiness checks."

Caching

A cache trades freshness for speed. Effective latency is h⋅tcache+(1−h)⋅tdbh \cdot t_{cache} + (1 - h) \cdot t_{db}: at a 95% hit rate, 1 ms cache and 10 ms database, reads average 1.45 ms. Every percent of hit rate you gain removes load from the database, which is usually the real goal.

WhereHoldsReach for it whenCosts
Client (browser, app)assets, API responses (Cache-Control, ETag)static assets, per-user data that changes rarelyyou can't purge it: use content-hashed file names
CDN / edgestatic files, images, video segments, cacheable HTML and API responsesglobal users, large or popular filespurge delays, cost per GB, personalized data needs care
Application (in-process LRU)hot config, small lookupstiny, very hot dataeach instance holds its own copy: inconsistent, lost on restart
Distributed cache (Redis, Valkey, Memcached)objects, sessions, computed views, countersshared across instances, read-heavy dataa network hop, memory cost, one more thing to fail
Database (buffer pool, materialized views)pages, precomputed query resultsexpensive aggregatesrefresh cost; staleness of the view
Read check the cache first app cache database 1 GET key 2 miss 3 SELECT 4 row 5 SET key, TTL on a hit, step 2 returns the value: done only requested keys get cached (lazy) Write database first, then invalidate app cache database 1 UPDATE 2 commit ok 3 DEL key the next read misses and reloads delete, not update: no races The app owns both paths; the cache never calls the DB (that is read-through). TTL is the backstop if a delete is lost.
Cache-aside, the default pattern: fill on read, invalidate on write.
PatternReadWriteReach for it whenCosts
Cache-aside (lazy)app checks cache, on miss loads DB and fillsapp writes DB, deletes the keythe default; the cache can fail without breaking correctnessfirst read is slow; a small race can re-cache stale data
Read-throughcache library loads from DB on a missas cache-asideyou want the loading logic in one placeneeds a cache that supports loaders
Write-throughfrom cachewrite cache and DB together, synchronouslydata read right after it's writtenslower writes; caches data nobody reads
Write-back (write-behind)from cachewrite cache, flush to DB later in batcheswrite-heavy counters, metricsdata loss if the cache dies before flushing
Write-aroundfrom cache (cache-aside on miss)write DB only, skip the cachedata written once and rarely read back (logs, uploads)first read of new data always misses
EvictionEvictsGood for
LRUleast recently usedgeneral default, recency-heavy traffic
LFUleast frequently usedstable popularity (top products); resists one-off scans
TTLanything past its expirybounding staleness; combine with LRU/LFU
FIFO / randomoldest / anycheap, surprisingly fine for large caches

Redis picks victims only when maxmemory is reached and the default policy is noeviction (writes fail), so a cache needs allkeys-lru or allkeys-lfu set explicitly (Redis eviction docs (opens in a new tab)).

ProblemWhat happensFix
Invalidationdata changes but the cache still serves the old valuedelete after commit, short TTL as a backstop, versioned keys, or CDC-driven deletes
Stampede (dogpile)a hot key expires and thousands of requests hit the DB at oncesingle-flight per key, a lease so only one caller reloads (Facebook memcache (opens in a new tab)), serve stale while one refreshes, early refresh with jitter
Hot keyone key gets so many reads that one cache node saturatesin-process cache in front, replicate the key as key#1..key#N and read a random copy
Penetrationrequests for keys that don't exist always misscache "not found" briefly; Bloom filter of existing keys

More fixes and a single-flight snippet: System design: Caching & CDNs.

Say it like this

  • Cache-aside: "Reads go cache-aside with a TTL; on write we update Postgres, then delete the key, so the next read reloads it."
  • Why a cache: "The read:write ratio is 100:1 and the data tolerates a few seconds of staleness, so a cache takes most load off the database."
  • Stampede: "To stop a stampede on the celebrity profile, one request refreshes the key while the others get the stale value."
  • CDN: "Video segments and thumbnails are immutable, so the CDN caches them forever under content-hashed names."

Data stores

TypeExamplesPick whenAvoid when
RelationalPostgres, MySQL; distributed: CockroachDB, Spanner, Vitess, Citusthe default: relations, constraints, transactions, ad-hoc querieswrite volume or data size outgrows one primary and sharding by hand hurts
Key-valueRedis, DynamoDB, etcdlookups by a known key: sessions, carts, counters, feature flagsyou need queries on anything but the key
DocumentMongoDB, Firestore, Couchbaseself-contained aggregates read and written together; flexible schemamany-to-many relations and cross-document transactions
Wide-columnCassandra, ScyllaDB, Bigtable, HBasehuge write rates, queries known upfront (partition key + clustering order): messages, eventsad-hoc queries, joins, strong consistency by default
GraphNeo4j, Neptunemany-hop relationship queries: friends of friends, fraud ringssimple lookups; data that isn't a graph
Time-seriesTimescaleDB, InfluxDB, Prometheusappend-only measurements with time-range queries, rollups, retentiongeneral application data
Search indexElasticsearch, OpenSearch, Meilisearchfull-text relevance, fuzzy matching, facetsas the source of truth (it's a projection)
Object / blobS3, GCS, R2files, images, video, backups, data lakesmall records, frequent in-place updates
Vectorpgvector, Qdrant, Pinecone, Milvusnearest-neighbor search over embeddings (semantic search, RAG) using approximate indexes (HNSW, IVF)exact matching; recall must be 100%

SQL vs NoSQL is really "general-purpose queries vs design-for-known-access-patterns". Relational databases let you ask new questions later and keep invariants with transactions; NoSQL stores make you model tables around the queries you already know, and in return partition and replicate automatically. Distributed SQL (Spanner, CockroachDB) narrows the gap at the cost of write latency. Default to Postgres and name the specific access pattern that pushes you elsewhere.

Indexes: B-tree vs LSM-tree

B-treeLSM-tree
Writesupdate the page in place (plus the WAL)append to the WAL and an in-memory table; flush sorted SSTables; merge them in the background (compaction)
ReadsO(log n) page reads, predictablecheck the memtable, then SSTables newest to oldest; Bloom filters skip files that can't hold the key
Wins whenread-heavy, range scans, predictable latencywrite-heavy ingest, time-series, data much larger than RAM
Costsrandom writes, page splitscompaction spikes, read amplification, space amplification
Used byPostgres, MySQL InnoDB, SQLiteRocksDB, LevelDB, Cassandra, ScyllaDB, HBase
Secondary index on a sharded storeHowCost
Local (per shard)each shard indexes only its own rowsreads must ask every shard (scatter-gather)
Global (partitioned by the indexed value)the index is its own sharded tablewrites touch several shards; usually updated asynchronously

Schema, normalization and composite-index design: Database design.

Say it like this

  • Default: "I'll start with Postgres: the data is relational and we need transactions; one primary with replicas covers this scale."
  • Wide-column: "Messages are append-heavy and always read by conversation, so Cassandra with conversation_id as the partition key and time as the clustering key fits."
  • LSM: "This is write-heavy telemetry, so an LSM-based store absorbs writes as sequential appends."
  • Search: "Postgres stays the source of truth; Elasticsearch is a projection fed by CDC, so it can be rebuilt."

Scaling data

StepWhatReach for it whenCosts
Vertical scalingbigger machinefirst; modern boxes go far (hundreds of cores, TBs of RAM)a ceiling, a single point of failure, price at the top end
Read replicasfollowers serve readsread-heavy loadreplication lag: stale reads, read-your-writes needs care
Cachingsee Cachinghot, repeated readsinvalidation
Sharding (partitioning)split rows across nodes by a shard keywrites or data outgrow one primarycross-shard queries and transactions, rebalancing, operational load
Denormalizationstore data pre-joined or duplicated for a read patha hot read needs several joinswrites update several copies; copies can drift
Leader-follower one writer Multi-leader a writer per region Leaderless quorum N=3, W=2, R=2 client leader writes follower follower reads: any replica async lag = stale reads log client EU client US leader EU leader US follower follower replicate both ways writes can conflict: LWW, merge or CRDT client r1 ack r2 ack r3 late write 3, done at W=2 read R=2, newest wins R + W > N → sets overlap solid: client writes dashed: replication (usually async) default: leader-follower · multi-region writes: multi-leader highest write availability: leaderless
Leader-follower, multi-leader and leaderless (quorum) replication.
ReplicationHowReach for it whenCosts
Synchronousleader waits for the follower before acknowledgingno data loss on failoverwrite latency; a slow follower stalls writes (use one sync follower, the rest async)
Asynchronousleader acknowledges, followers catch upthe common defaultfailover can lose the last writes; stale reads
Leader-followerall writes through one leaderalmost alwaysthe leader caps write throughput; failover needs care
Multi-leadera leader per region (or per device), syncing both waysmulti-region writes, offline-first appswrite conflicts: last-writer-wins loses data; merge functions or CRDTs keep it
Leaderless (Dynamo-style)client writes to N replicas, waits for W, reads Rhigh write availability, tolerating node lossquorum tuning, read repair, conflicts via version vectors or LWW

Read-your-writes on async replicas: after a user writes, read their data from the leader for a few seconds, or send the write's log position and have the replica wait until it has caught up.

Sharding strategies

StrategyReach for it whenCosts
Range (by key or time)range scans mattersequential keys (timestamps) make one hot shard
Hash (hash(key) mod N)even spread, point lookupsno range scans; changing N moves almost every key
Consistent hashing + virtual nodesnodes join and leave often (caches, Dynamo-style stores)ring management
Fixed partitions (many more than nodes)planned growth: create 1,000 partitions, move whole partitions (Kafka, Elasticsearch, Riak)partition count chosen upfront
Directory (lookup table)tenant placement, moving big customersthe lookup service is on every request path
Geodata residency, latency per regionuneven region sizes

Detail and shard-key advice: System design: Sharding.

Consistent hashing

A#0 B#0 C#0 A#1 C#1 B#1 A#2 C#2 B#2 hash("user:42") hash ring 0 … 2³² − 1, wraps Place each server at many points: hash("A#0"), hash("A#1") … (virtual nodes) Look up hash the key, walk clockwise to the first virtual node: user:42 → A Add or remove a server only arcs next to its points change owner: ~1/N of keys move, not almost all of them (hash mod N would move ~all) server A server B server C colored arcs: the key range each virtual node owns (from the previous point up to it)
Keys go clockwise to the next virtual node; each server owns many small arcs.

Hash servers and keys onto the same ring (0 … 2³² − 1). A key belongs to the first server point clockwise from it. Adding a server takes over only the arcs just before its points, about 1/N of keys; with hash mod N almost every key would move. Virtual nodes (each server at 100–200 points) even out the arcs, spread a dead server's load across all survivors, and let a bigger machine take more points.

// FNV-1a over UTF-8 bytes, then a murmur3 finalizer
// so near-identical strings land far apart
function hash32(s: string, seed = 0x811c9dc5): number {
  let h = seed;
  for (const b of new TextEncoder().encode(s)) {
    h = Math.imul(h ^ b, 0x01000193);
  }
  h = Math.imul(h ^ (h >>> 16), 0x85ebca6b);
  h = Math.imul(h ^ (h >>> 13), 0xc2b2ae35);
  return (h ^ (h >>> 16)) >>> 0;
}
 
type VNode = { pos: number; node: string };
 
class HashRing {
  private ring: VNode[] = [];
 
  constructor(private readonly vnodes = 100) {}
 
  add(node: string): void {
    for (let i = 0; i < this.vnodes; i++) {
      const pos = hash32(`${node}#${i}`);
      this.ring.push({ pos, node });
    }
    this.ring.sort((a, b) => a.pos - b.pos);
  }
 
  remove(node: string): void {
    this.ring = this.ring.filter((v) => v.node !== node);
  }
 
  nodeFor(key: string): string {
    if (this.ring.length === 0) throw new Error("empty");
    const h = hash32(key);
    let lo = 0;
    let hi = this.ring.length; // first pos ≥ h
    while (lo < hi) {
      const mid = (lo + hi) >>> 1;
      if (this.ring[mid].pos < h) lo = mid + 1;
      else hi = mid;
    }
    return this.ring[lo % this.ring.length].node; // wrap
  }
}

add: O(Vlog⁡V)O(V \log V) for VV virtual nodes in total (re-sort); nodeFor: O(log⁡V)O(\log V); space O(V)O(V). With 4 servers × 100 virtual nodes and 20,000 keys, adding a fifth server moved 25% of keys (ideal: 20%), all of them to the new server.

Hot partitions and resharding

ProblemFix
Celebrity key (one user, one product, one hashtag)split the key with a suffix (post:123#0..#9) and merge on read; cache it; give it a dedicated shard
Time-ordered key (all writes hit "now")hash-prefix the key, or shard by (tenant, time bucket)
Uneven tenantsdirectory-based placement; move big tenants to their own shard
Resharding a live systemdouble-write old and new layout, backfill, compare, switch reads, then stop old writes; or fixed partitions from day one so you only move whole partitions

Say it like this

  • Replicas: "Reads outnumber writes 50 to 1, so I'll add read replicas and route read-after-write for the author to the leader."
  • Shard key: "I'll shard by user_id: it's in every query, high-cardinality, and spreads writes evenly."
  • Consistent hashing: "The cache tier uses consistent hashing with virtual nodes, so adding a node only moves about 1/N of the keys."
  • Hot partition: "A celebrity's counter would melt one shard, so I split it into 10 sub-keys and sum on read."

Consistency & transactions

CAP, stated correctly: when a network partition separates replicas, a system must either refuse some requests (keep Consistency, meaning linearizability) or answer them with possibly stale data (keep Availability). Partitions aren't optional, so "pick two of three" is misleading: the real choice is what happens during a partition. PACELC adds the everyday case: if Partitioned, choose A or C; else choose Latency or Consistency, because waiting for replicas costs time even when nothing is broken.

ModelGuaranteeReach for it whenCosts
Linearizable (strong)every read sees the latest completed write; behaves like one copybalances, inventory, unique usernames, locks, leader electioncoordination on every operation: latency, unavailable in a partition
Causaleffects never appear before their causes, for everyonecomments and replies, chattracking dependencies
Read-your-writesyou always see your own writesprofile edits, posting then viewingrouting or waiting per user
Monotonic readsyou never see time go backwardany replica-served UIpin a user to one replica
Eventualreplicas converge once writes stoplikes, view counts, feeds, DNStemporary anomalies users may notice

Definitions in depth: Jepsen: consistency models (opens in a new tab).

Quorums

With NN replicas, a write waits for WW acknowledgments and a read asks RR replicas. If R+W>NR + W > N, every read set overlaps every write set, so at least one replica in the read has the latest acknowledged write (pick it by version). N=3,W=2,R=2N = 3, W = 2, R = 2 tolerates one slow or dead node for both reads and writes. W=N,R=1W = N, R = 1 gives fast reads and fragile writes; W=1,R=NW = 1, R = N the reverse. Quorums alone are not linearizable: sloppy quorums (writes land on stand-in nodes during failures), concurrent writes and partially failed writes break the guarantee (DDIA ch. 5).

Transactions and isolation

ACID: Atomic (all or nothing), Consistent (invariants hold; mostly the app's job), Isolated (concurrent transactions don't see each other's partial work), Durable (committed data survives a crash).

AnomalyWhat goes wrong
Dirty readyou read another transaction's uncommitted write
Non-repeatable readthe same row read twice returns different values
Phantomthe same query returns new rows the second time
Lost updatetwo read-modify-write cycles race; one overwrites the other
Write skewtwo transactions read the same data, write different rows, and together break a rule ("at least one doctor on call")
LevelPreventsStill allowsNotes
Read committeddirty reads and writesnon-repeatable reads, phantoms, lost updates, write skewPostgres default
Repeatable read / snapshot isolationthe above + non-repeatable readswrite skew (phantoms in the SQL standard)MySQL InnoDB default; Postgres aborts lost updates here, InnoDB doesn't
Serializableall of the abovenothing; conflicting transactions abort and retryPostgres uses SSI (optimistic); others use two-phase locking

Fixes without going serializable: atomic updates (SET n = n + 1), SELECT … FOR UPDATE, compare-and-set on a version column, unique constraints. Postgres behavior: transaction isolation docs (opens in a new tab).

Across services

Two-phase commit (2PC)Saga
Howcoordinator asks all participants to prepare, then commita chain of local transactions; each failure runs compensating actions (refund, release)
Guaranteeatomic across databaseseventually consistent; intermediate states are visible
Reach for it whenone organization, few participants, databases that support XAmicroservices, long-running business flows (order → payment → shipping)
Costsblocking: a coordinator crash leaves participants holding locks; latencycompensations to design; no isolation between steps

Publishing events reliably alongside a DB write: the transactional outbox.

Say it like this

  • CAP: "For the inventory counter I'd choose consistency during a partition: better to reject an order than oversell."
  • PACELC: "Even without partitions, synchronous cross-region replication adds about 70 ms per write, so the feed stays eventually consistent."
  • Quorum: "With N = 3, W = 2, R = 2 any read overlaps the last write and we survive one node down."
  • Isolation: "Two bookings could race for the last seat, so I'd use SELECT … FOR UPDATE on the seat row, or a unique constraint on (flight, seat)."
  • Saga: "Checkout is a saga: reserve stock, charge, ship; if charging fails, we release the reservation."

Async processing

Queue (SQS, RabbitMQ)Log / stream (Kafka, Kinesis, Redpanda)
Modelmessages deleted once acknowledgedappend-only partitioned log; consumers track an offset
Consumerscompeting workers share one queueconsumer groups: each partition goes to one consumer in the group; many groups read independently
Replaynoyes, within retention
OrderingFIFO queues per group ID; standard queues best-effortwithin a partition only
Reach for it whenbackground jobs, task distribution, per-message retriesevent streams, CDC, analytics, several independent consumers of the same events
DeliveryMeaningHow
At most oncemay lose, never duplicatesacknowledge before processing
At least oncenever lose, may duplicateacknowledge after processing (the usual default)
Exactly onceeach message affects state onceonly inside one system (Kafka transactions for read-process-write); anywhere else it's effectively once: at-least-once + an idempotent consumer
ConcernPractice
Idempotent consumerrecord processed message IDs in the same transaction as the side effect; or make writes naturally idempotent (upsert, set instead of increment)
Orderingchoose the partition key so related events share a partition (order_id); parallelism is capped by the partition count
Poison messagesafter N failed attempts move to a dead-letter queue, alert, fix, replay
Visibility timeout / ack deadlinelonger than processing time, or the message is redelivered mid-work
Backpressurebounded queues; watch consumer lag; slow producers or shed load instead of buffering forever
Retriesbackoff with delay queues or retry topics, not tight loops

Kafka vs RabbitMQ vs SQS side by side: System design: Queues & streams.

Fan-out on write vs on read

Fan-out on write (push) precompute every follower's feed amy posts fan-out workers queue feed:bob feed:cy feed:dan 1 write per follower bob reads 1 key lookup: fast reads cost: a post by someone with 10M followers means 10M writes Fan-out on read (pull) build the feed when it is requested posts:amy posts:eve posts:max amy posts 1 write, nothing else merge + sort query each followee bob reads cost: slow reads if bob follows 1,000 Hybrid (the usual answer): push for ordinary authors, pull for celebrities, and merge the two at read time. Bars under each feed = cached post IDs.
Push costs writes per follower; pull costs work per read. Real feeds mix both.
Fan-out on write (push)Fan-out on read (pull)
On postappend the post ID to every follower's feed liststore the post once
On readread one precomputed listfetch recent posts of every followee, merge by time
Reach for it whenmost users, most readscelebrities (millions of followers), inactive readers
Costswrite amplification; wasted work for inactive followersslow, expensive reads for users who follow many accounts

Twitter's timeline service used the hybrid: fan-out on write into cached timelines, with very large accounts merged in at read time (Timelines at Scale (opens in a new tab)).

Say it like this

  • Queue: "Sending email doesn't need to block the request, so the API enqueues a job and returns 202."
  • Kafka: "Order events go to Kafka keyed by order_id, so each order's events stay in order and analytics can replay them."
  • Exactly once: "Delivery is at least once, so the consumer dedupes by event ID in the same transaction as its write."
  • Fan-out: "Push to followers' feeds on write, except accounts over about 10k followers, whose posts get merged in at read time."

Coordination

BlockWhat it isReach for it whenCosts
Leader electionexactly one node acts as leader (scheduler, partition owner, primary)one writer or one cron runner is neededelection time is downtime; two leaders at once is the failure to prevent
Consensus (Raft, Paxos, ZAB)a majority agrees on an ordered log of decisionsleader election, config, metadata, locks2f+1 nodes to survive f failures; every write is a majority round trip
ZooKeeper / etcd / Consulsmall, strongly consistent stores built on consensus, with watches and leasescluster metadata, service discovery, locks (Kubernetes stores its state in etcd)not for application data or high write rates
Distributed lock with leasea lock that expires unless renewedmaking sure only one worker does a joba paused holder can wake up after expiry and still act: needs fencing
Fencing tokena number that increases with each lock grant; storage rejects writes with an older tokenany lock protecting a writethe storage must check the token

Raft in one paragraph. Each node is a follower, candidate or leader. A follower that hears no heartbeat within a randomized election timeout (150–300 ms in the paper) becomes a candidate, increments the term and asks for votes; a node grants one vote per term, and only to candidates whose log is at least as up to date as its own. A majority makes a leader. The leader appends client commands to its log and replicates them; an entry is committed once stored on a majority, then applied to each node's state machine. Because any two majorities overlap, a committed entry survives every future leader change (Raft paper (opens in a new tab)). Use 3 nodes to survive 1 failure, 5 to survive 2.

Fencing, the classic bug: worker 1 takes a lease, pauses for GC longer than the lease, worker 2 takes the lock and writes, worker 1 wakes up and writes too. With fencing tokens worker 1 holds token 33, worker 2 holds 34, and storage refuses 33 after seeing 34 (Kleppmann: How to do distributed locking (opens in a new tab)).

Clocks

ClockWhat it givesLimit
Wall clock (NTP-synced)real timeskew of milliseconds or more between machines; can jump backward; never order events across machines by it
Monotonic clockdurations on one machinemeaningless across machines
Lamport clockcounter: max(local, received) + 1; if a caused b, then L(a) < L(b)can't tell whether two events were concurrent
Vector clocka counter per node; compares as before, after or concurrentsize grows with the number of nodes
Hybrid logical clock / TrueTimephysical time plus a logical counter (CockroachDB) or bounded uncertainty with commit-wait (Spanner)HLC needs bounded skew; TrueTime needs GPS and atomic clocks

Unique IDs: use time-ordered IDs for index locality: Snowflake (64-bit: timestamp, worker, sequence), ULID and UUIDv7 (128-bit, millisecond-sortable). Layouts and trade-offs: System design: IDs.

Say it like this

  • Consensus: "Leader election and config live in etcd, which uses Raft, so with 3 nodes we survive one failure without split-brain."
  • Locks: "The lock has a 30-second lease and a fencing token, so a stale holder can't overwrite newer work."
  • Clocks: "I won't order messages by server timestamps because clocks skew; the partition's sequence number defines order."
  • IDs: "Message IDs are Snowflake-style: sortable by time, generated locally without coordination."

Specialized structures

StructureAnswersAccuracy and spaceWhere it shows up
Bloom filter"definitely not in the set" or "maybe"no false negatives; about 9.6 bits per item for 1% false positives; no deletes (a counting variant allows them)LSM reads skip SSTables, "username taken?", crawler's seen URLs, cache penetration
Count-min sketchapproximate count per itemonly overestimates; error ≤ εN with probability 1 − δ using ⌈e/ε⌉ × ⌈ln(1/δ)⌉ countersheavy hitters, trending hashtags, top-k with a heap
HyperLogLogapproximate number of distinct itemsstandard error ≈ 1.04/m1.04 / \sqrt{m} for mm registers; Redis uses m=16384m = 16384: at most 12 KB per key, 0.81% standard error, mergeable (Redis docs (opens in a new tab))unique visitors per page per day
Geohasha base32 string per cell; shared prefix ≈ nearby5 chars ≈ 4.9 × 4.9 km, 6 chars ≈ 1.2 × 0.6 km (at the equator)proximity search in a plain B-tree or Redis GEOSEARCH
Quadtreea tree that splits a cell into 4 when it holds more than k pointsadapts to density (Manhattan splits deep, deserts don't)in-memory proximity index
S2 / H3hierarchical cells on the sphere: S2 squares on a Hilbert curve (Google), H3 hexagons (Uber)uniform cells without pole distortion; H3 neighbors are all equidistantmaps, ride dispatch, surge pricing zones
Trie (prefix tree)all keys with a given prefixmemory-heavy; store the top-k completions at each nodetypeahead
Inverted indexterm → sorted list of document IDsintersection of posting lists answers AND queriesfull-text search (Lucene, Elasticsearch, Postgres GIN)
Skip listsorted set with O(log n) expected search and insert via random "express lanes"simpler than balanced trees, lock-friendlyRedis sorted sets, LSM memtables
SSTableimmutable file of sorted key-value pairs plus a sparse indexmerge-friendlyLSM-tree storage (see Data stores)
Merkle treea tree of hashes; equal roots mean equal datafind differing ranges by comparing O(log n) hashes, not all dataanti-entropy repair between replicas (Dynamo, Cassandra), Git, certificate transparency

Geohash cells cut through neighborhoods: two points 10 m apart can sit on either side of a boundary with no common prefix, so always search the cell and its 8 neighbors (Geohash precision table (opens in a new tab)).

Bloom filter

With mm bits, kk hash functions and nn items inserted, the false-positive rate is p≈(1−e−kn/m)kp \approx (1 - e^{-kn/m})^k. Sizing for nn items at a target pp: m=−nln⁡p(ln⁡2)2m = -\frac{n \ln p}{(\ln 2)^2} bits and k=mnln⁡2k = \frac{m}{n} \ln 2 hash functions. The kk positions come from two hashes as h1+i⋅h2h_1 + i \cdot h_2 (double hashing). Uses hash32 from the ring above.

class BloomFilter {
  private readonly bits: Uint8Array;
 
  constructor(
    private readonly m: number, // bits
    private readonly k: number, // hash functions
  ) {
    this.bits = new Uint8Array(Math.ceil(m / 8));
  }
 
  static forCapacity(n: number, p: number): BloomFilter {
    const m = Math.ceil((-n * Math.log(p)) / Math.LN2 ** 2);
    const k = Math.max(1, Math.round((m / n) * Math.LN2));
    return new BloomFilter(m, k);
  }
 
  private indexes(item: string): number[] {
    const h1 = hash32(item);
    const h2 = hash32(item, 0x9747b28c);
    return Array.from(
      { length: this.k },
      (_, i) => (h1 + i * h2) % this.m, // double hashing
    );
  }
 
  add(item: string): void {
    for (const i of this.indexes(item)) {
      this.bits[i >> 3] |= 1 << (i & 7);
    }
  }
 
  mightContain(item: string): boolean {
    return this.indexes(item).every(
      (i) => (this.bits[i >> 3] & (1 << (i & 7))) !== 0,
    );
  }
}

add and mightContain: O(k)O(k) time; space mm bits. For 1,000 items at 1%: 9,586 bits (1.2 KB) and 7 hashes; measured false-positive rate on 100,000 absent keys: 1.04%.

Say it like this

  • Bloom filter: "Before hitting the database for 'is this username taken', a Bloom filter answers most 'no' cases from memory."
  • HyperLogLog: "Unique viewers per video per day go into a HyperLogLog: 12 KB per counter, under 1% error, and they merge across days."
  • Count-min: "Trending hashtags: a count-min sketch per minute window plus a min-heap of the top 100."
  • Geo: "Businesses are indexed by geohash; a search covers the user's cell and the 8 neighbors, then filters by exact distance."
  • Merkle: "Replicas compare Merkle trees in the background and only exchange the ranges whose hashes differ."

Reliability & availability

AvailabilityDowntime / yearDowntime / 30 daysDowntime / day
99%3.65 days7.2 h14.4 min
99.5%1.83 days3.6 h7.2 min
99.9%8.76 h43.2 min1.44 min
99.95%4.38 h21.6 min43.2 s
99.99%52.6 min4.32 min8.6 s
99.999%5.26 min25.9 s0.86 s

Computed as (1−A)×(1 - A) \times period with a 365-day year. In series, availabilities multiply (two 99.9% hops give 99.8%); in parallel, redundant copies fail together only with probability (1−A)n(1 - A)^n (two independent 99% replicas give 99.99%), but real failures are often correlated (same deploy, same region). See the SRE book on embracing risk (opens in a new tab).

BlockWhat it isReach for it whenCosts
Redundancy (N+1, N+2)spare capacity across zonesalways for anything user-facingidle capacity
Active-passivea standby takes over when the primary failsdatabases, anything with one writerfailover takes seconds to minutes; the standby may lag; must fence the old primary (split-brain)
Active-activeall sites serve trafficstateless tiers; multi-region readsdata conflicts if all sites write; capacity for a site's loss must already be there
Multi-regiondeploy in several regions: read-local/write-home, or partitioned by user's home regionregional disasters, latency for global users, data residencycross-region replication lag, cost, complexity
Timeoutsan upper bound on every network call, derived from the caller's deadlineevery calltoo short: false failures; too long: threads pile up
Retries + exponential backoff + jitterretry transient failures with growing random delaysidempotent operationsamplification: 3 layers × 3 attempts = 27 calls; use retry budgets and retry at one layer (AWS: backoff and jitter (opens in a new tab))
Circuit breakerafter N failures, fail fast (open), probe later (half-open), recover (closed)a dependency that is down or slowtuning thresholds; needs a fallback
Bulkheadseparate thread pools, connection pools or queues per dependencyone slow dependency must not take everything downless pooling efficiency
Rate limitingcap requests per client over timeprotect capacity, fairness, abuselegitimate bursts get 429s
Graceful degradationserve cached, partial or default contentnon-critical features (recommendations, counts)product decisions about what can be dropped
Load sheddingreject early (503) when over capacity, lowest priority firstoverload; queues growing faster than they drainsome users see errors so the rest don't
Rate limit algorithmIn one line
Token buckettokens refill at rate r up to b; a request spends one; allows bursts up to b
Leaky bucketrequests queue and drain at a fixed rate; smooth output, added latency
Fixed windowcounter per key per minute; up to 2x burst at window edges
Sliding window logexact, stores a timestamp per request
Sliding window counterweight the previous window's count by overlap; cheap and close

Distributed versions, headers and a comparison: System design: Rate limiting.

Token bucket

class TokenBucket {
  private tokens: number;
  private last: number;
 
  constructor(
    private readonly capacity: number, // burst size
    private readonly perSec: number, // refill rate
    now: number = performance.now(),
  ) {
    this.tokens = capacity;
    this.last = now;
  }
 
  tryTake(cost = 1, now = performance.now()): boolean {
    const elapsed = (now - this.last) / 1000; // ms → s
    this.tokens = Math.min(
      this.capacity,
      this.tokens + elapsed * this.perSec,
    );
    this.last = now;
    if (this.tokens < cost) return false;
    this.tokens -= cost;
    return true;
  }
}

O(1)O(1) time and space per key: refill lazily from elapsed time instead of running a timer. Across many servers, keep tokens and last per key in Redis and update both in one Lua script so the read-modify-write is atomic. The now parameter makes it testable with a fake clock.

Say it like this

  • Nines: "Four nines is about 52 minutes a year, so failover has to be automatic; a human can't be paged and fix it in time."
  • Retries: "Retries use exponential backoff with full jitter and a retry budget, only at the edge, so we don't multiply load during an outage."
  • Circuit breaker: "If the recommendations service trips its breaker, the page renders without recommendations instead of timing out."
  • Load shedding: "Past 80% of capacity we shed anonymous traffic first and keep checkout working."
  • Multi-region: "Each user has a home region for writes; reads are served locally and cross-region replication is async."

Batch vs stream processing

BatchStream
Inputa bounded dataset (yesterday's logs)an unbounded flow of events
Latencyminutes to hoursmilliseconds to seconds
ToolsMapReduce, Spark, SQL warehouses (BigQuery, Snowflake)Flink, Kafka Streams, Spark Structured Streaming
Reach for it whenreports, ML training sets, backfills, reindexingfraud detection, live counters, alerting, feeds
Costsstale resultsstate management, late events, harder to test and replay

MapReduce: map turns each input record into key-value pairs, the framework shuffles so all values for a key reach the same reducer, reduce combines them. Failed tasks just rerun because inputs are immutable (MapReduce paper (opens in a new tab)).

Stream processing adds windows (tumbling, sliding, session), event time vs processing time, watermarks (how long to wait for late events) and checkpointed state so a restart resumes exactly where it left off.

ArchitectureIdeaTrade-off
Lambdaa batch layer computes accurate views; a speed layer covers recent data; queries merge bothtwo codebases computing the same thing
Kappaeverything is a stream; to recompute, replay the log through a new version of the jobneeds a replayable log with long retention
Change data capture (CDC)stream the database's write-ahead log as events (Debezium)keeps caches, search indexes and warehouses in sync without dual writes

Say it like this

  • Batch vs stream: "Billing reports can be a nightly batch; fraud scoring must be streaming because it has to block the payment."
  • Kappa: "I'd keep events in Kafka with long retention, so fixing a bug means replaying the topic through the new job."
  • CDC: "The search index is fed by CDC from Postgres, so there's no dual write that can drift."

Security & observability

BlockWhat it isReach for it whenCosts
Authentication (authN)who are you: passwords, passkeys, SSOevery user-facing systemaccount recovery, MFA, credential stuffing defenses
Authorization (authZ)what may you do: roles (RBAC), attributes (ABAC), relationships (ReBAC, Zanzibar-style)every request, checked on the serverpolicy sprawl; checks on every hop
OAuth 2.0 / OpenID ConnectOAuth delegates access with scoped tokens; OIDC adds identity (an ID token) on top"Sign in with Google", third-party API accessredirect flows, token storage
JWTsigned, self-contained claims a service can verify without a lookupstateless verification between services, short-lived access tokenshard to revoke before expiry: keep them short (minutes) with refresh tokens
Session IDrandom ID pointing at server-side statebrowser apps you controla session store lookup per request
Encryption in transitTLS everywhere, mTLS between servicesalwayscertificate management
Encryption at restdisk or column encryption with keys in a KMS; envelope encryption (data key encrypted by a master key)always for user data, required by most compliance regimeskey rotation, access policies
Secrets managementVault, cloud secret managers; injected at runtime, rotatedany credentialnever in code, images or logs
ObservabilityWhat it is
Metricsnumbers over time: rate, errors, duration (RED), saturation; cheap, good for alerts
Logsstructured events with detail; include trace IDs
Tracesone request's path across services with timing per hop (OpenTelemetry)
SLIthe measured indicator: "fraction of requests under 300 ms and successful"
SLOthe internal target for an SLI: "99.9% over 30 days"
SLAthe contract with customers, with penalties; looser than the SLO
Error budget1−SLO1 - \text{SLO}: 99.9% of 10M requests a month allows 10,000 failures; when it's spent, freeze risky launches

Dashboards and alerting detail: System design: Observability.

Say it like this

  • AuthN vs authZ: "The gateway authenticates the JWT; each service still authorizes the action against the resource owner."
  • JWT: "Access tokens live 15 minutes; revocation happens at refresh, which checks the session store."
  • SLO: "Our SLO is 99.9% of feed loads under 500 ms; alerts fire on error-budget burn rate, not on single spikes."

Trade-off → phrase

Trade-offPhrase
Consistency vs availability"During a partition I'd rather ___ than ___, because being wrong here costs ___."
Latency vs consistency"Reads come from the nearest replica; the cost is up to a second of staleness, which the feed tolerates."
Read vs write cost"Reads outnumber writes 100:1, so I'll pay on write (fan-out, denormalize) to make reads one lookup."
Freshness vs load"A 60-second TTL takes 95% of reads off the database; nobody notices a like count a minute old."
Simplicity vs scale"One Postgres primary handles this for years; I'd shard only when writes pass what one node sustains, and I'd pick user_id then."
Sync vs async"The user needs an answer for ___ now; everything else goes on a queue."
Exactly once vs cost"At-least-once plus idempotent consumers gives effectively-once without distributed transactions."
Accuracy vs memory"An approximate count within 1% is fine for a dashboard, so HyperLogLog instead of storing every ID."
Build vs buy"Queues and auth aren't our differentiator, so managed services; our effort goes into ranking."
Cost vs availability"Going from three to four nines means multi-region active-active, roughly doubling cost; is 43 minutes a month acceptable?"
Normalize vs denormalize"Normalized for writes, plus a denormalized read model rebuilt from events."

Recipes

"Which database?" checklist

Walk these in order and say each answer out loud.

  • What are the top 3 queries, and by which key? Point lookups, ranges, full text, graph hops, nearest neighbor?
  • Do you need multi-row transactions or constraints (money, inventory, uniqueness)? → relational.
  • Reads and writes per second at peak, and data size in 3 years? Does it fit one primary with replicas?
  • Is the shape a self-contained aggregate (document), append-only by key and time (wide-column, time-series), or files (object store)?
  • Which consistency does each query need (see Consistency & transactions)?
  • Default: Postgres; add Redis (cache), an object store (files) and a search index (projection) as named needs appear.

"Make it faster" checklist

  • Measure first: which endpoint, p50 or p99, where does the time go (trace)?
  • Cache the hot reads (cache-aside, CDN for static and immutable data).
  • Add the missing index; fix N+1 queries; paginate with cursors.
  • Precompute: denormalize, materialized views, fan-out on write.
  • Move work off the request path: queue it and return 202.
  • Cut round trips: batch requests, HTTP/2 or HTTP/3, keep connections alive, colocate services with data.
  • Serve from closer: CDN, edge compute, read replicas per region.
  • Scale out the stateless tier; then read replicas; shard last.

"Make it reliable" checklist

  • Remove single points of failure: N+1 instances across 3 zones, replicated data, automatic failover.
  • Timeouts on every call, retries with backoff and jitter only for idempotent operations, idempotency keys.
  • Circuit breakers and bulkheads around each dependency; a fallback for each.
  • Queues between tiers to absorb spikes; dead-letter queues for poison messages.
  • Rate limiting at the edge; load shedding under overload.
  • Backups you have actually restored; defined RPO (data loss allowed) and RTO (downtime allowed).
  • SLOs with burn-rate alerts; health checks separating liveness from readiness; practiced failover drills.

Pick a consistency level per operation

  • Money, inventory, uniqueness, locks → linearizable (single leader, consensus, or SERIALIZABLE).
  • A user's own recent writes → read-your-writes (read from leader briefly, or wait for replica position).
  • Conversations, comment threads → causal (per-conversation sequence numbers).
  • Counts, likes, feeds, analytics → eventual; show "about" numbers.

Explain any building block in 30 seconds

  1. What it is, in one sentence.
  2. Why this system needs it (tie to a requirement or number).
  3. What it costs and what you'll do about it.
  4. What you'd monitor to know it's working (hit rate, lag, error rate).

Example: "A Redis cache in front of profiles. Reads are 95% of traffic and profiles change rarely, so a 5-minute TTL takes most load off Postgres. The cost is staleness, so we delete the key on update. I'd watch hit rate and p99."

Scale a hot write path

  • Batch: buffer writes and flush every N ms or N items.
  • Make writes append-only (log, LSM store) instead of updates in place.
  • Split hot keys into sub-keys; aggregate on read or in a periodic job.
  • Queue writes and let workers drain at the database's pace.
  • Shard by a key that spreads writes; keep related rows on the same shard.

References