00 to 84Bold ideas Brighter designs Bigger systems
System design,from first array to Netflix scale
Start with Big O and the handful of data structures every system is built from, learn the algorithms hiding inside caches and load balancers, then scale, store, cache, queue and replicate until you can design Bitly, Twitter, WhatsApp, Netflix, Uber and a payment system out loud.
By Shree Kumar Sharma, with Claude Design, for Backend Engineering. Progress is saved in this browser only.
A catalog of every keyword in this guide. Each entry says what the word means in plain language, what it means technically, and what it is like in everyday life, then links to the module that teaches it.
In detail
Use it as a reference rather than a lesson. Filter by kind to see only data structures, storage words or distributed systems terms, or type to find one. Come back whenever a case study uses a word you have not met yet.
Big O
Algorithms
In plain words
How much slower something gets as the input grows.
Technically
Asymptotic upper bound on time or space as a function of input size n, ignoring constants.
Think of it as
Whether a queue grows by one minute or doubles every time a person joins.
Every big system is a few simple data structures at enormous scale. A cache is a hash table, a message queue is a queue, a database index is a tree, and a feed ranks items with a heap.
In detail
System design interviews rarely ask you to code a red black tree, but they constantly ask why a choice is fast or slow. Knowing what each structure costs lets you explain why Redis answers in a millisecond, why an index makes a query fast, and why a queue protects a service from spikes. Start here and every later module will feel like a familiar shape wearing a bigger coat.
Big O describes how the work grows as the input grows. O(1) stays flat, O(log n) grows slowly, O(n) grows in step, and O(n squared) explodes.
In detail
We ignore constants and keep the fastest growing term, because at a million users only the shape of the curve matters. Time complexity counts steps; space complexity counts memory. In system design the same idea scales up: a full table scan is O(n) per request, an index lookup is O(log n), a cache hit is O(1). Amortised cost means an operation is usually cheap even if it is occasionally expensive, like a dynamic array resizing.
Steps needed as n grows
Complexity
n = 10
n = 1,000
n = 1 million
Example
O(1)
1
1
1
Hash lookup
O(log n)
3
10
20
Binary search, B-tree
O(n)
10
1,000
1,000,000
Scan a list
O(n log n)
33
10,000
20,000,000
Good sort
O(n²)
100
1,000,000
1012, hours
Compare every pair
Cost of each operation
Hash lookupO(1)
Binary searchO(log n)
Scan a listO(n)
Good sortO(n log n)
Compare all pairsO(n^2)
complexity.tsTypeScript
/** O(n) time, O(1) space: check every item. */exportfunctioncontainsLinear(items: number[], x: number): boolean {for (const item of items) {if (item === x) returntrue; }returnfalse;}
/** O(n) time, O(1) space: check every item. */exportfunctioncontainsLinear(items: number[], x: number): boolean {for (const item of items) {if (item === x) returntrue; }returnfalse;}/** O(log n) time: halve the search space each step. */exportfunctioncontainsSorted(items: number[], x: number): boolean {let lo = 0;let hi = items.length;while (lo < hi) {const mid = (lo + hi) >>> 1;if (items[mid] < x) lo = mid + 1;else hi = mid; }return lo < items.length && items[lo] === x;}/** O(1) average time: hash straight to the bucket. */exportfunctioncontainsSet(items: Set<number>, x: number): boolean {return items.has(x);}/** O(n^2): compares every pair. Fine for 100 items, hopeless for 10 million. */exportfunctionhasDuplicatePair(items: number[]): boolean {for (let i = 0; i < items.length; i++) {for (let j = i + 1; j < items.length; j++) {if (items[i] === items[j]) returntrue; } }returnfalse;}
Why it matters The same question asked three ways costs n, log n or 1 steps. Choosing the structure is choosing the curve.
An array stores items side by side in memory, so reading any position is instant but inserting in the middle shifts everything after it.
In detail
Because elements are contiguous, the CPU cache loves arrays: scanning them is the fastest loop in computing. Dynamic arrays (Python lists, Java ArrayList) double their capacity when full, so appends are O(1) amortised. Strings are arrays of characters; building strings in a loop by concatenation can be O(n squared), which is why we join lists. Two pointer and sliding window techniques solve many array problems in one pass, and the same sliding window idea reappears in rate limiting.
index lookup jumps straight to a slot
17[0]
4[1]
42[2]
8[3]
23[4]
15[5]
9[6]
Cost of each operation
Read by indexO(1)
AppendO(1)
Insert in middleO(n)
Search unsortedO(n)
sliding-window.tsTypeScript
let window = nums.slice(0, k).reduce((a, b) => a + b, 0);let best = window;for (let i = k; i < nums.length; i++) { window += nums[i] - nums[i - k]; // slide: add the new item, drop the oldest best = Math.max(best, window);}return best;
/** Largest sum of any k consecutive items in O(n) instead of O(n*k). */functionmaxSumWindow(nums: number[], k: number): number {if (k > nums.length) {thrownewRangeError("window larger than input"); }let window = nums.slice(0, k).reduce((a, b) => a + b, 0);let best = window;for (let i = k; i < nums.length; i++) { window += nums[i] - nums[i - k]; // slide: add the new item, drop the oldest best = Math.max(best, window); }return best;}/** Two pointers on a sorted array: O(n) time, O(1) space. */functiontwoSumSorted(nums: number[], target: number): [number, number] | null {let lo = 0;let hi = nums.length - 1;while (lo < hi) {const s = nums[lo] + nums[hi];if (s === target) {return [lo, hi]; }if (s < target) { lo++; } else { hi--; } }returnnull;}console.log(maxSumWindow([3, 1, 4, 1, 5, 9, 2, 6], 3)); // 17 (9 + 2 + 6)console.log(twoSumSorted([1, 3, 4, 6, 9], 10)); // [0, 4]
Why it matters Sliding windows turn repeated sums into one pass, the same trick a rate limiter uses over time.
A hash table turns a key into a bucket number with a hash function, so finding a value takes the same time whether you store ten items or ten million.
In detail
A good hash function spreads keys evenly. When two keys land in the same bucket (a collision) the table chains them in a small list or probes to the next slot. When it fills past a load factor it resizes and rehashes. Python dicts, Java HashMaps, Redis, Memcached and every database hash index use this idea. At system scale the same trick partitions data across machines, which is exactly where consistent hashing comes in.
Keys hashed into 5 buckets; bucket 3 shows a collision chained in one slot
A linked list stores each item with a pointer to the next one. Inserting or removing next to a node you already hold is instant, but finding the fifth item means walking from the start.
In detail
Singly linked lists point forward; doubly linked lists point both ways, which lets you remove a node in O(1) when you have it. That property is why an LRU cache pairs a hash map with a doubly linked list. Linked lists use more memory per item and are unfriendly to CPU caches, so in practice arrays win most scans, but the pointer idea reappears everywhere: log segments, skip lists in Redis sorted sets, and blockchain style hash chains.
Each node points to the next; the last points to nothing
A stack hands back the last thing you put in, a queue hands back the first. Undo buttons and function calls are stacks; print jobs, request buffers and message brokers are queues.
In detail
Stacks power recursion, expression parsing and depth first search. Queues power breadth first search, task scheduling and every buffer between a fast producer and a slow consumer. A deque supports both ends. A circular buffer is a fixed size queue that overwrites the oldest entry, which is how logs and metrics rings keep memory bounded. Message queues such as SQS, RabbitMQ and Kafka are queues made durable and shared across machines.
First in, first out: enqueue at the back, dequeue from the front
// Stack: last in, first outconst undo: string[] = [];undo.push("type 'hello'");undo.push("bold");console.log(undo.pop()); // "bold"// Queue: first in, first out, buffering bursts between producer and consumerconst MAX_QUEUED = 1000; // bounded: protects memoryconst requests: string[] = [];for (let i = 0; i < 5; i++) {if (requests.length === MAX_QUEUED) requests.shift(); // full: drop the oldest requests.push(`req-${i}`);}while (requests.length > 0) {const handle = requests.shift();}// Balanced brackets, the classic stack problemfunctionbalanced(s: string): boolean {const pairs: Record<string, string> = { ")": "(", "]": "[", "}": "{" };const stack: string[] = [];for (const ch of s) {if ("([{".includes(ch)) { stack.push(ch); } elseif (ch in pairs) {if (stack.length === 0 || stack.pop() !== pairs[ch]) returnfalse; } }return stack.length === 0;}console.log(balanced("{[()]}")); // true
Why it matters Why it matters Capping the queue at a fixed size and dropping the oldest entry makes it a bounded buffer: a bounded queue is the simplest form of back pressure.
A tree organises data in parent and child levels. A balanced search tree keeps keys sorted so lookups, inserts and range scans all take O(log n).
In detail
A binary search tree keeps smaller keys left and larger keys right; it must stay balanced (AVL, red black) or it degrades into a list. Databases use B-trees and B+ trees instead: each node holds hundreds of keys so the tree is only three or four levels deep, which means three or four disk reads to find any row among millions. Range queries walk the linked leaves in order. Traversals (in order, pre order, post order, level order) are the vocabulary of every tree problem.
A heap keeps the smallest (or largest) item on top, so you can always grab the most urgent thing in O(1) and add or remove items in O(log n).
In detail
A binary heap is stored as a plain array where the children of index i sit at 2i plus 1 and 2i plus 2. Priority queues schedule jobs, run Dijkstra's shortest path, merge sorted streams and keep the top K items of an endless stream with a fixed size heap. Leaderboards, trending topics, timers in event loops and the next expiring key in a cache all rely on this structure.
Min heap: every parent is smaller than its childrenServiceCacheData
Cost of each operation
Peek minO(1)
PushO(log n)
Pop minO(log n)
Build from listO(n)
top-k.tsTypeScript
heap.push([count, word]);if (heap.size > k) { heap.pop(); // drop the smallest of the current top k}
/** A small binary min heap: push and pop in O(log n), smallest item on top. */classMinHeap<T> {privatereadonly items: T[] = [];privatereadonly less: (a: T, b: T) => boolean;constructor(less: (a: T, b: T) => boolean) {this.less = less; }getsize(): number {returnthis.items.length; }push(item: T): void {const a = this.items; a.push(item);for (let i = a.length - 1; i > 0; ) {const p = (i - 1) >> 1;if (!this.less(a[i], a[p])) break; [a[i], a[p]] = [a[p], a[i]]; i = p; } }pop(): T | undefined {const a = this.items;const top = a[0];const last = a.pop();if (a.length > 0 && last !== undefined) { a[0] = last;for (let i = 0; ; ) {const l = 2 * i + 1;let m = i;if (l < a.length && this.less(a[l], a[m])) m = l;if (l + 1 < a.length && this.less(a[l + 1], a[m])) m = l + 1;if (m === i) break; [a[i], a[m]] = [a[m], a[i]]; i = m; } }return top; }toArray(): T[] {return [...this.items]; }}/** Top k by frequency using a min heap of size k: O(n log k). */functiontopKWords(stream: string[], k: number): [string, number][] {const counts = newMap<string, number>();for (const word of stream) counts.set(word, (counts.get(word) ?? 0) + 1);// Ordered by count, then by word, the same as comparing [count, word] pairsconst heap = newMinHeap<[number, string]>( (x, y) => x[0] < y[0] || (x[0] === y[0] && x[1] < y[1]), );for (const [word, count] of counts) { heap.push([count, word]);if (heap.size > k) { heap.pop(); // drop the smallest of the current top k } }return heap .toArray() .map(([c, w]): [string, number] => [w, c]) .sort((x, y) => y[1] - x[1]);}/** k way merge, how databases merge sorted runs and SSTables. */functionmergeSorted(...lists: number[][]): number[] {// [value, which list, position in that list]const heap = newMinHeap<[number, number, number]>((x, y) => x[0] < y[0]); lists.forEach((list, li) => {if (list.length > 0) heap.push([list[0], li, 0]); });const out: number[] = [];for (let top = heap.pop(); top !== undefined; top = heap.pop()) {const [value, li, pos] = top; out.push(value);if (pos + 1 < lists[li].length) heap.push([lists[li][pos + 1], li, pos + 1]); }return out;}console.log(topKWords("a b a c b a d a b".split(" "), 2)); // [["a", 4], ["b", 3]]console.log(mergeSorted([1, 4, 9], [2, 3, 10], [5])); // [1, 2, 3, 4, 5, 9, 10]
Why it matters A size k heap uses O(k) memory no matter how long the stream is, which is why it scales to trending topics.
A graph is a set of nodes joined by edges. Social networks, road maps, service dependencies and the web itself are graphs, and two simple walks explore them: breadth first and depth first.
In detail
Store graphs as adjacency lists for sparse data. BFS explores level by level with a queue and finds shortest paths in unweighted graphs, which is how you compute friends of friends. DFS dives deep with a stack or recursion and finds cycles and connected components. Dijkstra adds weights for routing, topological sort orders build steps or task dependencies, and PageRank scores web pages by the links between them. A web crawler is BFS over the internet.
BFS spreads out one hop at a timeServiceEdgeData
Cost of each operation
BFS or DFSO(V + E)
Dijkstra with heapO(E log V)
Adjacency check, listO(degree)
bfs.tsTypeScript
constqueue: [string, number][] = [[start, 0]];const seen = newSet([start]);for (let head = 0; head < queue.length; head++) {const [node, dist] = queue[head];for (const nxt of graph[node]) {
typeGraph = Record<string, string[]>;const graph: Graph = { asha: ["ben", "chen"], ben: ["asha", "dev"], chen: ["asha", "dev", "eli"], dev: ["ben", "chen"], eli: ["chen"],};// BFS: shortest number of hops in an unweighted graph.functiondegreesOfSeparation(graph: Graph, start: string, goal: string): number {constqueue: [string, number][] = [[start, 0]];const seen = newSet([start]);// An array read through a moving head index pops from the front in O(1).for (let head = 0; head < queue.length; head++) {const [node, dist] = queue[head];if (node === goal) {return dist; }for (const nxt of graph[node]) {if (!seen.has(nxt)) { seen.add(nxt);queue.push([nxt, dist + 1]); } } }return -1;}functionfriendsOfFriends(graph: Graph, user: string): Set<string> {const direct = newSet(graph[user]);const fofs = [...direct].flatMap((f) => graph[f]);returnnewSet(fofs.filter((fof) => !direct.has(fof) && fof !== user));}console.log(degreesOfSeparation(graph, "asha", "eli")); // 2console.log(friendsOfFriends(graph, "asha")); // Set(2) { 'dev', 'eli' }
Why it matters The seen set prevents loops; a crawler needs the same thing, only stored in a distributed set or bloom filter.
A trie stores strings letter by letter so words with the same prefix share a path. Finding every word that starts with a prefix takes time proportional to the prefix, not the dictionary.
In detail
Each node holds children keyed by the next character and a marker for word endings. Search autocomplete keeps a trie of popular queries and caches the top results at each node, so a keystroke is answered by walking a few nodes. Routers use a compressed variant (radix tree) to match URL paths and IP prefixes. Tries trade memory for speed, so large systems shard them by first letter or precompute results into a key value store.
Words starting with ne share one pathServiceEdgeData
Sorting arranges data so it can be searched by halving. Good general sorts take O(n log n), and binary search then finds anything in O(log n).
In detail
Merge sort splits, sorts halves and merges; it is stable and the basis of external sorting when data does not fit in memory, which is how databases and MapReduce sort terabytes. Quicksort is usually fastest in memory. Counting and radix sorts beat n log n for small integer ranges. Binary search is far more general than finding a number: you can binary search a version history for the first bad deploy or a capacity setting for the largest safe value.
Bars rise into sorted order, then a binary search halves the range each step
5
2
8
3
9
1
6
4
Cost of each operation
Merge sortO(n log n)
Quicksort averageO(n log n)
Binary searchO(log n)
Bubble sortO(n^2)
binary-search.tsTypeScript
let lo = 1;let hi = n;while (lo < hi) {const mid = Math.floor((lo + hi) / 2);
functionmergeSort(a: number[]): number[] {if (a.length <= 1) return a;const mid = Math.floor(a.length / 2);const left = mergeSort(a.slice(0, mid));const right = mergeSort(a.slice(mid));const out: number[] = [];let i = 0;let j = 0;while (i < left.length && j < right.length) {if (left[i] <= right[j]) out.push(left[i++]);else out.push(right[j++]); }return [...out, ...left.slice(i), ...right.slice(j)];}/** Binary search over deploys: O(log n) checks instead of n. */functionfirstBadVersion(n: number, isBad: (version: number) => boolean): number {let lo = 1;let hi = n;while (lo < hi) {const mid = Math.floor((lo + hi) / 2);if (isBad(mid)) hi = mid;else lo = mid + 1; }return lo;}console.log(mergeSort([38, 27, 43, 3, 9, 82, 10]));console.log(firstBadVersion(1000, (v) => v >= 734)); // 734 in about 10 checks
Why it matters git bisect is first_bad_version over commits. Knowing the pattern makes debugging production regressions fast.
02
Phase 02, modules 12 to 20
Algorithms that run the internet
Hashing, limiting, caching and IDs
The small, clever algorithms hiding inside load balancers, caches, databases and URL shorteners.
Consistent hashing places servers and keys on a ring. Each key belongs to the next server clockwise, so adding or removing a server only moves the keys next to it.
In detail
With plain modulo hashing, changing the number of servers from 4 to 5 moves almost every key, which would flush a cache cluster. On a ring only about 1 over n of the keys move. Virtual nodes give each server many points on the ring so load spreads evenly and a big machine can take more points. DynamoDB, Cassandra, Discord, Akamai and most distributed caches use this idea for partitioning and replication.
Each key belongs to the next node clockwise; adding a node only steals keys from its neighbourBCDAk1→Dk2→Bk3→Ck4→Ck5→C
Cost of each operation
Find nodeO(log v)
Add nodeO(v log v)
Keys moved on changeabout K/n
hash-ring.tsTypeScript
const idx = bisectRight(this.keys, h(key)) % this.keys.length; // next point clockwisereturnthis.ring.get(this.keys[idx]) asstring;
import { createHash } from"node:crypto";const h = (value: string): bigint =>BigInt(`0x${createHash("md5").update(value).digest("hex")}`);/** Index of the first element greater than x, in a sorted array. */functionbisectRight(sorted: bigint[], x: bigint): number {let lo = 0;let hi = sorted.length;while (lo < hi) {const mid = (lo + hi) >>> 1;if (x < sorted[mid]) hi = mid;else lo = mid + 1; }return lo;}exportclassHashRing {privatereadonly ring = newMap<bigint, string>();privatereadonly keys: bigint[] = [];constructor( nodes: string[],privatereadonly vnodes = 100, ) {for (const node of nodes) this.add(node); }add(node: string): void {for (let i = 0; i < this.vnodes; i++) {// many points per serverconst point = h(`${node}#${i}`);this.ring.set(point, node);this.keys.splice(bisectRight(this.keys, point), 0, point); } }remove(node: string): void {for (let i = 0; i < this.vnodes; i++) {const point = h(`${node}#${i}`);this.ring.delete(point);this.keys.splice(this.keys.indexOf(point), 1); } }nodeFor(key: string): string {const idx = bisectRight(this.keys, h(key)) % this.keys.length; // next point clockwisereturnthis.ring.get(this.keys[idx]) asstring; }}const ring = newHashRing(["cache-a", "cache-b", "cache-c"]);const keys = Array.from({ length: 10_000 }, (_, i) => `user:${i}`);const before = newMap(keys.map((k) => [k, ring.nodeFor(k)]));ring.add("cache-d");const moved = keys.filter((k) => before.get(k) !== ring.nodeFor(k)).length;console.log(`${(moved / 100).toFixed(1)}% of keys moved`); // about 25%, not 75%
Why it matters Adding a fourth node moves roughly a quarter of the keys, exactly the share the new node should own.
Rate limiters decide how many requests a client may make in a period. Token bucket, leaky bucket, fixed window and sliding window are the four algorithms you will be asked to compare.
In detail
Token bucket refills tokens at a steady rate and lets clients burst up to the bucket size; it is the most common choice (AWS, Stripe). Leaky bucket drains a queue at a fixed rate, smoothing output. Fixed window counters are simple but allow double bursts at window edges. Sliding window log is exact but stores every timestamp; sliding window counter blends two fixed windows for accuracy with tiny memory. In a cluster the counters live in Redis with atomic operations.
An LRU cache evicts the item used longest ago when it is full. A hash map finds items in O(1) and a doubly linked list keeps them in usage order.
In detail
On every get, move the item to the front; on every put beyond capacity, drop the item at the back. LFU evicts the least frequently used item instead, which keeps steady favourites but adapts slowly to new trends. Redis offers approximated LRU and LFU policies by sampling keys rather than keeping a perfect list. Choosing the eviction policy is a real design decision: LRU suits recency driven data like sessions, LFU suits stable popular content.
Cost of each operation
GetO(1)
PutO(1)
EvictO(1)
lru.tsTypeScript
this.data.delete(key); // mark as most recently usedthis.data.set(key, value);if (this.data.size > this.capacity) {const oldest = this.data.keys().next().value asstring;this.data.delete(oldest); // evict least recently used
// A Map keeps insertion order, so delete and set again moves a key to the end.classLRUCache {readonly data = newMap<string, unknown>(); hits = 0; misses = 0;constructor(readonly capacity: number) {}get(key: string): unknown {if (!this.data.has(key)) {this.misses += 1;returnundefined; }const value = this.data.get(key);this.data.delete(key); // mark as most recently usedthis.data.set(key, value);this.hits += 1;return value; }put(key: string, value: unknown): void {this.data.delete(key);this.data.set(key, value);if (this.data.size > this.capacity) {const oldest = this.data.keys().next().value asstring;this.data.delete(oldest); // evict least recently used } }gethitRatio(): number {const total = this.hits + this.misses;return total ? this.hits / total : 0; }}const cache = newLRUCache(2);cache.put("a", 1);cache.put("b", 2);cache.get("a"); // a becomes most recentcache.put("c", 3); // evicts bconsole.log([...cache.data.keys()]); // [ 'a', 'c' ]
Why it matters Hit ratio is the number to watch in production: a cache at 50% hits halves database load, at 99% it removes it.
A bloom filter answers is this item in the set with either definitely not or probably yes, using a tiny fraction of the memory a real set would need.
In detail
It sets k bits for each item using k hash functions; a lookup checks those bits. False positives happen, false negatives never do. Databases like Cassandra and RocksDB check a bloom filter before reading an SSTable from disk, crawlers use one to skip seen URLs, and CDNs use one to avoid caching one hit wonders. Its cousins answer other questions cheaply: HyperLogLog counts unique visitors in 12 KB, count min sketch estimates frequencies in a stream.
Cost of each operation
AddO(k)
CheckO(k)
Memory per itemabout 10 bits
bloom.tsTypeScript
for (let i = 0; i < this.k; i++) {this.bits[this.hash(item, i)] = 1;
import { createHash } from"node:crypto";exportclassBloomFilter {readonly m: number;readonly k: number;readonly bits: Uint8Array;constructor(expectedItems: number, falsePositiveRate = 0.01) {this.m = Math.ceil((-expectedItems * Math.log(falsePositiveRate)) / Math.log(2) ** 2);this.k = Math.max(1, Math.round((this.m / expectedItems) * Math.log(2)));this.bits = newUint8Array(this.m); }privatehash(item: string, i: number): number {const digest = createHash("sha256").update(`${i}:${item}`).digest();returnNumber(digest.readBigUInt64BE(0) % BigInt(this.m)); }add(item: string): void {for (let i = 0; i < this.k; i++) {this.bits[this.hash(item, i)] = 1; } }mightContain(item: string): boolean {for (let i = 0; i < this.k; i++) {if (!this.bits[this.hash(item, i)]) returnfalse; }returntrue; }}const seen = newBloomFilter(1_000_000, 0.01);console.log(Math.floor(seen.m / 8) / 1024 / 1024, "MiB for a million URLs"); // about 1.1 MiBseen.add("https://example.com/a");console.log(seen.mightContain("https://example.com/a")); // trueconsole.log(seen.mightContain("https://example.com/b")); // almost certainly false
Why it matters About 1.2 MB tracks a million URLs at 1% error; storing the URLs themselves would take around 100 MB.
A Merkle tree hashes data blocks, then hashes the hashes up to a single root. If two roots match, the data matches; if not, you can find the differing block in O(log n) comparisons.
In detail
Replicas in Cassandra and DynamoDB compare Merkle trees during anti entropy repair, exchanging only the branches that differ instead of whole datasets. Git commits, Bitcoin blocks, IPFS and certificate transparency logs all use hash trees. The idea is perfect whenever two machines must check large data is identical over a slow link.
Follow only the branch whose hash differsServiceEdgeDataCache
Geohash turns a latitude and longitude into a short string where nearby places share a prefix. Quadtrees split the map into four boxes again and again until each box holds a manageable number of places.
In detail
Finding restaurants near me cannot scan every restaurant. With geohash you compute the user's cell and its eight neighbours, then query places whose geohash starts with those prefixes. Quadtrees adapt to density: Manhattan splits into tiny cells, the desert stays one big cell. Uber built H3, a hexagonal grid, because hexagons have equal distance to every neighbour. Databases expose these as geospatial indexes (PostGIS, Redis GEO, Elasticsearch geo_point).
Cost of each operation
Encode a pointO(precision)
Prefix lookupO(log n)
Quadtree queryO(log n + k)
geohash.tsTypeScript
const rng = even ? lonRng : latRng; // interleave lon and lat bitsconst value = even ? lon : lat;const mid = (rng[0] + rng[1]) / 2;if (value >= mid) { ch = (ch << 1) | 1;
const BASE32 = "0123456789bcdefghjkmnpqrstuvwxyz";exportfunctiongeohash(lat: number, lon: number, precision = 7): string {const latRng: [number, number] = [-90, 90];const lonRng: [number, number] = [-180, 180];let even = true;let bitCount = 0;let ch = 0;let out = "";while (out.length < precision) {const rng = even ? lonRng : latRng; // interleave lon and lat bitsconst value = even ? lon : lat;const mid = (rng[0] + rng[1]) / 2;if (value >= mid) { ch = (ch << 1) | 1; rng[0] = mid; } else { ch = ch << 1; rng[1] = mid; } even = !even; bitCount += 1;if (bitCount === 5) { out += BASE32[ch]; bitCount = 0; ch = 0; } }return out;}// Two cafes a few hundred metres apart in Delhi share a long prefixconsole.log(geohash(28.6139, 77.209)); // ttnfucjconsole.log(geohash(28.615, 77.2105)); // ttnfucn, same first 6 characters// Cell size by precision: 5 chars ~ 4.9 km, 6 ~ 1.2 km, 7 ~ 153 m
Why it matters Prefix equals proximity, so a plain B-tree index on the geohash column answers nearby queries.
Top K finds the most frequent items in a stream, such as trending hashtags or most played songs, without keeping a counter for everything ever seen.
In detail
For exact results on bounded data, count with a hash map and keep a size K heap. For endless streams, approximate: a count min sketch estimates counts in fixed memory, and the heap keeps candidates. Systems compute top K per time window (last five minutes) on many machines, then merge partial results. The pattern appears in trending topics, top searched queries for autocomplete and top viewed videos for recommendations.
Distributed systems need IDs that are unique without asking a central database every time. UUIDs, Snowflake IDs and ticket servers are the usual answers.
In detail
Auto increment only works on one database. UUID v4 is random and needs no coordination but is 128 bits and sorts randomly, which fragments indexes; UUID v7 fixes sorting by putting time first. Twitter's Snowflake packs a timestamp, a machine id and a per millisecond sequence into 64 bits, so IDs are unique, compact and roughly time ordered. A ticket server hands out ranges of numbers to each app server to cut round trips.
const EPOCH_MS = 1704067200000; // 2024-01-01, keeps the timestamp field small// 41 bits time | 10 bits machine | 12 bits sequence = 4096 ids per ms per machine.// The id needs 63 bits, more than a number holds exactly, so it is a bigint.classSnowflake {privatereadonly machine: bigint;private lastMs = -1;private sequence = 0;constructor(machineId: number) {if (!(Number.isInteger(machineId) && machineId >= 0 && machineId < 1024)) {thrownewRangeError("machineId must fit in 10 bits"); }this.machine = BigInt(machineId); }// Synchronous, so no lock is needed: nothing else runs on this thread meanwhile.nextId(): bigint {let ms = Date.now();if (ms < this.lastMs) {thrownewError("clock moved backwards, refuse to issue ids"); }if (ms === this.lastMs) {this.sequence = (this.sequence + 1) & 0xfff;if (this.sequence === 0) {// 4096 used this ms, wait for the nextwhile (ms <= this.lastMs) { ms = Date.now(); } } } else {this.sequence = 0; }this.lastMs = ms;return (BigInt(ms - EPOCH_MS) << 22n) | (this.machine << 12n) | BigInt(this.sequence); }}const gen = newSnowflake(7);const a = gen.nextId();const b = gen.nextId();// true true: ordered and fits a signed 64 bit columnconsole.log(a < b, a.toString(2).length <= 63);
Why it matters Refusing to issue IDs when the clock goes backwards is what keeps Snowflake IDs unique after NTP corrections.
Base62 writes a number using 0 to 9, a to z and A to Z, so a huge database id becomes a short, URL safe code. Seven characters cover 3.5 trillion links.
In detail
A URL shortener can turn each new row id into base62 (unique by construction, no collision checks) or hash the long URL and take the first characters (needs collision handling). Sequential ids make codes guessable; mixing in a Feistel shuffle or a random offset hides the pattern. The same encoding shows up in YouTube video ids and invite codes.
Back of the envelope
Alphabet62 symbols0-9 a-z A-Z
6 characters56.8 billioncodes
7 characters3.5 trillioncodes
8 characters218 trillioncodes
base62.tsTypeScript
while (n > 0) { out.push(ALPHABET[n % 62]); n = Math.floor(n / 62);
const ALPHABET = "0123456789abcdefghijklmnopqrstuvwxyzABCDEFGHIJKLMNOPQRSTUVWXYZ";exportfunctionencode(n: number): string {if (n === 0) {return ALPHABET[0]; }const out: string[] = [];while (n > 0) { out.push(ALPHABET[n % 62]); n = Math.floor(n / 62); }return out.reverse().join("");}exportfunctiondecode(code: string): number {let n = 0;for (const ch of code) { n = n * 62 + ALPHABET.indexOf(ch); }return n;}console.log(62 ** 7); // 3521614606208 possible 7 character codesconsole.log(encode(125_000_000_000)); // 2crtgcg: 7 charactersconsole.log(decode(encode(987654321)) === 987654321); // true
Why it matters Encoding an id needs no lookup and can never collide, which is why most shorteners prefer it to hashing.
03
Phase 03, modules 21 to 30
Thinking like an architect
System design fundamentals
Requirements, estimates, scale, availability and the trade-offs every design is really about.
What system design is, and the interview framework
System design is deciding which parts a product needs, how they talk, where data lives and how it all keeps working as users grow. Interviews test whether you can reason through those choices out loud.
In detail
A strong answer follows a rhythm: clarify requirements, estimate scale, sketch a high level design, define the API and data model, then deep dive into the hardest parts and discuss trade-offs and failure modes. There is rarely one right answer; there are justified choices. Interviewers listen for numbers, bottlenecks and what you would monitor, not for memorised diagrams.
1. Requirements 2. Estimates 3. High level design4. API + data model 5. Deep dives 6. Bottlenecks and trade-offs
# 45 minute system design interview, timed0-5 Clarify: who uses it, core features, what is out of scope5-10 Non functional: scale, latency target, availability, consistency10-15 Estimate: users, QPS (average and peak), storage per year, bandwidth15-25 High level: clients, load balancer, services, data stores, caches, queues25-30 API and data model: endpoints, request and response, tables or documents30-40 Deep dive: the hardest part (hot keys, fan out, consistency, ordering)40-45 Wrap up: bottlenecks, failure modes, monitoring, what you would do next
Why it matters Saying the plan aloud in the first minute shows structure and lets the interviewer steer you.
Functional requirements are what the system does: shorten a URL, send a message. Non functional requirements are how well it does it: how fast, how available, how consistent, at what scale.
In detail
Write both lists before drawing anything. Non functional requirements drive the architecture far more than features do: 100 ms p99 latency forces caching, 99.99% availability forces redundancy across zones, strong consistency rules out some databases. Also state what is out of scope so the interview stays focused. Read heavy versus write heavy is the single most useful classification.
requirements.mdNotes
Functional: create short link, redirect, custom alias, expiryNon functional: 100M links/month, p99 redirect < 50 ms, 99.99% available
## URL shortener: requirements### Functional- Create a short link for a long URL (optional custom alias)- Redirect a short link to its long URL- Links can expire; owners can delete them- Basic click analytics per link### Non functional- Scale: 100M new links per month, 10B redirects per month (100:1 read heavy)- Latency: redirect p99 under 50 ms worldwide- Availability: 99.99% for redirects; link creation can be 99.9%- Durability: a created link must never be lost- Consistency: a new link may take a second to resolve everywhere (eventual is fine)### Out of scope- User accounts beyond an API key, link previews, spam scanning
Why it matters The 100 to 1 read ratio alone tells you to cache aggressively and keep redirects off the write path.
Rough numbers for users, requests per second, storage and bandwidth tell you whether one database will do or whether you need a fleet. Round aggressively and show your working.
In detail
Start from daily active users and actions per user, divide by 86,400 seconds (call it 100,000) for average QPS, multiply by 2 to 5 for peak. Storage is items per day times size times retention. Bandwidth is QPS times payload size. Then compare with what one machine handles: a well tuned PostgreSQL does thousands of simple writes per second, Redis around 100,000 operations per second, a single server tens of thousands of HTTP requests.
Reading memory takes nanoseconds, an SSD takes microseconds, a network round trip across a data centre takes half a millisecond, and crossing an ocean takes over a hundred. Designs follow these gaps.
In detail
Latency is time for one request; throughput is requests per second. Measure percentiles, not averages: p99 is what your unluckiest regular users feel, and with fan out the slowest of many calls decides the total. Keep hot data in memory, avoid sequential network calls on the request path, and move far away users closer with CDNs and regional replicas.
Latency numbers, log scaleeach step to the right is about ten times slower
Vertical scaling buys a bigger server; horizontal scaling adds more servers behind a load balancer. The first is simple but capped, the second is unlimited but needs stateless services and partitioned data.
In detail
Scale up first: it is cheaper in engineering time than you think. When you scale out, keep application servers stateless (sessions in Redis, files in object storage) so any server can handle any request. Data is the hard part: read replicas scale reads, caching removes reads, sharding scales writes. Every extra machine adds coordination, so scale the part that is actually the bottleneck.
Vertical, scale up
64 cores, 512 GB
Buy a bigger machine. Simple, no code changes, but there is a ceiling and it is still one point of failure.
Horizontal, scale out
Add more ordinary machines behind a load balancer. Near unlimited, survives failures, but needs stateless services and sharded data.
stateless.tsTypeScript
// Bad: session in process memory, breaks with 2+ servers// SESSIONS.set(token, userId);// Good: stored outside the process, so any app server behind the load balancer can read itawait r.setex(`session:${token}`, 3600, JSON.stringify({ userId } satisfiesSession));
import { Redis } from"ioredis";const r = newRedis({ host: "sessions.internal", port: 6379 });typeSession = { userId: string };exportasyncfunctionlogin(userId: string, token: string): Promise<void> {// Bad: session in process memory, breaks with 2+ servers// SESSIONS.set(token, userId);// Good: stored outside the process, so any app server behind the load balancer can read itawait r.setex(`session:${token}`, 3600, JSON.stringify({ userId } satisfiesSession));}exportasyncfunctioncurrentUser(token: string): Promise<string | null> {const raw = await r.get(`session:${token}`);return raw ? (JSON.parse(raw) asSession).userId : null;}// Files: never write uploads to the local disk of one server// const Key = `avatars/${userId}.webp`;// await s3.send(new PutObjectCommand({ Bucket: "uploads", Key, Body: data }));
Why it matters Once state lives in shared stores, adding a server is just starting another identical process.
Availability is the share of time a system works. 99.9% allows about 43 minutes of downtime a month, 99.99% about 4 minutes. Every extra nine costs roughly ten times more.
In detail
Components in series multiply availability (two 99.9% services in a chain give 99.8%); redundant components in parallel improve it. Remove single points of failure with replicas across availability zones, health checks and automatic failover. Define service level indicators (what you measure), objectives (the target) and agreements (the promise with penalties). An error budget, the allowed unreliability, lets teams trade speed against stability with data.
What each extra nine of availability allows
Availability
Down per year
Down per month
Down per week
Typical for
99%
3.65 days
7.3 hours
1.7 hours
Internal tools
99.9%
8.8 hours
43.8 minutes
10.1 minutes
Most web apps
99.95%
4.4 hours
21.9 minutes
5 minutes
Business critical APIs
99.99%
52.6 minutes
4.4 minutes
1 minute
Payments, core platforms
99.999%
5.3 minutes
26 seconds
6 seconds
Telecom, cloud control planes
nines.tsTypeScript
for (const p of parts) total *= p;return1 - (1 - part) ** copies;console.log(`chain: ${pct(serial(0.999, 0.999, 0.999), 2)}`); // 99.70%console.log(`redundant: ${pct(serial(...tiers), 4)}`); // 99.9997%
const MONTH_MIN = 30 * 24 * 60;const pct = (x: number, digits: number): string => `${(x * 100).toFixed(digits)}%`;functiondowntimePerMonth(availability: number): number {return (1 - availability) * MONTH_MIN;}for (const nines of [0.99, 0.999, 0.9999, 0.99999]) {const minutes = downtimePerMonth(nines).toFixed(1).padStart(7); console.log(`${pct(nines, 5)} -> ${minutes} minutes per month`);}functionserial(...parts: number[]): number {let total = 1;for (const p of parts) total *= p;return total;}functionparallel(part: number, copies: number): number {return1 - (1 - part) ** copies;}// LB -> app -> DB, each 99.9%console.log(`chain: ${pct(serial(0.999, 0.999, 0.999), 2)}`); // 99.70%// Same chain with every tier doubled across two zonesconst tiers = Array.from({ length: 3 }, () => parallel(0.999, 2));console.log(`redundant: ${pct(serial(...tiers), 4)}`); // 99.9997%
Why it matters Redundancy at every tier turns three 99.9% parts into a system near five nines, as long as failover actually works.
Fault tolerance means the system keeps serving when parts fail. Disks die, networks split and deploys go wrong, so good designs assume failure and plan the response.
In detail
Use redundancy (more than one of everything), isolation (a failure in one feature should not take down another), graceful degradation (show cached results when the recommender is down), and fast detection with automatic recovery. Chaos engineering, popularised by Netflix's Chaos Monkey, breaks things on purpose in production to prove the recovery works. Replication protects against machine loss; backups protect against mistakes such as a bad migration.
When a network partition cuts replicas apart, a distributed system must choose: keep answering with possibly stale data (availability) or refuse until it can be sure (consistency).
In detail
Partitions are not optional in real networks, so the practical choice is CP or AP during a partition. PACELC adds the everyday trade-off: else, when there is no partition, you trade latency against consistency. Banks and inventory counters lean CP; social feeds, shopping carts and DNS lean AP. Many databases let you choose per query with tunable quorum reads and writes.
CAPCA: single nodeCP: etcd, Spanner, HBaseAP: Cassandra, DynamoDB, CouchDB
Networks always partition eventually, so in practice the choice is between C and A while the partition lasts. PACELC adds: else, choose between latency and consistency.
quorum.tsTypeScript
const overlap = r + w > n; // read and write sets must overlap
/** Dynamo style quorums: N replicas, wait for W writes and R reads. */functionquorumProfile(n: number, w: number, r: number): string {const overlap = r + w > n; // read and write sets must overlapconst survivesW = n - w; // replicas that may be down and writes still succeedconst survivesR = n - r;const kind = overlap ? "strong reads (latest write visible)" : "eventual reads (may be stale)";return (`N=${n} W=${w} R=${r}: ${kind}; ` +`writes survive ${survivesW} down, reads survive ${survivesR} down` );}console.log(quorumProfile(3, 2, 2)); // strong, tolerates 1 down either wayconsole.log(quorumProfile(3, 1, 1)); // fast and available, but may read stale dataconsole.log(quorumProfile(3, 3, 1)); // fast reads, but one node down blocks all writes
Why it matters Tuning R and W is CAP made concrete: overlap buys consistency, small quorums buy availability and latency.
A consistency model promises what a reader may see after a write. Strong consistency shows the latest value everywhere at once; eventual consistency promises replicas will agree soon.
In detail
Between the extremes sit useful middle grounds: read your own writes (you always see your own post), monotonic reads (you never go back in time), causal consistency (replies appear after the message they answer). Linearizability is the strongest single object guarantee and needs consensus or a single leader. Pick the weakest model the product tolerates, because each step up costs latency and availability.
System design is mostly trade-offs: latency against consistency, cost against availability, simplicity against flexibility. Naming the trade-off out loud is half of a good answer.
In detail
Common pairs: SQL versus NoSQL, push versus pull, cache everything versus always fresh, monolith versus microservices, synchronous versus asynchronous, normalised versus denormalised, strong versus eventual consistency. For each choice say what you gain, what you give up, and under which numbers you would switch. The table below is a quick reference you can return to before any interview.
DNS turns a name like netflix.com into an IP address. It is hierarchical, heavily cached and also a simple traffic tool: it can send users to the nearest or healthiest region.
In detail
A resolver asks the root servers, then the .com servers, then the domain's authoritative servers, caching each answer for its time to live. Low TTLs make changes fast but increase lookups. GeoDNS and latency based routing answer differently per user location, and health checked records drop a failed region. Because DNS caching is outside your control, it is a slow failover tool; anycast and load balancers react faster.
A cold lookup walks three levels, then cachesClientCacheExternalService
# Resolve a name (answers come from your resolver's cache when possible)dig +short netflix.com# Watch the full walk: root -> .com -> authoritative serversdig netflix.com +trace# See the TTL that controls caching (second column)dig netflix.com A +noall +answer# Different record typesdig netflix.com AAAA # IPv6dig netflix.com MX # mail serversdig www.netflix.com CNAME # alias to another name
Why it matters The TTL column tells you how long a change could take to reach everyone.
Clients send requests and servers send responses. HTTP rides on TCP for reliable, ordered delivery; UDP drops those guarantees for speed; QUIC brings reliability back on top of UDP for HTTP/3.
In detail
A TCP connection needs a handshake, and TLS adds more round trips, so reusing connections (keep alive, pooling) matters. HTTP/2 multiplexes many requests over one connection; HTTP/3 over QUIC removes head of line blocking and survives network changes on phones. Use UDP when late data is useless: video calls, games, DNS. Know status code families: 2xx success, 3xx redirect, 4xx client error, 5xx server error.
Transport and HTTP versions
Protocol
Built on
Strength
Weakness
Typical use
TCP
IP
Reliable, ordered
Handshake cost, head of line blocking
Almost everything
UDP
IP
No setup, low overhead
No delivery guarantee
Video calls, games, DNS
HTTP/1.1
TCP
Simple, universal
One request at a time per connection
Legacy APIs
HTTP/2
TCP + TLS
Multiplexed streams, header compression
TCP head of line blocking
Modern web, gRPC
HTTP/3
QUIC over UDP
No head of line blocking, fast setup
Newer, some networks block UDP
Mobile, lossy networks
http.shShell
curl -sI https://example.com | head -5curl --http3 -sI https://cloudflare.com
# Status line and headers onlycurl -sI https://example.com# Timing breakdown: DNS, TCP connect, TLS, first byte, totalcurl -s -o /dev/null https://example.com -w \"dns %{time_namelookup}s tcp %{time_connect}s tls %{time_appconnect}s ttfb %{time_starttransfer}s total %{time_total}s\n"# Which HTTP version was negotiatedcurl -s -o /dev/null -w"%{http_version}\n" https://www.google.com
Why it matters The timing line shows how much of a request is setup; that is the cost connection reuse saves.
A CDN keeps copies of static files, and sometimes whole pages, on servers near users. Images, video and scripts load from a nearby city instead of crossing an ocean.
In detail
Pull CDNs fetch from your origin on the first miss and cache by TTL; push CDNs receive content ahead of time, which suits large known files such as video. Cache keys, Cache-Control headers and versioned filenames decide what stays fresh. Netflix built Open Connect, appliances inside internet providers that serve most of its traffic from within the user's own ISP. CDNs also absorb DDoS attacks and terminate TLS close to users.
Hits stay at the edge, misses go to origin onceClientCacheService
# Fingerprinted assets: the name changes when the content changes, cache foreverGET /assets/app.3f9a2c.jsCache-Control: public, max-age=31536000, immutable# HTML entry point: always revalidate so new deploys show upGET /index.htmlCache-Control: no-cacheETag: "v2025-10-06-1"# Personalised API responses: never on a shared cacheGET /api/meCache-Control: private, no-store# Let the CDN serve stale briefly while it refreshes in the backgroundGET /api/trendingCache-Control: public, s-maxage=30, stale-while-revalidate=60
Why it matters Versioned filenames plus long TTLs give near 100% CDN hit rates without ever serving outdated code.
A load balancer spreads requests across healthy servers, hides failures and lets you add capacity without clients noticing. It works at layer 4 (TCP) or layer 7 (HTTP).
In detail
Layer 4 balancers route by IP and port and are extremely fast; layer 7 balancers read HTTP, so they can route by path, header or cookie and terminate TLS. Algorithms include round robin, weighted round robin, least connections and consistent hashing for stickiness. Health checks remove broken instances. Run balancers in pairs or as a managed service so the balancer itself is not a single point of failure.
Requests spread across healthy instancesClientEdgeService
A reverse proxy sits in front of servers; an API gateway is a reverse proxy that also handles authentication, rate limits, routing to microservices, request shaping and analytics.
In detail
Clients see one host while the gateway routes /orders to the order service and /users to the user service. Centralising cross cutting concerns keeps services simple, but the gateway becomes critical infrastructure that must scale and stay thin. The backend for frontend pattern gives mobile and web their own gateway tuned to their needs. Examples include Kong, Envoy, NGINX, AWS API Gateway and Apigee.
REST exposes resources over HTTP verbs, GraphQL lets clients ask for exactly the fields they need, and gRPC sends compact binary messages over HTTP/2 for fast service to service calls.
In detail
REST is simple, cacheable and universal, so it suits public APIs. GraphQL removes over fetching and many round trips for rich clients but makes caching and rate limiting harder. gRPC gives typed contracts from Protocol Buffers, streaming and low latency between internal services. Whatever the style, design idempotent writes, cursor pagination and versioning from day one.
REST, GraphQL and gRPC
Style
Shape
Strength
Weakness
Best for
REST
Resources and HTTP verbs, JSON
Simple, cacheable, universal
Over and under fetching
Public APIs
GraphQL
One endpoint, client written queries
Exactly the fields needed
Caching and cost control are harder
Rich frontends, many clients
gRPC
Protobuf over HTTP/2
Fast, typed, streaming
Not browser native
Service to service
WebSocket
Bidirectional messages
Real time both ways
Stateful connections to scale
Chat, live updates
api.httpNotes
GET /v1/links?cursor=abc&limit=20POST /v1/links Idempotency-Key: 7f3c...
### Create a short link (idempotent: retrying with the same key returns the same link)POST /v1/linksContent-Type: application/jsonIdempotency-Key: 7f3c2a8e-4b1d-4c55-9a1e-2f0d1b6c9e10{ "long_url": "https://example.com/a/very/long/path", "custom_alias": null, "expires_at": "2026-12-31T00:00:00Z" }### 201 Created{ "code": "aZ3kQ9x", "short_url": "https://sho.rt/aZ3kQ9x", "long_url": "https://example.com/a/very/long/path" }### Cursor pagination: stable under inserts, unlike offset paginationGET /v1/links?limit=20&cursor=eyJpZCI6IDEyMzR9### 200 OK{ "items": [ ... ], "next_cursor": "eyJpZCI6IDEyNTR9" }### RedirectGET /aZ3kQ9x### 301 Moved Permanently (or 302 if every click must be counted by the server)Location: https://example.com/a/very/long/path
Why it matters Idempotency keys and cursors are the two API habits that save the most production incidents.
To push updates to clients you can poll, long poll, stream with server sent events, or hold a two way WebSocket. Chat, live scores and collaborative editing all depend on this choice.
In detail
Short polling is simple but wasteful. Long polling holds the request until there is news. SSE streams one way from server to client over plain HTTP with automatic reconnects, ideal for feeds and notifications. WebSockets give full duplex messaging for chat and games. Persistent connections are stateful, so you need a connection registry and a pub/sub layer to route a message to whichever server holds the recipient's socket.
Relational databases store tables with strict schemas, joins and transactions. NoSQL databases trade some of that for flexible models and easier horizontal scale.
In detail
Choose SQL (PostgreSQL, MySQL) by default when data is relational and correctness matters: payments, inventory, bookings. Choose a key value store (Redis, DynamoDB) for simple lookups at huge scale, a document store (MongoDB) for nested records read together, a wide column store (Cassandra) for massive write throughput and time ordered data, and a graph database (Neo4j) for deep relationship queries. Many systems use several, each for what it does best.
Picking a database
Kind
Examples
Great at
Weak at
Relational
PostgreSQL, MySQL
Transactions, joins, constraints
Horizontal write scaling
Key value
Redis, DynamoDB
Fast lookups by key, huge scale
Queries on other fields
Document
MongoDB, Firestore
Flexible nested records
Cross document transactions and joins
Wide column
Cassandra, ScyllaDB, HBase
Massive write throughput, time series
Ad hoc queries
Graph
Neo4j, Neptune
Relationship traversal
Bulk analytics on everything
Search
Elasticsearch, OpenSearch
Full text and faceted search
Being the source of truth
Time series
Prometheus, InfluxDB, TimescaleDB
Metrics, compression, rollups
General purpose data
schema.sqlSQL
CREATETABLE links ( code TEXT PRIMARYKEY, long_url TEXT NOTNULL, owner_id BIGINT, created_at TIMESTAMPTZ DEFAULTnow());
-- Relational: strict schema, constraints, joins, transactionsCREATETABLE users ( id BIGINT PRIMARYKEY, email TEXT UNIQUE NOTNULL, created_at TIMESTAMPTZ NOTNULLDEFAULTnow());CREATETABLE links ( code TEXT PRIMARYKEY, -- base62 short code long_url TEXT NOTNULL, owner_id BIGINT REFERENCES users(id), expires_at TIMESTAMPTZ, created_at TIMESTAMPTZ NOTNULLDEFAULTnow());CREATEINDEX links_owner_created ON links (owner_id, created_at DESC);-- The same link in a key value store (DynamoDB style): one item, one key-- PK = "LINK#aZ3kQ9x"-- { "long_url": "...", "owner_id": 42, "expires_at": 1798675200, "ttl": 1798675200 }
Why it matters The redirect path only ever looks up by code, which is why a key value store is a strong fit for that one table.
An index is a sorted structure that points to rows, turning a full table scan into a few page reads. B-trees favour reads; LSM trees favour heavy writes.
In detail
Index the columns you filter, join and sort by, in the order your queries use them (composite index on owner_id then created_at serves where owner_id = ? order by created_at). Every index slows writes and uses space. LSM trees (Cassandra, RocksDB, ScyllaDB) write sequentially to a memtable and flush sorted files, compacting them later; bloom filters avoid checking files that cannot contain a key. EXPLAIN shows whether a query actually uses the index.
B-tree and LSM tree storage engines
Aspect
B-tree
LSM tree
Writes
Update pages in place
Append to memtable, flush sorted files
Reads
One path down the tree
May check several files, Bloom filters help
Write throughput
Moderate
Very high
Space and compaction
Some fragmentation
Background compaction merges files
Used by
PostgreSQL, MySQL InnoDB
Cassandra, RocksDB, LevelDB, ScyllaDB
Cost of each operation
B-tree lookupO(log n)
Full scanO(n)
LSM writeO(1) append
LSM readO(log n) per level
explain.sqlSQL
CREATEINDEX links_owner_created ON links (owner_id, created_at DESC);EXPLAIN ANALYZE SELECT * FROM links WHERE owner_id = 42ORDERBY created_at DESCLIMIT20;
-- Before: sequential scan over every rowEXPLAIN ANALYZESELECT code, long_url FROM links WHERE owner_id = 42ORDERBY created_at DESCLIMIT20;-- Seq Scan on links (rows=50000000) Execution Time: 4100 ms-- Composite index matching the filter, then the sortCREATEINDEX CONCURRENTLY links_owner_created ON links (owner_id, created_at DESC);-- After: index scan reads only the 20 rows it needsEXPLAIN ANALYZESELECT code, long_url FROM links WHERE owner_id = 42ORDERBY created_at DESCLIMIT20;-- Index Scan using links_owner_created (rows=20) Execution Time: 0.3 ms-- Covering index: answer from the index alone, no table visitCREATEINDEX links_owner_cover ON links (owner_id, created_at DESC) INCLUDE (code, long_url);
Why it matters CONCURRENTLY builds the index without locking writes, which matters on a live production table.
A transaction groups operations so they all succeed or all fail. ACID promises atomicity, consistency, isolation and durability, and isolation levels decide how concurrent transactions see each other.
In detail
Read committed (PostgreSQL's default) prevents dirty reads; repeatable read gives each transaction a stable snapshot; serializable makes concurrent transactions behave as if run one by one, at the cost of retries. Lost updates and double booking come from read then write races: fix them with SELECT FOR UPDATE, atomic updates, unique constraints or optimistic version checks. Keep transactions short; long ones hold locks and bloat storage.
Isolation levels and the anomalies they allow
Level
Dirty read
Non repeatable read
Phantom read
Lost update
Read uncommitted
Possible
Possible
Possible
Possible
Read committed
Prevented
Possible
Possible
Possible
Repeatable read
Prevented
Prevented
Depends on engine
Depends on engine
Snapshot
Prevented
Prevented
Prevented
Write skew possible
Serializable
Prevented
Prevented
Prevented
Prevented
book_seat.sqlSQL
UPDATE seats SET status = 'held', held_by = $1WHERE id = $2AND status = 'free'; -- 0 rows = someone else got it
BEGIN;-- Optimistic and atomic: only succeeds if the seat is still freeUPDATE seatsSET status = 'held', held_by = $user_id, held_until = now() + interval'10 minutes'WHERE id = $seat_id AND status = 'free';-- 1 row updated: we hold it. 0 rows: someone else got there first, tell the user.INSERTINTO bookings (seat_id, user_id, status) VALUES ($seat_id, $user_id, 'pending');COMMIT;-- Pessimistic alternative: lock the row while decidingBEGIN;SELECT status FROM seats WHERE id = $seat_id FORUPDATE; -- others wait here-- ... check status, then update ...COMMIT;-- Constraint as a safety net: impossible to book one seat twiceALTER TABLE bookings ADD CONSTRAINT one_booking_per_seat UNIQUE (seat_id, event_id);
Why it matters The conditional UPDATE is a compare and set in SQL: the database decides the race, not your application code.
Replication keeps copies of data on several machines for availability, durability and read scale. Leader follower, multi leader and leaderless are the three shapes.
In detail
With a single leader, writes go to the leader and stream to followers; reads can go to followers but may lag. Synchronous replication waits for a follower so no write is lost on failover; asynchronous is faster but can lose recent writes. Multi leader setups accept writes in several regions and must resolve conflicts. Leaderless systems (Dynamo, Cassandra) write to several nodes with quorums and repair differences with Merkle trees and read repair.
Writes go to the leader, reads can fan outServiceDataCache
SELECT client_addr, state, replay_lag FROM pg_stat_replication;
-- On the primary: who is following and how far behindSELECT client_addr, state, sync_state, write_lag, flush_lag, replay_lagFROM pg_stat_replication;-- Make one follower synchronous: commits wait until it has the WALALTER SYSTEM SET synchronous_standby_names = 'FIRST 1 (replica_a, replica_b)';SELECT pg_reload_conf();-- On a follower: how stale am I right nowSELECTnow() - pg_last_xact_replay_timestamp() AS replica_delay;
Why it matters Alert on replay_lag: a lagging follower silently serves stale reads long before it causes an outage.
Sharding splits one big dataset across many databases by a shard key, so writes and storage scale beyond a single machine.
In detail
Range sharding keeps nearby keys together (good for range scans, risky for hot spots); hash sharding spreads load evenly; directory based sharding keeps a lookup table for full control. The shard key is the hardest decision: it must spread load and keep related data together, because cross shard joins and transactions are slow. Plan for hot keys (a celebrity's account), resharding with consistent hashing, and per shard replication.
The shard key decides where each row livesEdgeData
user NehaRouterEdgeShard 0: A-FDataShard 1: G-MDataShard 2: N-SDataShard 3: T-ZData
Sharding strategies
Strategy
How keys map
Strength
Weakness
Range
Key ranges per shard, A to F, G to M
Range scans stay local
Hot spots on recent or popular ranges
Hash
hash(key) mod shards
Even spread
Range queries hit every shard
Consistent hash
Hash ring with virtual nodes
Little data moves on resize
More moving parts
Directory
Lookup table key to shard
Fully flexible moves
The directory is critical infrastructure
Geographic
By region or tenant
Data near users, compliance
Uneven region sizes
router.tsTypeScript
return SHARDS[digest.readUInt32BE(0) % SHARDS.length];const db = connect(shardFor(userId));await db.query("INSERT INTO posts (user_id, text) VALUES ($1, $2)", [userId, text]);
import { createHash } from"node:crypto";import pg from"pg";const SHARDS = ["pg-shard-0", "pg-shard-1", "pg-shard-2", "pg-shard-3"];const pools = newMap<string, pg.Pool>();typePost = { id: number; user_id: number; text: string };/** One connection pool per shard host, opened on first use. */functionconnect(host: string): pg.Pool {let pool = pools.get(host);if (!pool) { pool = new pg.Pool({ host, database: "app" }); pools.set(host, pool); }return pool;}/** Hash shard by user: all of one user's data lives together. */exportfunctionshardFor(userId: number): string {const digest = createHash("sha1").update(String(userId)).digest();return SHARDS[digest.readUInt32BE(0) % SHARDS.length];}exportasyncfunctioncreatePost(userId: number, text: string): Promise<void> {const db = connect(shardFor(userId));await db.query("INSERT INTO posts (user_id, text) VALUES ($1, $2)", [userId, text]);}exportasyncfunctionuserTimeline(userId: number): Promise<Post[]> {// Single shard query: fastconst { rows } = awaitconnect(shardFor(userId)).query<Post>("SELECT * FROM posts WHERE user_id = $1 ORDER BY id DESC LIMIT 50", [userId], );return rows;}exportasyncfunctionglobalSearch(term: string): Promise<Post[]> {// Scatter gather across every shard: slow, avoid on hot pathsconst results = awaitPromise.all( SHARDS.map((s) =>connect(s).query<Post>("SELECT * FROM posts WHERE text ILIKE $1 LIMIT 10", [`%${term}%`]), ), );return results.flatMap((res) => res.rows);}
Why it matters Modulo by shard count makes adding a shard move most rows; production routers use consistent hashing or a directory.
NoSQL is four families: key value stores for lookups by key, document stores for nested records, wide column stores for huge write heavy tables, and graph stores for relationships.
In detail
Model NoSQL data around access patterns, not entities: list the queries first, then design keys so each query reads one partition. Cassandra wants a partition key plus clustering columns (messages by conversation, ordered by time). DynamoDB single table design stores several entity types under shared keys. Graph databases shine for recommendations and fraud rings where queries hop many relationships.
-- Cassandra: design the table for the query "latest messages in a conversation"CREATETABLE messages_by_conversation ( conversation_id uuid, sent_at timeuuid, message_id uuid, sender_id uuid, body text,PRIMARYKEY ((conversation_id), sent_at, message_id)) WITHCLUSTERINGORDERBY (sent_at DESC);-- One partition read, already sorted newest firstSELECT sender_id, body FROM messages_by_conversationWHERE conversation_id = ? LIMIT50;-- A second table for a second query, data duplicated on purposeCREATETABLE conversations_by_user ( user_id uuid, last_message_at timeuuid, conversation_id uuid, title text,PRIMARYKEY ((user_id), last_message_at, conversation_id)) WITHCLUSTERINGORDERBY (last_message_at DESC);
Why it matters In NoSQL one query maps to one table; duplication is the price of single partition reads.
Object storage such as Amazon S3 keeps files of any size behind a key, with eleven nines of durability and no servers to manage. Videos, images, backups and logs live here, not in databases.
In detail
Databases store metadata (owner, size, key) while bytes go to object storage. Clients upload directly with presigned URLs so files never pass through app servers. Large files use multipart uploads, which also enables resume. Lifecycle rules move old data to cheaper tiers. Object storage plus a CDN in front is the standard way to serve user generated media.
Full text search uses an inverted index: a map from each word to the documents that contain it. Elasticsearch, OpenSearch and Lucene power search bars, logs and product catalogs.
In detail
Text is tokenised, lower cased and stemmed, then each term points to a posting list of document ids. Queries intersect posting lists and rank results with BM25, boosts and freshness. The index is a secondary, eventually consistent copy fed from the main database through change data capture or events, and it is sharded and replicated like any data store.
inverted-index.tsTypeScript
for (const [docId, text] of docs) {for (const term oftokenize(text)) { index.get(term)!.add(docId);
const STOP = newSet(["the", "a", "and", "of", "to", "in"]);functiontokenize(text: string): string[] {return (text.toLowerCase().match(/[a-z0-9]+/g) ?? []).filter((t) => !STOP.has(t));}const docs = newMap<number, string>([ [1, "Designing a URL shortener like Bitly"], [2, "Designing Netflix video streaming at scale"], [3, "Scale a chat app like WhatsApp"],]);const index = newMap<string, Set<number>>();for (const [docId, text] of docs) {for (const term oftokenize(text)) {if (!index.has(term)) index.set(term, newSet()); index.get(term)!.add(docId); }}exportfunctionsearch(query: string): number[] {const terms = tokenize(query);if (terms.length === 0) {return []; }const lists = terms.map((t) => index.get(t) ?? newSet<number>());const hits = lists.reduce((acc, s) => newSet([...acc].filter((id) => s.has(id)))); // AND queryreturn [...hits].sort((a, b) => a - b);}console.log(search("designing scale")); // [ 2 ]
Why it matters Posting list intersection is the core of every search engine; ranking and sharding are layered on top.
06
Phase 06, modules 46 to 52
Remember and relay
Caching and messaging
Keeping hot data close, and letting services talk without waiting on each other.
A cache stores a copy of data somewhere faster than its source. Caches sit at every layer: the browser, the CDN, the API server's memory, a shared Redis cluster and the database's own buffer pool.
In detail
Each layer trades freshness for speed. Browser and CDN caches remove whole requests; in process caches avoid a network hop but are per instance; a shared cache such as Redis or Memcached is consistent across instances and survives deploys. Measure hit ratio: at 95% hits, the database only sees one request in twenty. Hot keys, cold starts after a flush and memory limits are the classic problems.
Each layer only forwards its missesClientEdgeServiceCacheData
missmiss12BrowserClientCDNEdgeApp + local cacheServiceRedisCacheDatabaseData
Back of the envelope
L1 hitabout 1 nsCPU cache
Local memoryabout 100 nsin process
Redis round tripabout 0.5 mssame zone
Database read1 to 10 msindexed
layers.tsTypeScript
if (hit && hit.expiresAt > Date.now()) return hit.value; // ~100 nslet value = awaitthis.redis.get(key); // ~0.5 msvalue = awaitthis.db.load(key); // ~5 ms
importtype { Redis } from"ioredis";interfaceDatabase {load(key: string): Promise<string>;}/** In process Map in front of a shared Redis, in front of the database. */classTwoLevelCache {privatereadonly local = newMap<string, { value: string; expiresAt: number }>();privatereadonly redis: Redis;privatereadonly db: Database;privatereadonly localTtl: number;privatereadonly sharedTtl: number;constructor(redis: Redis, db: Database, localTtl = 5, sharedTtl = 300) {this.redis = redis;this.db = db;this.localTtl = localTtl;this.sharedTtl = sharedTtl; }asyncget(key: string): Promise<string> {const hit = this.local.get(key);if (hit && hit.expiresAt > Date.now()) return hit.value; // ~100 nslet value = awaitthis.redis.get(key); // ~0.5 msif (value === null) { value = awaitthis.db.load(key); // ~5 msawaitthis.redis.set(key, value, "EX", this.sharedTtl); }this.local.set(key, { value, expiresAt: Date.now() + this.localTtl * 1000 });return value; }}
Why it matters A short local TTL absorbs hot keys without letting instances drift apart for long.
Cache aside, read through, write through, write back and write around decide who loads and updates the cache. Invalidation and expiry decide how stale data may get.
In detail
Cache aside is the default: the app reads the cache, falls back to the database and fills the cache. Write through updates cache and database together; write back writes to the cache and flushes later, fast but risky. Expire with TTLs plus jitter, delete keys on writes, and protect against stampedes with request coalescing or a short lock. Phil Karlton's joke stands: cache invalidation is one of the two hard things.
A message queue lets a producer hand off work and move on while consumers process it at their own pace. It absorbs spikes, retries failures and decouples services.
In detail
Producers publish messages; the broker stores them durably; consumers pull, process and acknowledge. Unacknowledged messages are redelivered, so consumers must be idempotent. Messages that keep failing go to a dead letter queue. Delivery is usually at least once; exactly once is approximated with idempotency keys. RabbitMQ, Amazon SQS and Kafka are common choices, Kafka being a log rather than a classic queue.
Work waits safely until a worker is freeServiceQueueExternal
In publish subscribe, a message goes to every interested subscriber, not just one worker. A log based system such as Kafka keeps messages in an ordered, replayable log split into partitions.
In detail
Topics fan a message out to many independent consumer groups: billing, email and analytics each get every order event. Kafka partitions a topic by key, keeps order within a partition, stores messages for days and lets each consumer group track its own offset, so a new service can replay history. Throughput comes from partitions; ordering only holds within one.
One event, every subscriber gets a copyServiceQueue
Event driven architecture, CQRS and event sourcing
In an event driven system, services announce facts such as OrderPlaced and others react. Event sourcing stores those facts as the source of truth; CQRS keeps separate models for writing and reading.
In detail
Events decouple teams: the order service does not know who listens. The outbox pattern writes the event in the same database transaction as the state change, then a relay publishes it, so you never lose or invent events. Event sourcing rebuilds state by replaying events, giving a full audit trail at the cost of complexity. Choreography lets services react on their own; orchestration has one coordinator drive the flow.
-- One transaction: state change and event are saved together or not at allBEGIN;UPDATE orders SET status = 'paid', paid_at = now() WHERE id = 42;INSERTINTO outbox (id, topic, payload, created_at)VALUES (gen_random_uuid(), 'order.paid', '{"order_id": 42, "amount": 1999}', now());COMMIT;-- A relay process publishes and marks rows, at least onceSELECT id, topic, payload FROM outboxWHERE published_at ISNULLORDERBY created_atLIMIT100FORUPDATESKIPLOCKED;UPDATE outbox SET published_at = now() WHERE id = ANY($1);
Why it matters Publishing straight to Kafka inside a request can lose the event if the process dies after the commit; the outbox cannot.
Batch processing crunches a large, bounded dataset on a schedule. Stream processing handles an unbounded flow of events continuously, producing results in seconds.
In detail
MapReduce and Spark split a batch job across many machines: map transforms records, shuffle groups by key, reduce aggregates. Stream processors such as Flink and Kafka Streams keep state per key, group events into windows (tumbling, sliding, session) and deal with late events using watermarks. Many companies run both: streams for live dashboards, batch nightly to correct and backfill.
Batch and stream processing
Aspect
Batch
Stream
Data
Bounded, a day of logs
Unbounded, events as they happen
Latency
Minutes to hours
Milliseconds to seconds
Tools
Spark, Hadoop MapReduce, dbt
Flink, Kafka Streams, Spark Structured Streaming
Correctness
Easy to recompute
Late and out of order events
Typical use
Reports, model training, backfills
Fraud alerts, live dashboards, trending
wordcount.tsTypeScript
const sorted = [...pairs].sort(([a], [b]) => (a < b ? -1 : a > b ? 1 : 0)); // the shufflefor (const [k, v] of sorted) { out[k] = (out[k] ?? 0) + v;
// Batch: classic MapReduce word countfunction* mapPhase(lines: string[]): Generator<[string, number]> {for (const line of lines) {for (const word of line.toLowerCase().split(/\s+/).filter(Boolean)) {yield [word, 1]; } }}functionreducePhase(pairs: Iterable<[string, number]>): Record<string, number> {const sorted = [...pairs].sort(([a], [b]) => (a < b ? -1 : a > b ? 1 : 0)); // the shuffleconst out: Record<string, number> = {};for (const [k, v] of sorted) { out[k] = (out[k] ?? 0) + v; }return out;}console.log(reducePhase(mapPhase(["to be or", "not to be"]))); // { be: 2, not: 1, or: 1, to: 2 }// Stream: tumbling one minute windows of clicks per URLtypeClick = { url: string; ts: number };const windows = newMap<string, number>(); // key is `${url}@${window}`exportfunctiononClick(event: Click): void {const window = Math.floor(event.ts / 60);const key = `${event.url}@${window}`; windows.set(key, (windows.get(key) ?? 0) + 1); // emit when the window closes}
Why it matters The shuffle, grouping every value for one key onto one machine, is what makes MapReduce scale and what makes it slow.
An idempotent operation has the same effect whether it runs once or ten times. It is what makes retries safe in a world where networks drop replies.
In detail
Delivery is at most once (may lose), at least once (may duplicate) or effectively exactly once (at least once plus deduplication). Clients send an idempotency key; the server stores the first result for that key and returns it on repeats. Natural keys, upserts and conditional writes (version numbers) give idempotency without extra tables. Payments, orders and anything charged must be idempotent.
Machines' clocks drift, so wall clock time cannot reliably order events across servers. Logical clocks order events by cause instead of by time.
In detail
NTP keeps clocks within milliseconds, but leap seconds, pauses and drift break assumptions. Lamport clocks give a total order consistent with causality; vector clocks detect concurrent updates, which Dynamo style stores use to keep conflicting versions. Google Spanner uses TrueTime, clocks with a known error bound, and waits out the uncertainty to give global ordering.
lamport.tsTypeScript
tick(): number {this.time += 1;returnthis.time;}receive(remoteTime: number): number {this.time = Math.max(this.time, remoteTime) + 1;
exportclassLamportClock { time = 0;tick(): number {// a local eventthis.time += 1;returnthis.time; }send(): number {// attach to outgoing messagereturnthis.tick(); }receive(remoteTime: number): number {// merge on incoming messagethis.time = Math.max(this.time, remoteTime) + 1;returnthis.time; }}const a = newLamportClock();const b = newLamportClock();const t = a.send(); // a: 1b.tick(); // b: 1 (concurrent)b.receive(t); // b: 2, now ordered after a's sendconsole.log(a.time, b.time); // 1 2
Why it matters If event x caused event y, x always has a smaller Lamport time; the reverse is not guaranteed, which is what vector clocks add.
Consensus lets a group of machines agree on one value or one leader even when some fail. Raft and Paxos power etcd, ZooKeeper, Consul and the metadata of most distributed databases.
In detail
Raft elects a leader by majority vote; the leader appends commands to a replicated log and commits an entry once a majority has stored it. With five nodes, two can fail. Quorums generalise this: with N replicas, writing to W and reading from R where W + R > N guarantees a read sees the latest write. Split brain is prevented because only a majority can act.
When one business action spans several services or databases, you need a way to keep them consistent. Two phase commit locks everyone until all agree; sagas use a chain of local transactions with compensating undo steps.
In detail
Two phase commit has a coordinator ask every participant to prepare, then commit; it is strongly consistent but blocks if the coordinator dies. Sagas suit microservices: book flight, then hotel, then car, and if the car fails, cancel the hotel and flight. Sagas are eventually consistent and every step needs an idempotent compensation. Orchestrated sagas have a central coordinator; choreographed ones react to events.
saga.tsTypeScript
for (const step of steps) {await step.do(ctx);} catch (error) {for (const finished of done.reverse()) {await finished.undo(ctx); // must be idempotent and retriedthrownewSagaFailed(step.name, { cause: error });
In a cluster, instances start, die and move constantly. Service discovery keeps a live registry of healthy addresses so callers find each other without hard coded IPs.
In detail
Client side discovery has callers query a registry such as Consul or etcd and pick an instance; server side discovery puts a load balancer in front. Kubernetes Services and DNS do this for you. Health checks remove dead instances, and a service mesh such as Istio or Linkerd adds retries, timeouts and mutual TLS as sidecars. The same consensus backed stores also hold feature flags and dynamic config.
Autoscaling adds instances when load rises and removes them when it falls, so you pay for what you use and survive spikes. Capacity planning decides the floor, ceiling and headroom.
In detail
Scale on a signal that predicts saturation: CPU, request rate, queue depth or p99 latency. Stateless services scale horizontally easily; stateful ones need sharding. New instances take time to boot, so keep headroom (often 30 to 50 percent), pre warm before known events, and use cooldowns to avoid flapping. Queue length is the best signal for workers.
Running in several regions cuts latency for distant users and survives a whole region failing. It is also where consistency, cost and complexity peak.
In detail
Active passive keeps a warm standby region and fails over by DNS; active active serves traffic everywhere and must handle concurrent writes, often with per region ownership of data or conflict free replicated data types. Recovery point objective (how much data you may lose) and recovery time objective (how long you may be down) drive the design. Rehearse failovers; untested ones fail.
Traffic follows health; data follows replicationClientEdgeService
# Disaster recovery plan, active passiveservice: checkoutprimary_region: ap-south-1secondary_region: eu-west-1objectives: rpo: 1m # at most one minute of writes may be lost rto: 15m # back up within fifteen minutesdata: database: async replica in secondary, lag alarm at 30s object_storage: cross region replication on cache: rebuilt on failover, no replicationtraffic: dns: health checked failover record, ttl 60 warm_capacity: 25% of primary, autoscale on promotiondrills: frequency: quarterly last_run: 2026-07-14
Why it matters Writing the RPO and RTO down turns vague hopes into numbers you can test and budget for.
A monolith is one deployable unit; microservices split a system into independently deployed services that own their data. Neither is better by default; team size and change rate decide.
In detail
Monoliths are simpler to build, test and run, and a modular monolith with clear boundaries gets most of the benefit. Microservices let teams deploy independently and scale parts separately, but add network failures, distributed data, versioning and heavy operations. Split along business capabilities, not technical layers, and avoid a distributed monolith where every change touches five services.
Monolith, modular monolith and microservices
Aspect
Monolith
Modular monolith
Microservices
Deploy
All at once
All at once
Each service alone
Data
One database
One database, owned schemas
Database per service
Calls
In process
In process through interfaces
Over the network
Operations
Simple
Simple
Heavy: tracing, discovery, CI per service
Team size
One team
A few teams
Many independent teams
boundaries.mdNotes
Split by capability: Orders | Payments | Catalog | ShippingEach owns its database; talk via APIs and events
# Signs you should split a service out- A team blocks on another team's release train- One part needs very different scaling (search, video encoding)- One part needs a different language or runtime- Failure in one area keeps taking down unrelated features# Signs you split too early- Every feature changes three services at once- Services share a database "for now"- Most calls are chatty synchronous chains- One team owns all the services anyway# Rule of thumbStart with a modular monolith. Extract a service when a modulehas its own owners, data and scaling story.
Why it matters Shared databases between services are the clearest sign of a distributed monolith.
Resilience patterns stop one slow or broken dependency from dragging the whole system down: timeouts bound waiting, retries with backoff ride out blips, circuit breakers stop hammering a dead service, bulkheads isolate resources.
In detail
Every network call needs a timeout. Retry only idempotent calls, with exponential backoff and jitter, and cap total retries to avoid retry storms. A circuit breaker opens after repeated failures, fails fast for a while, then lets a trial request through. Bulkheads give each dependency its own pool of threads or connections. Graceful degradation serves a simpler response, such as cached recommendations, rather than an error.
breaker.tsTypeScript
if (this.openedAt !== null && Date.now() - this.openedAt < this.cooldownMs) {returnfallback(); // open: fail fast}try {const result = awaitfn(); // closed, or half open trial} catch {this.fails += 1;
Observability is being able to ask new questions about a live system. Metrics show trends, logs explain single events, and traces follow one request across many services.
In detail
Track the four golden signals per service: latency, traffic, errors and saturation. Use RED (rate, errors, duration) for services and USE (utilisation, saturation, errors) for resources. Distributed tracing with OpenTelemetry propagates a trace id through every hop. Define service level objectives, such as 99.9% of requests under 300 ms, and alert on error budget burn rather than on every spike.
LatencyHow long requests take, at p50, p95 and p99.
TrafficRequests per second, messages per second.
ErrorsRate of failed requests, including wrong answers.
SaturationHow full the system is: CPU, memory, queue depth.
Security: authentication, authorization and encryption
Security design covers who you are (authentication), what you may do (authorization), keeping data secret in transit and at rest (encryption), and limiting damage when something leaks.
In detail
Use TLS everywhere, including between services, and encrypt data at rest with keys in a key management service. OAuth 2.0 and OpenID Connect handle login; short lived JWTs or opaque tokens carry identity; every service checks authorization per request. Apply least privilege, rotate secrets, validate input, rate limit public endpoints and log security events. Assume breach: segment networks and keep blast radius small.
Safe delivery means small changes, released gradually, with a fast way back. Blue green swaps whole environments; canaries send a small slice of traffic to the new version first; feature flags decouple deploying code from releasing features.
In detail
Watch the canary's error rate and latency against the baseline and roll back automatically if they regress. Database changes use expand and contract: add new columns, write to both, backfill, switch reads, then remove the old ones, so every step is backward compatible. Most outages start with a change, so deploy frequency and rollback speed are reliability features.
canary.tsTypeScript
for (const pct of STAGES) {await lb.setWeight(version, pct);awaitsleep(soakMinutes * 60_000);if (canary.errorRate > base.errorRate * 1.5 || canary.p99 > base.p99 * 1.2) {await lb.setWeight(version, 0); // instant rollbackthrownewRolloutAborted(`regressed at ${pct}%`);
Turn a long URL into a seven character code and redirect anyone who opens it, billions of times a month, in a few milliseconds.
In detail
It is read heavy, about 100 redirects per link created. Generate codes from a unique 64 bit id encoded in base62 (no collisions, no lookups) or by hashing with collision checks. Store code to URL in a key value store or a sharded SQL table keyed by code. Redirects hit a cache first; most traffic goes to a small set of popular links. Use 301 for cacheable permanent links or 302 when you need every click counted, and push click events to a queue for analytics rather than writing them on the redirect path.
Redirects read the cache first; clicks go asyncClientEdgeServiceCacheDataExternalQueue
GET /aB3x9Q1 lookup2 on misson createclick eventClientClientLoad balancerEdgeLink serviceServiceRedis cacheCacheLink storeDataID generatorExternalClick queueQueue
Back of the envelope
New links100 M per monthabout 40 per second
Redirects10 B per monthabout 4,000 per second, 20k peak
Limit how many requests each user, key or IP can make per window, consistently across every server, adding under a millisecond to each request.
In detail
Run the limiter in the API gateway or as middleware backed by Redis, so all instances share counters. Token bucket allows short bursts; sliding window counter is accurate and cheap. Make each check atomic with a Lua script or INCR plus EXPIRE. Return 429 with Retry-After and rate limit headers. Decide to fail open (allow traffic if Redis is down) or closed (block) per endpoint. Keep local in memory pre checks for very hot keys.
Every gateway node shares one counterClientEdgeCacheService
Send push, SMS and email notifications for many products at millions per minute, respecting user preferences, rate limits and quiet hours, without duplicates.
In detail
Producers call one notification API with a user, template and data. The service checks preferences and limits, renders per channel, and enqueues to separate queues per channel so a slow SMS provider never delays push. Channel workers call APNs, FCM, an SMS gateway or an email provider with retries, and record delivery status. An idempotency key per notification stops duplicates on retry. Scheduled and digest notifications go through a delay queue.
One queue per channel keeps providers isolatedServiceDataQueueExternal
Show each user a ranked feed of posts from people they follow, in under 200 milliseconds, for hundreds of millions of users, including celebrities with millions of followers.
In detail
Fan out on write pushes each new post id into every follower's cached feed list, making reads instant but costly for celebrities. Fan out on read pulls from followees at read time, cheap to write but slow to read. The hybrid: fan out on write for normal users, merge celebrity posts at read time. Feeds hold only post ids in Redis lists; posts and media are fetched in bulk and cached separately; a ranking service orders candidates.
Push for most, pull for celebritiesClientServiceDataQueueCache
savenew postpush idsreadhydrateUser postsClientPost serviceServicePosts DBDataFan out queueQueueFan out workersServiceFeed cacheCacheFeed serviceService
Deliver one to one and group messages in real time, keep them in order, show delivery and read receipts and presence, and sync history across devices.
In detail
Clients hold a WebSocket to a chat gateway; a session service maps users to gateway servers. A message is stored first, then routed to the recipient's gateway or, if offline, a push notification. Per conversation sequence numbers keep order. Messages are stored in a wide column store such as Cassandra partitioned by conversation and sorted by time. Large groups use pub sub fan out. Presence uses heartbeats with a short TTL. End to end encryption keeps servers from reading content.
Store first, then deliver or pushClientEdgeServiceCacheDataExternal
WebSocketsendstorewhere is Bobdeliverif offlineAliceClientBobClientGateway 1EdgeGateway 2EdgeChat serviceServiceSessionsCacheMessagesDataPushExternal
Back of the envelope
Daily users500 M40 messages each
Messages20 B per dayabout 230k per second
Storageabout 100 bytes each2 TB per day
Connections50 M concurrentabout 50k per gateway box
-- Cassandra: one partition per conversation per time bucketCREATETABLE messages ( conversation_id uuid, bucket int, -- e.g. days since epoch, keeps partitions bounded seq bigint, -- per conversation sequence, assigned on write sender_id uuid, body blob, -- end to end encrypted payload sent_at timestamp,PRIMARYKEY ((conversation_id, bucket), seq)) WITHCLUSTERINGORDERBY (seq DESC);-- Latest 50 messages in a conversation, one partition readSELECT seq, sender_id, body, sent_atFROM messagesWHERE conversation_id = ? AND bucket = ?LIMIT50;-- Per user, per conversation read pointer for receipts and unread countsCREATETABLE read_state ( user_id uuid, conversation_id uuid, last_read_seq bigint,PRIMARYKEY (user_id, conversation_id));
Why it matters Bucketing the partition by time stops a ten year old busy group chat from becoming one enormous, slow partition.
Suggest the top completions as someone types, within about 100 milliseconds per keystroke, ranked by popularity and freshness.
In detail
Precompute: a trie where each node stores its top k completions, built offline from query logs (batch) and refreshed with recent trends (stream). Serve the trie from memory, sharded by prefix range, with a CDN or browser cache for the most common short prefixes. Clients debounce keystrokes and cancel stale requests. Filter offensive suggestions at build time.
Built offline, served from memoryClientEdgeServiceCacheData
Download billions of pages, follow their links, avoid duplicates and traps, and never overload any single website.
In detail
A URL frontier holds what to fetch next, prioritised by importance and freshness and split into per host queues so politeness (one request per host every few seconds, robots.txt) is enforced. Fetchers download pages; parsers extract links and content; a seen URL set, often a Bloom filter, drops duplicates; content hashes catch mirrored pages. DNS results are cached. The system is a giant pipeline of queues and workers that can run for months.
A loop of queues and workersClientQueueServiceCacheData
politehtmlsavenew links?add newSeed URLsClientURL frontierQueueFetchersServiceParsersServiceSeen set (Bloom)CachePage storeData
Back of the envelope
Pages1 B per monthabout 400 per second
Page sizeabout 500 KB200 TB per month raw
Politeness1 request per host per 2 sper host queue
Seen URLs10 B entriesBloom filter about 12 GB
crawler.tsTypeScript
const url = frontier.nextPolite();const html = awaitfetchPage(url); // respects robots.txtfor (const link ofextract(html, url)) {if (!seenUrls.mightContain(link)) { seenUrls.add(link); frontier.add(link);
Upload, transcode and stream video to millions of devices at once, starting playback in under two seconds and adapting quality to every connection.
In detail
Uploads land in object storage, then a pipeline splits each video into chunks and transcodes them in parallel into many resolutions and codecs. Players use adaptive bitrate streaming (HLS or DASH): a manifest lists segments at each quality and the player switches every few seconds based on bandwidth. Nearly all bytes come from a CDN; Netflix places its own Open Connect boxes inside ISPs and pre fills them overnight with what each region will watch. Metadata, search, recommendations and watch history are separate services.
Transcode once, serve from the edge foreverClientDataQueueServiceEdge
new videochunksrenditionspre fillABR streambrowseCreator uploadClientRaw storageDataTranscode queueQueueTranscodersServiceSegments storeDataCDN / Open ConnectEdgeViewersClientCatalog + recsService
Back of the envelope
Concurrent viewers50 M at peakevening prime time
Average bitrateabout 5 MbpsHD
Peak egressabout 250 Tbpswhy CDN is everything
Renditions10 to 30 per titlecodecs times resolutions
#EXTM3U#EXT-X-VERSION:6# Master playlist: the player picks a rendition and can switch every segment#EXT-X-STREAM-INF:BANDWIDTH=400000,RESOLUTION=426x240,CODECS="avc1.4d401e,mp4a.40.2"240p/index.m3u8#EXT-X-STREAM-INF:BANDWIDTH=800000,RESOLUTION=640x360,CODECS="avc1.4d401e,mp4a.40.2"360p/index.m3u8#EXT-X-STREAM-INF:BANDWIDTH=2500000,RESOLUTION=1280x720,CODECS="avc1.4d401f,mp4a.40.2"720p/index.m3u8#EXT-X-STREAM-INF:BANDWIDTH=5000000,RESOLUTION=1920x1080,CODECS="avc1.640028,mp4a.40.2"1080p/index.m3u8#EXT-X-STREAM-INF:BANDWIDTH=15000000,RESOLUTION=3840x2160,CODECS="hvc1.2.4.L150,mp4a.40.2"2160p/index.m3u8# Each rendition lists 4 second segments: seg_00001.m4s, seg_00002.m4s, ...
Why it matters Because every rendition shares the same segment boundaries, the player can change quality mid stream without a stutter.
Let users upload photos, follow others, and scroll a feed of images that load instantly worldwide, with likes and comments.
In detail
Clients upload directly to object storage with a pre signed URL, then a worker creates thumbnails and sizes. Metadata (post id, owner, caption, image keys) lives in a sharded database keyed by user id. Images are served through a CDN with long cache lifetimes, since a published photo never changes. The feed reuses the news feed design. Likes and counts use sharded counters and are eventually consistent.
Bytes go around the API, not through itClientServiceDataEdge
Keep a folder identical across a user's devices and collaborators, uploading only what changed, working offline, and resolving conflicts sensibly.
In detail
Split files into content addressed chunks (around 4 MB, named by their hash). Uploading a file means sending only chunks the server does not already have, which gives deduplication and resumable uploads for free. A metadata service records each file version as a list of chunk hashes; a notification service tells other devices to pull changes. Conflicts produce a conflicted copy rather than silent loss.
Chunks move once; versions are just lists of hashesClientServiceData
new chunkscommit versionchangedpullfetch chunksLaptopClientPhoneClientMetadata serviceServiceBlock storeDataVersions DBDataChange notifierService
Back of the envelope
Users500 M100 M daily
Chunk size4 MBcontent addressed
Dedup savingsoften 30% or moresame files across users
Match riders with nearby drivers within seconds, track millions of moving cars live, price trips and handle the trip lifecycle safely.
In detail
Drivers send GPS every few seconds to a location service that keeps the latest position in memory, indexed by geohash or H3 cells. A rider request queries nearby cells, ranks candidate drivers by ETA, and offers the trip to one driver at a time with a short timeout. Trip state lives in a strongly consistent store because double assignment is unacceptable. Surge pricing compares demand and supply per cell. Payments and receipts are asynchronous.
Index by cell, match by ETA, assign atomicallyClientServiceCacheData
GPS every 4 supdate cellrequestnearbyassignofferDriver appsClientRider appClientLocation serviceServiceGeo index (H3)CacheMatchingServiceTrip storeData
Back of the envelope
Active drivers5 Mupdates every 4 s
Location writesabout 1.25 M per secondkept in memory
Design a proximity service like Yelp or Google Maps places
Find businesses within a radius of a user, filtered and ranked, at very high read rates, where the data itself changes slowly.
In detail
Places rarely move, so precompute a spatial index: geohash prefixes, a quadtree or S2 cells. A search looks up the user's cell plus neighbours, filters by exact distance and category, then ranks by rating and distance. Read replicas and caches by cell handle the load; writes (new businesses, edits) go to the primary and are indexed within minutes. Reviews and photos are separate services.
Back of the envelope
Places200 Mchange slowly
Search QPSabout 5kread only path
Index sizea few GBfits in memory
Radius0.5 to 20 kmpick cell precision to match
nearby.sqlSQL
SELECT id, name FROM placesWHERE geohash6 IN (:cell, :n1, :n2, ...) AND ST_DWithin(geo, :point, 2000);
-- Places table with both a coarse cell and an exact geography columnCREATETABLE places ( id bigintPRIMARYKEY, name textNOTNULL, category textNOTNULL, rating numeric(2,1), geohash6 char(6) NOTNULL, geo geography(Point, 4326) NOTNULL);CREATEINDEX places_cell ON places (geohash6, category);CREATEINDEX places_geo ON places USING gist (geo);-- 1. Cheap filter by the user's cell and its 8 neighbours-- 2. Exact radius check, 3. rankSELECT id, name, rating, ST_Distance(geo, :point) AS metersFROM placesWHERE geohash6 = ANY(:cells)AND category = 'cafe'AND ST_DWithin(geo, :point, 2000)ORDERBY rating DESC, metersLIMIT20;
Why it matters The coarse cell narrows millions of rows to hundreds; only then does the precise distance calculation run.
Rank millions of players by score, update instantly as games finish, and answer top 100 and what is my rank in milliseconds.
In detail
A Redis sorted set does exactly this: ZINCRBY updates a score in O(log n), ZREVRANGE returns the top k, ZREVRANK returns a player's rank. Keep one sorted set per leaderboard and season. For hundreds of millions of players, shard by score range or keep exact ranks for the top and approximate ranks elsewhere. Persist scores in a database as the source of truth and rebuild the set if Redis is lost.
Design ticket booking like BookMyShow or Ticketmaster
Let thousands of fans compete for the same seats at the moment sales open, without ever selling one seat twice, and keep the site up under the stampede.
In detail
Hold seats briefly: selecting a seat creates a reservation with a short expiry (five to ten minutes) using a conditional update or row lock, and payment confirms it. A virtual waiting room admits users at a controlled rate so the booking system sees a steady flow. Inventory for a show fits on one shard, which keeps transactions local. Seat maps are cached and refreshed often; the final check is always against the database.
Queue the crowd, hold, then confirmClientEdgeServiceDataExternal
-- Hold a seat: succeeds for exactly one user, even with thousands racingUPDATE seatsSET status = 'held', held_by = :user_id, hold_until = now() + interval'8 minutes', version = version + 1WHERE show_id = :show_idAND seat_no = :seat_noAND (status = 'free'OR (status = 'held'AND hold_until < now()))RETURNING version;-- 0 rows returned means someone else got it-- Confirm after payment, only if the hold is still oursUPDATE seatsSET status = 'sold', order_id = :order_idWHERE show_id = :show_idAND seat_no = :seat_noAND status = 'held'AND held_by = :user_idAND hold_until >= now();
Why it matters Putting the whole check inside one conditional UPDATE lets the database settle every race; no application lock is needed.
11
Phase 11, modules 78 to 82
Hard mode
Payments, key value stores, collaboration, monitoring and scheduling
Designs where correctness, consistency or sheer volume make every shortcut dangerous.
Accept payments through external processors, move money between accounts correctly, and reconcile everything, where correctness matters more than speed.
In detail
Every request carries an idempotency key. A payment service records intent, calls the payment service provider, and handles asynchronous webhooks for the final status. Money movements are recorded in a double entry ledger: every transaction writes balanced debit and credit rows, append only, never updated. Nightly reconciliation compares the ledger with the provider's settlement files. Retries, timeouts and unknown states are designed explicitly; an unknown is resolved by querying the provider, never by guessing.
pay + idem keyauthorisestatusrecordintentnightly checkCheckoutClientPayment serviceServicePSP (Stripe, Razorpay)ExternalLedgerDataWebhooksServiceReconciliationService
Back of the envelope
Payments10 M per dayabout 115 per second
Peak10x averagesales days
Ledger rows2 to 4 per paymentappend only
Availability99.99%but correctness first
ledger.sqlSQL
INSERTINTO ledger(txn_id, account, amount) VALUES (:t, 'customer:42', -1999), (:t, 'merchant:7', 1999); -- sums to zero
CREATETABLE ledger_entries ( id bigserialPRIMARYKEY, txn_id uuidNOTNULL, account_id textNOTNULL, amount bigintNOTNULL, -- minor units, never floats currency char(3) NOTNULL, created_at timestamptzNOTNULLDEFAULTnow());-- Append only: revoke UPDATE and DELETE for the application role-- A payment of 19.99 with a 0.59 fee: three rows that sum to zeroBEGIN;INSERTINTO ledger_entries (txn_id, account_id, amount, currency) VALUES (:txn, 'customer:42', -1999, 'INR'), (:txn, 'merchant:7', 1940, 'INR'), (:txn, 'platform:fees', 59, 'INR');COMMIT;-- Invariant checked continuouslySELECT txn_id FROM ledger_entries GROUPBY txn_id HAVING sum(amount) <> 0;
Why it matters Integers in minor units and a zero sum invariant make whole classes of rounding and lost money bugs impossible.
Store and retrieve values by key across hundreds of machines, staying available when nodes and networks fail, with tunable consistency.
In detail
Partition keys with consistent hashing and virtual nodes; replicate each key to N successors on the ring. Reads and writes use quorums (R and W) chosen per request. Conflicting versions are detected with vector clocks and resolved by the client or last writer wins. Hinted handoff covers temporarily down nodes, read repair and Merkle tree anti entropy heal replicas, and gossip spreads membership. Each node stores data in an LSM tree.
Each key belongs to the next node clockwise; adding a node only steals keys from its neighbourBCDEAk1→Dk2→Bk3→Ck4→Ck5→Ck6→C
Let several people edit the same document at once, seeing each other's changes in real time, with no lost edits and the same final result for everyone.
In detail
Each keystroke becomes an operation (insert or delete at a position). Operational transformation adjusts concurrent operations against each other through a central server that orders them; CRDTs give every character a unique, ordered id so replicas merge without a central authority and work offline. Clients connect via WebSocket to a document session server, operations are appended to a log, and periodic snapshots make loading fast. Cursors and presence are ephemeral and not stored.
Collect metrics from thousands of servers every few seconds, store them efficiently, query dashboards in milliseconds and fire alerts reliably.
In detail
Agents scrape or push metrics (name, labels, timestamp, value). A time series database such as Prometheus, VictoriaMetrics or M3 compresses points with delta of delta timestamps and XOR floats, achieving about 1 to 2 bytes per point. Data is downsampled with age: raw for days, 5 minute rollups for months. Label cardinality is the main cost driver. The alerting pipeline must be more reliable than the systems it watches, so it is often run separately.
Ingest, store compressed, query, alertClientEdgeQueueDataService
every 10 swriterulesqueryServers + agentsClientCollectorsEdgeIngest queueQueueTime series DBDataAlert managerServiceDashboardsClient
Back of the envelope
Series10 M activelabel cardinality
Samples1 M per second10 s scrape
Compressedabout 1.4 bytes per sampleGorilla encoding
Run millions of scheduled and delayed jobs (send this email in an hour, rebuild that report at midnight) on time, once, even when workers crash.
In detail
Store jobs with their next run time in a database indexed by that time. Scheduler nodes poll for due jobs in small batches using SELECT FOR UPDATE SKIP LOCKED or a leased claim, push them onto a queue, and workers execute them. A lease with a timeout returns stuck jobs; idempotent job handlers make the inevitable duplicate harmless. Time buckets or a delay queue (Redis sorted set by timestamp) suit very large volumes. Recurring jobs compute the next run after each execution.
Back of the envelope
Jobs100 M per dayabout 1,200 per second
Poll interval1 sbatches of 100
Lease5 minutesthen retried
Indexpartial on pendingstays small
claim.sqlSQL
SELECT id FROM jobs WHERE run_at <= now() AND status='pending'ORDERBY run_at LIMIT100FORUPDATESKIPLOCKED;
CREATETABLE jobs ( id bigserialPRIMARYKEY, kind textNOTNULL, payload jsonbNOTNULL, run_at timestamptzNOTNULL, status textNOTNULLDEFAULT'pending', -- pending, running, done, failed lease_until timestamptz, attempts intNOTNULLDEFAULT0);CREATEINDEX jobs_due ON jobs (run_at) WHERE status = 'pending';-- Each scheduler claims a batch without blocking the othersWITH due AS (SELECT id FROM jobsWHERE status = 'pending'AND run_at <= now()ORDERBY run_atLIMIT100FORUPDATESKIPLOCKED)UPDATE jobs SET status = 'running', lease_until = now() + interval'5 minutes', attempts = attempts + 1FROM due WHERE jobs.id = due.idRETURNING jobs.id, jobs.kind, jobs.payload;-- Reaper: return expired leases to the poolUPDATE jobs SET status = 'pending'WHERE status = 'running'AND lease_until < now();
Why it matters SKIP LOCKED lets many schedulers poll the same table in parallel without ever handing the same job to two of them.
12
Phase 12, modules 83 to 84
Mastery
The interview playbook and the capstone
A repeatable way to run any design conversation, and one final system that uses every phase.
A repeatable script for any design question: clarify, estimate, sketch, detail, deep dive and wrap up, while thinking out loud and naming trade-offs.
In detail
Spend the first five minutes on requirements and the next five on numbers; they decide everything after. Draw the simplest design that works, then evolve it under load. Pick one or two deep dives that matter for this product (hot keys for a feed, double booking for tickets). Mention failure modes and what you would monitor. Common mistakes: jumping to Kafka and microservices before stating requirements, no numbers, and going silent while thinking.
Requirements5 minEstimates5 minHigh level design10 minAPI and data5 minDeep dives15 minWrap up5 min
45 minutes in total; the first ten decide whether the rest goes well.
script.mdNotes
"Before I design, can I confirm the core features and the scale?""Reads outnumber writes 100 to 1, so I'll cache redirects."
# Phrases that show structureOpening- "Let me confirm the core use cases and what is out of scope."- "Roughly how many daily users, and is this read or write heavy?"Numbers- "That is about 1,200 writes per second on average, maybe 5,000 at peak."- "A year of data is around 40 TB, so one database will not hold it."Design- "I'll start simple and then scale the part that breaks first."- "The trade-off here is freshness versus latency; I'll take eventual consistency for the feed."Deep dive- "The hardest part is the celebrity fan out, so let me spend time there."Wrap up- "Single points of failure are X and Y; I'd add replicas and a second region."- "I'd alert on p99 latency, error rate and queue depth."
Why it matters Interviewers grade the reasoning they can hear, so narrating choices matters as much as the diagram.
Put every phase together: design a short video platform with uploads, a feed, chat, notifications, payments for creators and global delivery, then defend it.
In detail
Write the requirements and estimates, draw the high level design, then specify the API and data model for each service. Walk through one request path end to end, from DNS and the CDN to the database and back. Pick three deep dives: transcoding and adaptive streaming, feed fan out for creators with millions of followers, and creator payouts on a double entry ledger. Finish with failure scenarios, SLOs and a cost estimate.
Short answers to the questions that come up most while preparing for system design, each linked to the module that covers it in depth.
Do I need to know data structures before system design?
Yes, at least the basics. Hash tables, trees, heaps, queues and graphs are the parts every database, cache and queue is built from, and interviewers expect you to reason about their costs.
Around fifteen classic ones done properly beats fifty skimmed. The URL shortener, news feed, chat, video streaming, ride sharing and payments cover most patterns.
Only after requirements and estimates show a need, such as fan out to many consumers or independent team scaling. Reaching for them first is a common red flag.
During a network partition you choose between rejecting requests to stay consistent or serving possibly stale data to stay available. Most of the time you trade latency against consistency instead.
Every source linked from the modules above, grouped by the phase that uses it and then by where it lives. 220 links in total, all opening in a new tab.
The system design dictionary
Phase 00, Before you start
3
Complexity and data structures
Phase 01, Counting the cost
25
Hashing, limiting, caching and IDs
Phase 02, Algorithms that run the internet
22
System design fundamentals
Phase 03, Thinking like an architect
19
Networking, load balancing and APIs
Phase 04, The roads between machines
19
Databases and storage
Phase 05, Where data lives
23
Caching and messaging
Phase 06, Remember and relay
21
Distributed systems
Phase 07, Many machines, one truth
24
Observability, security and delivery
Phase 08, Keep it running
9
Bitly, rate limiter, notifications, feeds, chat, search and crawling
Phase 09, Classic designs, part one
18
Netflix, Instagram, Dropbox, Uber, Yelp, leaderboards and tickets
Phase 10, Classic designs, part two
19
Payments, key value stores, collaboration, monitoring and scheduling
Phase 11, Hard mode
15
The interview playbook and the capstone
Phase 12, Mastery
3
Credits
The technologies this roadmap teaches and the tools used to build the page. The people behind it are listed in the footer.