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
| Block | What it is | Reach for it when | Costs |
|---|---|---|---|
| DNS | name → IP; stub resolver → recursive resolver → root → TLD → authoritative server; answers cached for their TTL | geo routing, weighted and failover records, internal service names | failover is only as fast as the TTL (and some clients ignore it); a cold lookup costs a round trip |
| TCP | reliable, ordered byte stream; 3-way handshake (1 RTT), congestion control | nearly everything: HTTP/1.1, HTTP/2, database protocols | handshake latency; one lost packet stalls every byte behind it (head-of-line blocking) |
| UDP | bare datagrams: no handshake, order or retransmit | voice, video, games, DNS; QUIC is built on it | you rebuild reliability and ordering yourself; some networks throttle it |
| TLS 1.3 | encryption 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 services | a round trip, certificate rotation, a little CPU; 0-RTT data can be replayed, so only for idempotent requests |
| HTTP/1.1 | HTTP/2 | HTTP/3 | |
|---|---|---|---|
| Transport | TCP | TCP | QUIC over UDP |
| Concurrency | one request at a time per connection; browsers open about 6 per origin | many streams on one connection | many independent streams |
| Head-of-line blocking | per connection | at the TCP layer: one lost packet stalls all streams | per stream only |
| Headers | plain text, repeated | HPACK compression | QPACK compression |
| New connection | TCP + TLS 1.3 = 2 RTT | same | 1 RTT (0-RTT on resume) |
| Extra | server push, since removed from Chrome and Firefox | connection migration (Wi-Fi to cellular keeps the connection) | |
| Spec | RFC 9112 (opens in a new tab) | RFC 9113 (opens in a new tab) | RFC 9114 (opens in a new tab) |
Server-to-client updates
| Option | How it works | Reach for it when | Costs |
|---|---|---|---|
| Short polling | client asks every N seconds | updates are rare, simplicity wins | wasted requests; up to N seconds late |
| Long polling | server holds the request until there is data or a timeout; the client reconnects | streaming is blocked by proxies | a 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 streams | server → 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) |
| WebSockets | HTTP Upgrade to a full-duplex framed connection (RFC 6455 (opens in a new tab)) | both sides send often: chat, multiplayer, collaborative editing | stateful: connection servers, a pub/sub backplane to reach the right one, reconnect and resync logic, draining on deploy |
| Webhooks | your server POSTs to a URL the customer registered | telling 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
| REST | gRPC | GraphQL | |
|---|---|---|---|
| Shape | resources + HTTP verbs, usually JSON | RPC methods, Protobuf over HTTP/2 | one endpoint; the client picks the fields |
| Strengths | cacheable by browsers and CDNs, universal tooling, easy to debug | small and fast, typed code generation, streaming both ways, deadlines | no over- or under-fetching, nested data in one round trip, typed schema |
| Weaknesses | over-fetching; many round trips for nested data | browsers need gRPC-Web or a proxy; binary is harder to inspect | HTTP caching is hard, N+1 resolvers (batch with DataLoader), must cap query depth and cost |
| Reach for it when | public APIs, CRUD | service-to-service calls inside your network | many clients needing different shapes (web, mobile), one API over many services |
| Pagination | Offset (?page=3&limit=20) | Cursor / keyset (?after=token&limit=20) |
|---|---|---|
| Query | ORDER BY id LIMIT 20 OFFSET 40 | WHERE (created_at, id) < (:t, :id) ORDER BY created_at DESC, id DESC LIMIT 20 |
| Deep pages | slow: the database still reads and skips the offset rows | an index seek per page, constant cost |
| New rows arrive | pages shift: duplicates and skipped items | stable |
| Jump to page N | yes | no, next and previous only |
| Use for | admin tables, small result sets | feeds, infinite scroll, public APIs at scale |
| Concern | Practice |
|---|---|
| Idempotency keys | client 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 retries | GET, PUT, DELETE are idempotent by definition (RFC 9110 (opens in a new tab)); POST and PATCH need a key |
| Versioning | path (/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 |
| Errors | consistent envelope, machine-readable code, application/problem+json (RFC 9457 (opens in a new tab)) |
| Rate limits | per 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 operations | 202 Accepted + a job resource to poll, or a webhook when done |
| Large uploads | presigned 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 /paymentstakes an idempotency key and replays the stored result."
Load balancing & proxies
| Block | What it is | Reach for it when | Costs |
|---|---|---|---|
| L4 load balancer | routes TCP/UDP connections by IP and port | raw throughput, non-HTTP protocols, TLS passthrough | can't see paths or headers; a long-lived connection stays pinned to one backend |
| L7 load balancer | parses HTTP; routes by host, path, header; terminates TLS | microservices, canaries, per-request gRPC balancing, retries | more CPU; it sees plaintext; one more hop |
| Reverse proxy | server in front of your app: TLS, compression, caching, buffering slow clients (nginx, Envoy, Caddy) | always, even with one app server | another component to configure and monitor |
| API gateway | L7 proxy with API concerns: auth, rate limits, quotas, request shaping, aggregation | one public API over many services; a backend-for-frontend per client | business logic creeps in; must be replicated like any hot path |
| Service discovery | registry of healthy instances: DNS (Kubernetes Services), Consul, etcd; client-side or server-side lookup | autoscaled, short-lived instances | the registry must be highly available; stale entries need health checks and TTLs |
| Service mesh | a sidecar proxy per instance for mTLS, retries, telemetry (Istio, Linkerd) | many services, uniform security and traffic policy | extra latency per hop, operational weight |
| Algorithm | Picks | Good for |
|---|---|---|
| Round robin (weighted) | next backend in turn | uniform, short requests |
| Least connections | fewest open connections | long or uneven requests, WebSockets |
| Power of two choices | better of 2 random backends | big fleets; avoids every LB herding onto the same "least loaded" node |
| Consistent hashing | key'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
/apiand/wsto 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 : 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.
| Where | Holds | Reach for it when | Costs |
|---|---|---|---|
| Client (browser, app) | assets, API responses (Cache-Control, ETag) | static assets, per-user data that changes rarely | you can't purge it: use content-hashed file names |
| CDN / edge | static files, images, video segments, cacheable HTML and API responses | global users, large or popular files | purge delays, cost per GB, personalized data needs care |
| Application (in-process LRU) | hot config, small lookups | tiny, very hot data | each instance holds its own copy: inconsistent, lost on restart |
| Distributed cache (Redis, Valkey, Memcached) | objects, sessions, computed views, counters | shared across instances, read-heavy data | a network hop, memory cost, one more thing to fail |
| Database (buffer pool, materialized views) | pages, precomputed query results | expensive aggregates | refresh cost; staleness of the view |
| Pattern | Read | Write | Reach for it when | Costs |
|---|---|---|---|---|
| Cache-aside (lazy) | app checks cache, on miss loads DB and fills | app writes DB, deletes the key | the default; the cache can fail without breaking correctness | first read is slow; a small race can re-cache stale data |
| Read-through | cache library loads from DB on a miss | as cache-aside | you want the loading logic in one place | needs a cache that supports loaders |
| Write-through | from cache | write cache and DB together, synchronously | data read right after it's written | slower writes; caches data nobody reads |
| Write-back (write-behind) | from cache | write cache, flush to DB later in batches | write-heavy counters, metrics | data loss if the cache dies before flushing |
| Write-around | from cache (cache-aside on miss) | write DB only, skip the cache | data written once and rarely read back (logs, uploads) | first read of new data always misses |
| Eviction | Evicts | Good for |
|---|---|---|
| LRU | least recently used | general default, recency-heavy traffic |
| LFU | least frequently used | stable popularity (top products); resists one-off scans |
| TTL | anything past its expiry | bounding staleness; combine with LRU/LFU |
| FIFO / random | oldest / any | cheap, 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)).
| Problem | What happens | Fix |
|---|---|---|
| Invalidation | data changes but the cache still serves the old value | delete 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 once | single-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 key | one key gets so many reads that one cache node saturates | in-process cache in front, replicate the key as key#1..key#N and read a random copy |
| Penetration | requests for keys that don't exist always miss | cache "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
| Type | Examples | Pick when | Avoid when |
|---|---|---|---|
| Relational | Postgres, MySQL; distributed: CockroachDB, Spanner, Vitess, Citus | the default: relations, constraints, transactions, ad-hoc queries | write volume or data size outgrows one primary and sharding by hand hurts |
| Key-value | Redis, DynamoDB, etcd | lookups by a known key: sessions, carts, counters, feature flags | you need queries on anything but the key |
| Document | MongoDB, Firestore, Couchbase | self-contained aggregates read and written together; flexible schema | many-to-many relations and cross-document transactions |
| Wide-column | Cassandra, ScyllaDB, Bigtable, HBase | huge write rates, queries known upfront (partition key + clustering order): messages, events | ad-hoc queries, joins, strong consistency by default |
| Graph | Neo4j, Neptune | many-hop relationship queries: friends of friends, fraud rings | simple lookups; data that isn't a graph |
| Time-series | TimescaleDB, InfluxDB, Prometheus | append-only measurements with time-range queries, rollups, retention | general application data |
| Search index | Elasticsearch, OpenSearch, Meilisearch | full-text relevance, fuzzy matching, facets | as the source of truth (it's a projection) |
| Object / blob | S3, GCS, R2 | files, images, video, backups, data lake | small records, frequent in-place updates |
| Vector | pgvector, Qdrant, Pinecone, Milvus | nearest-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-tree | LSM-tree | |
|---|---|---|
| Writes | update 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) |
| Reads | O(log n) page reads, predictable | check the memtable, then SSTables newest to oldest; Bloom filters skip files that can't hold the key |
| Wins when | read-heavy, range scans, predictable latency | write-heavy ingest, time-series, data much larger than RAM |
| Costs | random writes, page splits | compaction spikes, read amplification, space amplification |
| Used by | Postgres, MySQL InnoDB, SQLite | RocksDB, LevelDB, Cassandra, ScyllaDB, HBase |
| Secondary index on a sharded store | How | Cost |
|---|---|---|
| Local (per shard) | each shard indexes only its own rows | reads must ask every shard (scatter-gather) |
| Global (partitioned by the indexed value) | the index is its own sharded table | writes 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_idas 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
| Step | What | Reach for it when | Costs |
|---|---|---|---|
| Vertical scaling | bigger machine | first; modern boxes go far (hundreds of cores, TBs of RAM) | a ceiling, a single point of failure, price at the top end |
| Read replicas | followers serve reads | read-heavy load | replication lag: stale reads, read-your-writes needs care |
| Caching | see Caching | hot, repeated reads | invalidation |
| Sharding (partitioning) | split rows across nodes by a shard key | writes or data outgrow one primary | cross-shard queries and transactions, rebalancing, operational load |
| Denormalization | store data pre-joined or duplicated for a read path | a hot read needs several joins | writes update several copies; copies can drift |
| Replication | How | Reach for it when | Costs |
|---|---|---|---|
| Synchronous | leader waits for the follower before acknowledging | no data loss on failover | write latency; a slow follower stalls writes (use one sync follower, the rest async) |
| Asynchronous | leader acknowledges, followers catch up | the common default | failover can lose the last writes; stale reads |
| Leader-follower | all writes through one leader | almost always | the leader caps write throughput; failover needs care |
| Multi-leader | a leader per region (or per device), syncing both ways | multi-region writes, offline-first apps | write conflicts: last-writer-wins loses data; merge functions or CRDTs keep it |
| Leaderless (Dynamo-style) | client writes to N replicas, waits for W, reads R | high write availability, tolerating node loss | quorum 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
| Strategy | Reach for it when | Costs |
|---|---|---|
| Range (by key or time) | range scans matter | sequential keys (timestamps) make one hot shard |
Hash (hash(key) mod N) | even spread, point lookups | no range scans; changing N moves almost every key |
| Consistent hashing + virtual nodes | nodes 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 customers | the lookup service is on every request path |
| Geo | data residency, latency per region | uneven region sizes |
Detail and shard-key advice: System design: Sharding.
Consistent hashing
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
}
}// FNV-1a over UTF-8 bytes, then a murmur3 finalizer
// so near-identical strings land far apart
function hash32(s, seed = 0x811c9dc5) {
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;
}
class HashRing {
#ring = [];
#vnodes;
constructor(vnodes = 100) {
this.#vnodes = vnodes;
}
add(node) {
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) {
this.#ring = this.#ring.filter((v) => v.node !== node);
}
nodeFor(key) {
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
}
}from bisect import bisect_left
M32 = 0xFFFFFFFF
def hash32(s: str, seed: int = 0x811C9DC5) -> int:
h = seed
for b in s.encode():
h = ((h ^ b) * 0x01000193) & M32
h = ((h ^ (h >> 16)) * 0x85EBCA6B) & M32
h = ((h ^ (h >> 13)) * 0xC2B2AE35) & M32
return h ^ (h >> 16)
class HashRing:
def __init__(self, vnodes: int = 100) -> None:
self.vnodes = vnodes
self.ring: list[tuple[int, str]] = [] # (pos, node)
def add(self, node: str) -> None:
for i in range(self.vnodes):
self.ring.append((hash32(f"{node}#{i}"), node))
self.ring.sort()
def remove(self, node: str) -> None:
self.ring = [v for v in self.ring if v[1] != node]
def node_for(self, key: str) -> str:
if not self.ring:
raise LookupError("empty")
h = hash32(key)
i = bisect_left(self.ring, h, key=lambda v: v[0])
return self.ring[i % len(self.ring)][1] # wrapadd: for virtual nodes in total (re-sort); nodeFor: ; space . 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
| Problem | Fix |
|---|---|
| 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 tenants | directory-based placement; move big tenants to their own shard |
| Resharding a live system | double-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/Nof 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.
| Model | Guarantee | Reach for it when | Costs |
|---|---|---|---|
| Linearizable (strong) | every read sees the latest completed write; behaves like one copy | balances, inventory, unique usernames, locks, leader election | coordination on every operation: latency, unavailable in a partition |
| Causal | effects never appear before their causes, for everyone | comments and replies, chat | tracking dependencies |
| Read-your-writes | you always see your own writes | profile edits, posting then viewing | routing or waiting per user |
| Monotonic reads | you never see time go backward | any replica-served UI | pin a user to one replica |
| Eventual | replicas converge once writes stop | likes, view counts, feeds, DNS | temporary anomalies users may notice |
Definitions in depth: Jepsen: consistency models (opens in a new tab).
Quorums
With replicas, a write waits for acknowledgments and a read asks replicas. If , every read set overlaps every write set, so at least one replica in the read has the latest acknowledged write (pick it by version). tolerates one slow or dead node for both reads and writes. gives fast reads and fragile writes; 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).
| Anomaly | What goes wrong |
|---|---|
| Dirty read | you read another transaction's uncommitted write |
| Non-repeatable read | the same row read twice returns different values |
| Phantom | the same query returns new rows the second time |
| Lost update | two read-modify-write cycles race; one overwrites the other |
| Write skew | two transactions read the same data, write different rows, and together break a rule ("at least one doctor on call") |
| Level | Prevents | Still allows | Notes |
|---|---|---|---|
| Read committed | dirty reads and writes | non-repeatable reads, phantoms, lost updates, write skew | Postgres default |
| Repeatable read / snapshot isolation | the above + non-repeatable reads | write skew (phantoms in the SQL standard) | MySQL InnoDB default; Postgres aborts lost updates here, InnoDB doesn't |
| Serializable | all of the above | nothing; conflicting transactions abort and retry | Postgres 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 | |
|---|---|---|
| How | coordinator asks all participants to prepare, then commit | a chain of local transactions; each failure runs compensating actions (refund, release) |
| Guarantee | atomic across databases | eventually consistent; intermediate states are visible |
| Reach for it when | one organization, few participants, databases that support XA | microservices, long-running business flows (order → payment → shipping) |
| Costs | blocking: a coordinator crash leaves participants holding locks; latency | compensations 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 UPDATEon 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) | |
|---|---|---|
| Model | messages deleted once acknowledged | append-only partitioned log; consumers track an offset |
| Consumers | competing workers share one queue | consumer groups: each partition goes to one consumer in the group; many groups read independently |
| Replay | no | yes, within retention |
| Ordering | FIFO queues per group ID; standard queues best-effort | within a partition only |
| Reach for it when | background jobs, task distribution, per-message retries | event streams, CDC, analytics, several independent consumers of the same events |
| Delivery | Meaning | How |
|---|---|---|
| At most once | may lose, never duplicates | acknowledge before processing |
| At least once | never lose, may duplicate | acknowledge after processing (the usual default) |
| Exactly once | each message affects state once | only inside one system (Kafka transactions for read-process-write); anywhere else it's effectively once: at-least-once + an idempotent consumer |
| Concern | Practice |
|---|---|
| Idempotent consumer | record processed message IDs in the same transaction as the side effect; or make writes naturally idempotent (upsert, set instead of increment) |
| Ordering | choose the partition key so related events share a partition (order_id); parallelism is capped by the partition count |
| Poison messages | after N failed attempts move to a dead-letter queue, alert, fix, replay |
| Visibility timeout / ack deadline | longer than processing time, or the message is redelivered mid-work |
| Backpressure | bounded queues; watch consumer lag; slow producers or shed load instead of buffering forever |
| Retries | backoff 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) | Fan-out on read (pull) | |
|---|---|---|
| On post | append the post ID to every follower's feed list | store the post once |
| On read | read one precomputed list | fetch recent posts of every followee, merge by time |
| Reach for it when | most users, most reads | celebrities (millions of followers), inactive readers |
| Costs | write amplification; wasted work for inactive followers | slow, 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
| Block | What it is | Reach for it when | Costs |
|---|---|---|---|
| Leader election | exactly one node acts as leader (scheduler, partition owner, primary) | one writer or one cron runner is needed | election time is downtime; two leaders at once is the failure to prevent |
| Consensus (Raft, Paxos, ZAB) | a majority agrees on an ordered log of decisions | leader election, config, metadata, locks | 2f+1 nodes to survive f failures; every write is a majority round trip |
| ZooKeeper / etcd / Consul | small, strongly consistent stores built on consensus, with watches and leases | cluster metadata, service discovery, locks (Kubernetes stores its state in etcd) | not for application data or high write rates |
| Distributed lock with lease | a lock that expires unless renewed | making sure only one worker does a job | a paused holder can wake up after expiry and still act: needs fencing |
| Fencing token | a number that increases with each lock grant; storage rejects writes with an older token | any lock protecting a write | the 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
| Clock | What it gives | Limit |
|---|---|---|
| Wall clock (NTP-synced) | real time | skew of milliseconds or more between machines; can jump backward; never order events across machines by it |
| Monotonic clock | durations on one machine | meaningless across machines |
| Lamport clock | counter: max(local, received) + 1; if a caused b, then L(a) < L(b) | can't tell whether two events were concurrent |
| Vector clock | a counter per node; compares as before, after or concurrent | size grows with the number of nodes |
| Hybrid logical clock / TrueTime | physical 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
| Structure | Answers | Accuracy and space | Where 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 sketch | approximate count per item | only overestimates; error ≤ εN with probability 1 − δ using ⌈e/ε⌉ × ⌈ln(1/δ)⌉ counters | heavy hitters, trending hashtags, top-k with a heap |
| HyperLogLog | approximate number of distinct items | standard error ≈ for registers; Redis uses : at most 12 KB per key, 0.81% standard error, mergeable (Redis docs (opens in a new tab)) | unique visitors per page per day |
| Geohash | a base32 string per cell; shared prefix ≈ nearby | 5 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 |
| Quadtree | a tree that splits a cell into 4 when it holds more than k points | adapts to density (Manhattan splits deep, deserts don't) | in-memory proximity index |
| S2 / H3 | hierarchical cells on the sphere: S2 squares on a Hilbert curve (Google), H3 hexagons (Uber) | uniform cells without pole distortion; H3 neighbors are all equidistant | maps, ride dispatch, surge pricing zones |
| Trie (prefix tree) | all keys with a given prefix | memory-heavy; store the top-k completions at each node | typeahead |
| Inverted index | term → sorted list of document IDs | intersection of posting lists answers AND queries | full-text search (Lucene, Elasticsearch, Postgres GIN) |
| Skip list | sorted set with O(log n) expected search and insert via random "express lanes" | simpler than balanced trees, lock-friendly | Redis sorted sets, LSM memtables |
| SSTable | immutable file of sorted key-value pairs plus a sparse index | merge-friendly | LSM-tree storage (see Data stores) |
| Merkle tree | a tree of hashes; equal roots mean equal data | find differing ranges by comparing O(log n) hashes, not all data | anti-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 bits, hash functions and items inserted, the false-positive rate is
. Sizing for items at a target : bits and
hash functions. The positions come from two hashes as (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,
);
}
}class BloomFilter {
#bits;
#m;
#k;
constructor(
m, // bits
k, // hash functions
) {
this.#m = m;
this.#k = k;
this.#bits = new Uint8Array(Math.ceil(m / 8));
}
static forCapacity(n, p) {
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);
}
#indexes(item) {
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) {
for (const i of this.#indexes(item)) {
this.#bits[i >> 3] |= 1 << (i & 7);
}
}
mightContain(item) {
return this.#indexes(item).every(
(i) => (this.#bits[i >> 3] & (1 << (i & 7))) !== 0,
);
}
}import math
class BloomFilter:
def __init__(self, m: int, k: int) -> None:
self.m, self.k = m, k # bits, hash functions
self.bits = bytearray((m + 7) // 8)
@classmethod
def for_capacity(cls, n: int, p: float) -> "BloomFilter":
m = math.ceil(-n * math.log(p) / math.log(2) ** 2)
k = max(1, round(m / n * math.log(2)))
return cls(m, k)
def _indexes(self, item: str) -> list[int]:
h1 = hash32(item)
h2 = hash32(item, 0x9747B28C)
return [(h1 + i * h2) % self.m
for i in range(self.k)]
def add(self, item: str) -> None:
for i in self._indexes(item):
self.bits[i >> 3] |= 1 << (i & 7)
def might_contain(self, item: str) -> bool:
return all(
self.bits[i >> 3] & (1 << (i & 7))
for i in self._indexes(item)
)add and mightContain: time; space 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
| Availability | Downtime / year | Downtime / 30 days | Downtime / day |
|---|---|---|---|
| 99% | 3.65 days | 7.2 h | 14.4 min |
| 99.5% | 1.83 days | 3.6 h | 7.2 min |
| 99.9% | 8.76 h | 43.2 min | 1.44 min |
| 99.95% | 4.38 h | 21.6 min | 43.2 s |
| 99.99% | 52.6 min | 4.32 min | 8.6 s |
| 99.999% | 5.26 min | 25.9 s | 0.86 s |
Computed as 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 (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).
| Block | What it is | Reach for it when | Costs |
|---|---|---|---|
| Redundancy (N+1, N+2) | spare capacity across zones | always for anything user-facing | idle capacity |
| Active-passive | a standby takes over when the primary fails | databases, anything with one writer | failover takes seconds to minutes; the standby may lag; must fence the old primary (split-brain) |
| Active-active | all sites serve traffic | stateless tiers; multi-region reads | data conflicts if all sites write; capacity for a site's loss must already be there |
| Multi-region | deploy in several regions: read-local/write-home, or partitioned by user's home region | regional disasters, latency for global users, data residency | cross-region replication lag, cost, complexity |
| Timeouts | an upper bound on every network call, derived from the caller's deadline | every call | too short: false failures; too long: threads pile up |
| Retries + exponential backoff + jitter | retry transient failures with growing random delays | idempotent operations | amplification: 3 layers × 3 attempts = 27 calls; use retry budgets and retry at one layer (AWS: backoff and jitter (opens in a new tab)) |
| Circuit breaker | after N failures, fail fast (open), probe later (half-open), recover (closed) | a dependency that is down or slow | tuning thresholds; needs a fallback |
| Bulkhead | separate thread pools, connection pools or queues per dependency | one slow dependency must not take everything down | less pooling efficiency |
| Rate limiting | cap requests per client over time | protect capacity, fairness, abuse | legitimate bursts get 429s |
| Graceful degradation | serve cached, partial or default content | non-critical features (recommendations, counts) | product decisions about what can be dropped |
| Load shedding | reject early (503) when over capacity, lowest priority first | overload; queues growing faster than they drain | some users see errors so the rest don't |
| Rate limit algorithm | In one line |
|---|---|
| Token bucket | tokens refill at rate r up to b; a request spends one; allows bursts up to b |
| Leaky bucket | requests queue and drain at a fixed rate; smooth output, added latency |
| Fixed window | counter per key per minute; up to 2x burst at window edges |
| Sliding window log | exact, stores a timestamp per request |
| Sliding window counter | weight 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;
}
}class TokenBucket {
#tokens;
#last;
#capacity;
#perSec;
constructor(
capacity, // burst size
perSec, // refill rate
now = performance.now(),
) {
this.#capacity = capacity;
this.#perSec = perSec;
this.#tokens = capacity;
this.#last = now;
}
tryTake(cost = 1, now = performance.now()) {
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;
}
}import time
class TokenBucket:
def __init__(self, capacity: float, per_sec: float,
now: float | None = None) -> None:
self.capacity = capacity # burst size
self.per_sec = per_sec # refill rate
self.tokens = capacity
self.last = time.monotonic() if now is None else now
def try_take(self, cost: float = 1,
now: float | None = None) -> bool:
now = time.monotonic() if now is None else now
elapsed = now - self.last # already seconds
self.tokens = min(
self.capacity,
self.tokens + elapsed * self.per_sec,
)
self.last = now
if self.tokens < cost:
return False
self.tokens -= cost
return True 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
| Batch | Stream | |
|---|---|---|
| Input | a bounded dataset (yesterday's logs) | an unbounded flow of events |
| Latency | minutes to hours | milliseconds to seconds |
| Tools | MapReduce, Spark, SQL warehouses (BigQuery, Snowflake) | Flink, Kafka Streams, Spark Structured Streaming |
| Reach for it when | reports, ML training sets, backfills, reindexing | fraud detection, live counters, alerting, feeds |
| Costs | stale results | state 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.
| Architecture | Idea | Trade-off |
|---|---|---|
| Lambda | a batch layer computes accurate views; a speed layer covers recent data; queries merge both | two codebases computing the same thing |
| Kappa | everything is a stream; to recompute, replay the log through a new version of the job | needs 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
| Block | What it is | Reach for it when | Costs |
|---|---|---|---|
| Authentication (authN) | who are you: passwords, passkeys, SSO | every user-facing system | account recovery, MFA, credential stuffing defenses |
| Authorization (authZ) | what may you do: roles (RBAC), attributes (ABAC), relationships (ReBAC, Zanzibar-style) | every request, checked on the server | policy sprawl; checks on every hop |
| OAuth 2.0 / OpenID Connect | OAuth delegates access with scoped tokens; OIDC adds identity (an ID token) on top | "Sign in with Google", third-party API access | redirect flows, token storage |
| JWT | signed, self-contained claims a service can verify without a lookup | stateless verification between services, short-lived access tokens | hard to revoke before expiry: keep them short (minutes) with refresh tokens |
| Session ID | random ID pointing at server-side state | browser apps you control | a session store lookup per request |
| Encryption in transit | TLS everywhere, mTLS between services | always | certificate management |
| Encryption at rest | disk 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 regimes | key rotation, access policies |
| Secrets management | Vault, cloud secret managers; injected at runtime, rotated | any credential | never in code, images or logs |
| Observability | What it is |
|---|---|
| Metrics | numbers over time: rate, errors, duration (RED), saturation; cheap, good for alerts |
| Logs | structured events with detail; include trace IDs |
| Traces | one request's path across services with timing per hop (OpenTelemetry) |
| SLI | the measured indicator: "fraction of requests under 300 ms and successful" |
| SLO | the internal target for an SLI: "99.9% over 30 days" |
| SLA | the contract with customers, with penalties; looser than the SLO |
| Error budget | : 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-off | Phrase |
|---|---|
| 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
- What it is, in one sentence.
- Why this system needs it (tie to a requirement or number).
- What it costs and what you'll do about it.
- 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
- Amazon Dynamo paper (2007) (opens in a new tab): consistent hashing, vnodes, quorums, hinted handoff, Merkle trees
- Raft paper (opens in a new tab) and raft.github.io (opens in a new tab): consensus with a visualization
- Google Bigtable paper (opens in a new tab): wide-column model, SSTables
- MapReduce paper (opens in a new tab): the batch model
- Scaling Memcache at Facebook (opens in a new tab): leases against stampedes and stale sets
- RFC 9110: HTTP semantics (opens in a new tab), RFC 9114: HTTP/3 (opens in a new tab), RFC 8446: TLS 1.3 (opens in a new tab), RFC 6455: WebSocket (opens in a new tab): protocol facts
- MDN: Server-sent events (opens in a new tab):
EventSource, reconnection, connection limits - PostgreSQL: transaction isolation (opens in a new tab): what each level prevents in Postgres
- Redis: HyperLogLog (opens in a new tab) and key eviction (opens in a new tab): sizes, error, eviction policies
- Kleppmann: How to do distributed locking (opens in a new tab): leases and fencing tokens
- Jepsen: consistency models (opens in a new tab): precise definitions
- Designing Data-Intensive Applications (opens in a new tab): the book behind most of this sheet (replication, partitioning, transactions, streams)
- The System Design Primer (opens in a new tab): broad index of the same topics
- Google SRE book (opens in a new tab): availability, SLOs, overload handling
- AWS: exponential backoff and jitter (opens in a new tab): why full jitter
- Twitter: Timelines at Scale (opens in a new tab): the hybrid fan-out in production
- System design (production reference): numbers, tables and the worked URL shortener