High Level Design Interview Program — 19 Case Studies to Crack System Design
A complete 10-week High Level Design curriculum: the interview framework, capacity estimation, CAP, load balancing, consistent hashing, caching, sharding, replication, Kafka, and 19 case studies (URL Shortener, feeds, WhatsApp, Uber, payments, ChatGPT, RAG) explained topic by topic.
Introduction
High level design is the discipline of reasoning about scale before you own it. The interviewer is not grading whether you reproduce a reference architecture — they are grading whether your architecture follows from the constraints you clarified. A diagram that nobody can derive from the requirements is decoration, not design.
This guide is the full 10-week HLD cohort. It is organised topic by topic: every concept is explained with what it is, how it works, when to choose it, and what it costs — and then applied in a case study. You can go deep on the object-model side in the LLD companion: Low Level Design Interview Program.
Every case study is answered through the same five questions, in this order:
- Clarify — functional and non-functional requirements.
- Estimate — QPS, storage, bandwidth from first principles.
- Model — API and data model.
- Design — components, each with a reason to exist.
- Trade-offs — bottlenecks, consistency costs, and what changes at 10×.
Read the theory weeks (1–3) as reference; read the case-study weeks (4–10) as worked examples that apply it.
← Back to System Design overview
Quick index
| # | Topic | Description |
|---|---|---|
| 1 | Foundations & Mental Models | Framework, capacity estimation, latency numbers, CAP, consistency, storage. |
| 2 | Core Infrastructure Building Blocks | Load balancers, consistent hashing, CDN/DNS, caching, Redis. |
| 3 | Database & Messaging Systems | Transactions, sharding, replication, indexing, Kafka and event-driven fan-out. |
| 4 | High-Traffic & Search Systems | URL Shortener, distributed rate limiter, search autocomplete. |
| 5 | Feed & Streaming Systems | Twitter/Instagram, YouTube/Netflix, WhatsApp, resilience patterns. |
| 6 | Storage & Notification Systems | Key-value store, distributed cache, notifications, realtime transport. |
| 7 | Location & Scheduling Systems | Uber, Google Maps, task scheduler, geo-indexing. |
| 8 | Payments & Reliability | Payment gateway, leaderboards, observability, distributed transactions. |
| 9 | AI System Design: ChatGPT & Q&A | Streaming assistants, RAG, inference economics, tenant isolation. |
| 10 | AI Agents, Gateways & Mock Interviews | Agent platforms, multi-model gateways, evaluation, mocks. |
| 11 | Full Case Studies List (19) | Every case study, its focus, and where it was asked. |
Week 1: Foundations & Mental Models
Before any component exists, two things must exist: a clear picture of what the system does, and a rough sense of how big it is. Everything else — caching, sharding, queues — is a response to one of those two facts.
The HLD interview framework
A system design interview is a 45–60 minute conversation, not a monologue. The structure that consistently works:
| Phase | Time | What you produce |
|---|---|---|
| Requirements | 5 min | Functional scope + non-functional targets (latency, scale) |
| Estimation | 5 min | QPS, storage, bandwidth, cache size |
| High-level design | 20 min | Components, APIs, data model, data flow |
| Deep dive | 15 min | One or two hard parts: consistency, hotspots, failure |
| Wrap-up | 5 min | Bottlenecks, trade-offs, what you would improve |
The interviewer evaluates how you reason, not the final diagram. Two signals matter most: do you ask clarifying questions before designing, and do you justify every component against a requirement?
Functional vs non-functional requirements
Functional requirements are the verbs of the system — "post a tweet", "follow a user", "fetch a feed". Non-functional requirements are the constraints — latency, availability, consistency, durability, scale. Non-functional requirements drive architecture; aspirational phrases like "must scale" are meaningless until quantified.
| Non-functional | Question to ask | Architectural consequence |
|---|---|---|
| Latency | p99 read under 100 ms? | Cache at the edge + read replicas |
| Availability | 99.99% or 99.9%? | Redundancy vs a single leader |
| Consistency | Can a feed be 5 seconds stale? | Eventual consistency is allowed |
| Durability | Can we lose a payment? | Synchronous replication + ledger |
| Scale | 100 QPS or 100,000 QPS? | One DB vs sharding + a queue |
Back-of-the-envelope capacity estimation
Estimation turns a vague prompt into concrete numbers that decide your design. The method is always the same chain:
- Users — total users, then daily active users (DAU).
- Actions — requests per user per day.
- QPS —
(DAU × actions) / 86,400, then multiply by 2–3 for peak. - Storage —
writes per day × record size × retention. - Bandwidth —
record size × QPS(and the read:write ratio). - Cache — size the hot set, usually a fraction of daily reads (the 80/20 rule).
Assume 200M users, 10% daily active, 20 requests/user/day:
DAU = 200M × 10% = 20M
Requests/day = 20M × 20 = 400M
Average QPS = 400M / 86,400 ≈ 4,600 QPS
Peak QPS = 2–3× average ≈ 12,000 QPSStorage for 400M writes/day at 1 KB/record: 400M × 1 KB ≈ 400 GB/day ≈ 146 TB/year. Bandwidth = record size × QPS. If the read:write ratio is 100:1, reads dominate and caching + read replicas are the first moves.
Numbers worth memorising so you never stall:
| Quantity | Value |
|---|---|
| Seconds per day | 86,400 (≈ 100,000 for mental math) |
| 1 million requests/day | ≈ 12 QPS |
| 1 billion requests/day | ≈ 12,000 QPS |
| Powers of two | 2^10 = 1K, 2^20 = 1M, 2^30 = 1B, 2^40 = 1T |
| Memory read | ~100 ns |
| SSD random read | ~100 µs |
| Disk seek | ~10 ms |
| Same-datacenter round trip | ~0.5 ms |
| Cross-continent round trip | ~150 ms |
CAP theorem
CAP states that a distributed data store can guarantee at most two of Consistency, Availability, and Partition tolerance. The subtlety interviewers test: network partitions are a fact of life, so P is not optional — you are really choosing between C and A during a partition.
- Consistency (C) — every read sees the most recent write, across all nodes.
- Availability (A) — every request gets a non-error response (not necessarily the latest data).
- Partition tolerance (P) — the system keeps working when nodes cannot talk to each other.
| Choice | Under partition | Examples |
|---|---|---|
| CP | Reject writes/reads it cannot make consistent | HBase, ZooKeeper, Spanner-like |
| AP | Keep serving, accept divergence, reconcile later | Cassandra, DynamoDB, DNS |
Consistency models
CAP describes the extremes; real systems sit on a spectrum. Knowing the named models lets you be precise about what you actually need.
| Model | Guarantee | Typical use |
|---|---|---|
| Strong | Reads always see the latest committed write | Payments, inventory |
| Linearizable | Strong + operations appear in real-time global order | Leader election, distributed locks |
| Causal | Writes that are causally related are seen in order | Comment threads, collaboration |
| Read-your-writes | A user sees their own writes | Profiles, posting then refreshing |
| Monotonic reads | A user never sees time go backwards | Feeds, timelines |
| Eventual | Replicas converge given no new writes | Likes, counters, view counts |
PACELC: the rest of the story
CAP only covers partition behaviour. PACELC extends it: if there is a Partition, trade Against C; Else (normal operation) trade Latency against Consistency. Most of your design decisions happen in the "else" branch — synchronous replication costs latency even when nothing is broken.
SQL vs NoSQL
This is a decision about access patterns and scale, not fashion. Relational databases give you joins, constraints, and transactions; NoSQL stores trade some of that for horizontal scale and flexible schemas.
| Aspect | SQL | NoSQL |
|---|---|---|
| Schema | Fixed, migration-driven | Flexible / schema-on-read |
| Scaling | Vertical, then read replicas/sharding | Horizontal by design |
| Transactions | Strong ACID across tables | Often single-item or eventual |
| Joins | Native and efficient | Usually denormalised / application-side |
| Best for | Relationships, ledger, reporting | High write throughput, huge key space, simple lookups |
NoSQL is not one thing:
- Key-value — get/put by key (Redis, DynamoDB). Fastest, simplest.
- Document — JSON documents (MongoDB). Flexible schema.
- Wide-column — massive rows across nodes (Cassandra). Write-heavy time series.
- Graph — traversal of relationships (Neo4j). Social/fraud graphs.
Reuse the constants and access-pattern explorer to reason about the data model:
Interactive ER Explorer
Click an entity to inspect fields, indexes, and relationships.
Selected: users
Fields
- • id (PK)
- • email (unique)
- • name
- • created_at
Indexes
- • PRIMARY (id)
- • UNIQUE (email)
- • INDEX (created_at)
Relationships
Week 2: Core Infrastructure Building Blocks
Once traffic exceeds one machine, you need three things: something to spread traffic (load balancer), something to place data and traffic deterministically (consistent hashing), and something to reduce work (cache). This week is those building blocks plus the edge layer.
Load balancing: what and why
A load balancer distributes incoming requests across multiple servers so that no single server is overwhelmed, unhealthy servers are skipped, and capacity can be added or removed without downtime. It is both a performance tool (spread load) and a reliability tool (route around failure).
L4 vs L7 is the first distinction:
| Layer | Operates on | Sees | Strengths | Trade-off |
|---|---|---|---|---|
| L4 | TCP/UDP connections | IP, port | Very fast, protocol-agnostic | Cannot route by URL/header |
| L7 | HTTP requests | Path, headers, cookies | Smart routing, TLS, retries, sticky | More CPU; terminates requests |
Common algorithms and when each fits:
| Algorithm | How it works | Use when |
|---|---|---|
| Round robin | Rotate evenly | Homogeneous servers |
| Least connections | Pick the least busy | Variable request cost |
| IP/consistent hash | Same key → same server | Caching, session affinity |
| Weighted | Bigger servers get more traffic | Mixed hardware |
How L4 and L7 load balancing actually work
The layer determines what the balancer can see, and therefore what it can do.
L4 (transport layer) works on TCP/UDP:
- The client opens a TCP connection to the balancer's virtual IP.
- The balancer picks a backend (e.g. least connections) and forwards packets — often by NAT, without terminating the connection.
- Because the connection is pinned to that backend for its lifetime, every packet of that connection goes to the same server.
- It only sees IP and port, so it cannot route by URL or header — but it is extremely fast and works for any protocol (HTTP, gRPC, MySQL, MQTT).
L7 (application layer) terminates and parses the request:
- It terminates the client's TCP and TLS, decrypts, and reads the full HTTP request.
- It routes by path, host, header, or cookie — and can send each request to a different backend.
- It can retry idempotent requests, rewrite headers, add authentication, compress, and pool connections to backends.
- The cost is CPU: parsing HTTP and doing TLS handshakes is far more expensive than forwarding packets.
L4 — forwards packets, connection pinned:
client ──TCP──▶ [ L4 LB ] ──▶ backend #3
(sees IP:port only; every packet of this connection → #3)
L7 — parses HTTP, routes per request:
client ──TCP/TLS──▶ [ L7 LB: terminates TLS, parses HTTP ]
│
GET /api/orders ────┼──▶ API pool (backend #2)
GET /logo.png ──────┴──▶ static pool (backend #7)| Path | L4 | L7 |
|---|---|---|
| Terminates TLS | Pass-through (cannot see content) | Yes (can inspect and modify requests) |
| Routing unit | Connection | HTTP request |
| Can retry | No (cannot tell a request from a packet) | Yes, for idempotent requests |
| Typical use | Raw TCP, databases, extreme throughput | HTTP APIs, microservices, gateways |
Health checks. Active checks probe /health on each backend on an interval; passive checks eject a backend after N consecutive errors. Either way an unhealthy backend stops receiving traffic — this is what makes rolling deploys safe (drain the instance, deploy, return it to the pool).
High availability of the balancer itself. A single balancer is a single point of failure, so production uses an active-passive pair with a floating virtual IP, or anycast/DNS failover. Sticky sessions keep a client on one server but reduce failover flexibility — externalise session state to Redis instead where possible.
Try the distribution and failover behavior:
Load Balancer Playground
Route traffic, test failover, and compare stickiness vs fair distribution.
Click a server card to simulate health-check failure and observe failover behavior.
Consistent hashing: why modulo breaks
The naive way to assign a key to a node is hash(key) % N. The flaw: when N changes — a node is added or dies — almost every key remaps, causing a storm of cache misses or data movement.
Consistent hashing fixes this. Imagine a ring of hash values from 0 to 2^32. Each node is placed on the ring by hashing its identifier, and each key travels clockwise to the first node it meets. When a node is added or removed, only the keys between it and its predecessor move — roughly K/N keys instead of all of them.
Virtual nodes solve the second problem: with few physical nodes, the ring is uneven and one node gets a disproportionate share. By placing each physical node at many points on the ring, load evens out, and when a node leaves its virtual nodes spread the reassignment across all survivors.
This is why Redis Cluster, DynamoDB, Cassandra, and most distributed caches use it. Say "virtual nodes" out loud if asked how to keep distribution fair.
CDN: how a request is served from the edge
A Content Delivery Network is a geographically distributed set of edge servers (points of presence). It caches content close to users, which does two things: lower latency (a nearby hop instead of a cross-continent one) and reduced origin load (the origin serves each object roughly once per edge, not once per user).
The full request lifecycle. A user in Dhaka requests an image hosted behind a CDN:
- DNS resolution. Your domain points (via CNAME) at the CDN. Anycast or geo-DNS routes the lookup to the nearest edge PoP, which returns an edge IP.
- Edge cache lookup. The edge computes a cache key from the URL (plus query string and encoding, sometimes device type) and checks its local cache.
- Hit. The edge serves the object directly — single-digit milliseconds, with no origin involvement.
- Miss. The edge becomes a client itself: it fetches from the origin (or a mid-tier cache), stores the object honouring the origin's
Cache-Control/ETag, and serves it. The next user at the same edge gets a hit. - Revalidation. When a TTL expires, the edge asks the origin
If-None-Match/If-Modified-Since; a304 Not Modifiedrefreshes the TTL without re-downloading the body. - Invalidation. An explicit purge removes an object by URL or surrogate key; versioned filenames make old entries harmless.
user (Dhaka) ──DNS──▶ nearest edge PoP
│ HIT ──▶ serve immediately (~5–20 ms)
│ MISS
▼
edge ──▶ origin (US): fetch once, store, serve
later users at this edge → HITPull vs push.
- Pull CDN — the edge fetches from origin on the first miss, then caches by TTL. Easy; the good default.
- Push CDN — you upload content to the edge proactively. Better for large, predictable assets (VOD catalogs, software releases) where you want them warm everywhere.
What to cache, and how long, is controlled by response headers:
Cache-Control: public, max-age=31536000, immutable # hashed static asset
Cache-Control: public, s-maxage=60, stale-while-revalidate=300 # semi-dynamic GET
Cache-Control: private, no-store # per-user data — never at a shared edgeThe hard part is invalidation: hashed/versioned filenames (app.abc123.js) make immutable assets trivially safe, while explicit purge by URL or surrogate key handles the rest. Never put per-user private data on a shared edge cache.
DNS and the resolution path
DNS maps names to IP addresses through a hierarchy — root, TLD, authoritative nameservers — with recursive resolvers caching results by TTL. Two mechanisms make it a load-balancing tool:
- Anycast — the same IP is advertised from many locations; the network routes to the nearest.
- Geo-routing / GSLB — return different IPs by region to hit the nearest datacenter.
A long TTL reduces lookup latency but slows failover; a short TTL does the opposite. That trade-off is the DNS answer interviewers look for.
Reverse proxy vs API gateway
A reverse proxy (Nginx, Envoy) sits in front of servers and handles TLS termination, compression, caching, and rate limiting. An API gateway is a reverse proxy specialised for APIs: authentication, request routing, quotas, and protocol translation (REST ↔ gRPC). Both hide the backend topology from clients.
Caching: strategies, eviction, and invalidation
Caching is the highest-leverage optimisation in system design. Cache rules: the cache-aside pattern is the default; write policies decide consistency; eviction decides what leaves; invalidation decides what is correct.
| Strategy | Behaviour | Consistency | Trade-off |
|---|---|---|---|
| Cache-aside | App reads cache, loads DB on miss, fills cache | Eventual | Simple; stale until invalidated |
| Read-through | Cache itself loads from DB on miss | Eventual | Cleaner app code; cache must know DB |
| Write-through | Write cache and DB together | Strong-ish | Higher write latency |
| Write-back | Write cache now, flush to DB later | Eventual, risk of loss | Fast writes; durability risk |
| Write-around | Write DB only, bypass cache | Eventual | Avoids caching write-once data |
Eviction policies: LRU (least recently used) fits access patterns where recency predicts reuse; LFU (least frequently used) favours stable hot items; FIFO is simple but naive; TTL bounds staleness regardless.
Two failure modes to name:
- Thundering herd / cache stampede — many requests miss the same hot key at once and hammer the DB. Fix with request coalescing (single-flight), locking, or staggered TTLs.
- Hot key — one key gets disproportionate traffic. Fix by replicating the key across shards or adding a small local cache.
Cache layers, from nearest to farthest: client/browser → CDN edge → application in-process → distributed cache (Redis) → database. Cache at the layer that removes the most work, and never cache auth decisions or balances without strict invalidation.
Explore hit ratio and eviction behavior:
Cache Strategy Simulator
Not every layer needs a cache — pick the right one and measure hit ratio.
Hit ratio
78%
Avg latency
48ms
Cache hits
0
Hot keys stored in memory with TTL invalidation.
Path: Client → API → Redis → Database · Misses: 0
Week 3: Database & Messaging Systems
When one database cannot hold the data or the write rate, you keep writes correct with transactions, split the data (sharding), copy it (replication), speed up lookups (indexing), and decouple writers from readers (messaging). This week is the correctness and async backbone of every large system.
Database transactions: ACID, isolation levels, and locks
A transaction is a unit of work that either fully commits or fully rolls back. Everything about correctness in a relational database flows from the four ACID guarantees:
| Property | Guarantee | How it is implemented |
|---|---|---|
| Atomicity | All statements succeed, or none do | Undo log / rollback of partial changes |
| Consistency | Constraints hold before and after | Foreign keys, unique, check constraints |
| Isolation | Concurrent transactions do not corrupt each other | Locks + MVCC snapshots |
| Durability | A committed write survives a crash | Write-ahead log flushed to disk (fsync) |
The canonical example — moving money must never lose or create it:
BEGIN;
UPDATE accounts SET balance = balance - 100 WHERE id = 1;
UPDATE accounts SET balance = balance + 100 WHERE id = 2;
COMMIT;If the second statement fails, the first is rolled back; if the process crashes right after COMMIT, the write-ahead log replays it. Atomicity prevents half a transfer, durability prevents a lost one.
Isolation is the subtle property, because "isolated" has levels. The SQL standard defines four levels as trade-offs against three anomalies:
| Isolation level | Dirty read | Non-repeatable read | Phantom read |
|---|---|---|---|
| Read uncommitted | Yes | Yes | Yes |
| Read committed | No | Yes | Yes |
| Repeatable read | No | No | Yes* |
| Serializable | No | No | No |
*In the SQL standard, repeatable read still allows phantom reads; PostgreSQL's implementation prevents them.
The anomalies, in plain terms:
- Dirty read — you read a row another transaction wrote but has not committed; it may yet roll back.
- Non-repeatable read — you read the same row twice in one transaction and get different values because another transaction committed in between.
- Phantom read — you run the same query twice and new rows appear that another transaction inserted.
- Lost update — two transactions read a value, both add to it, and the second write overwrites the first.
How databases implement isolation:
- Two-phase locking (2PL) — readers and writers acquire locks and hold them until commit. Correct and simple, but blocking and deadlock-prone.
- MVCC (multi-version concurrency control) — the database keeps multiple versions of a row; a reader sees a consistent snapshot and never blocks a writer. PostgreSQL, Oracle, and MySQL (InnoDB) use MVCC. Stronger isolation means more conflict detection, so
SERIALIZABLEtransactions can fail and must be retried.
The classic lost update, and three fixes:
-- Both transactions read 100, both write 110 — one +10 increment is lost:
SELECT balance FROM accounts WHERE id = 1;
UPDATE accounts SET balance = 100 + 10 WHERE id = 1;
-- Fix 1: atomic read-modify-write in one statement
UPDATE accounts SET balance = balance + 10 WHERE id = 1;
-- Fix 2: pessimistic lock — the row is held until commit
SELECT balance FROM accounts WHERE id = 1 FOR UPDATE;
-- Fix 3: optimistic version check — fails the second writer
UPDATE accounts SET balance = balance + 10, version = version + 1
WHERE id = 1 AND version = 5;Sharding: splitting data across nodes
Sharding (horizontal partitioning) splits a large dataset across multiple machines by a shard key, so no single node holds everything and the write rate spreads across nodes. It removes the vertical ceiling — one machine's CPU, RAM, and disk — but introduces cross-shard complexity.
| Method | How it works | Strength | Weakness |
|---|---|---|---|
| Range | Contiguous key ranges per shard | Great for range scans | Hot spots on sequential keys |
| Hash | hash(key) → shard | Even spread | Range queries scatter; resharding costly |
| Directory | Lookup service maps key → shard | Flexible, easy rebalance | Extra hop; lookup is a SPOF |
Worked example. 1 billion users at ~1 KB per row ≈ 1 TB — too large and too hot for one node. Choose user_id as the shard key and 4 shards:
- Range: shard 0 = ids 0–250M, shard 1 = 250M–500M, and so on. Reads by id range are efficient, but every new signup lands on the last shard — a write hotspot, because ids grow monotonically.
- Hash:
shard = hash(user_id) % 4. Requests spread evenly, but adding a fifth shard changes the modulus and remaps almost every key. - Consistent hashing (the resharding fix): place shard nodes (with virtual nodes) on a hash ring; adding a fifth shard steals only the slice between it and its neighbours — roughly
1/Nof keys move, not all.
function shardFor(key: string, shards: number): number {
let h = 0
for (const ch of key) h = (h * 31 + ch.charCodeAt(0)) >>> 0
return h % shards
}Routing. Something must map a key to a shard: the application computes it, or a proxy/router such as Vitess or Citus does. The alternative is a directory service — a lookup table of key → shard — which is flexible and easy to rebalance, but adds a network hop and must itself be highly available.
The shard-key decision is the whole design. It must (1) spread load evenly, (2) avoid hot partitions, and (3) keep data for common queries on one shard. A bad key shows up as the celebrity / hot-partition problem — choosing country when one country dominates, or a single viral user_id saturating one shard. Fixes are a composite key (user_id + bucket) or salting a hot key across several shards.
Cross-shard queries are the tax. A query like "all orders in the last hour" on user-sharded data must hit every shard and merge results (scatter-gather); cross-shard joins and transactions are the same problem. This is why the shard key must match the dominant access pattern, not just distribute writes.
Replication: copies for durability and reads
Replication keeps copies of data on multiple nodes for durability, read scaling, and failover. The trade-off is always consistency vs latency.
- Leader-follower — one node accepts writes, followers replicate and serve reads. Simple and read-scalable; fails over by promoting a follower. Asynchronous replication introduces replication lag (a follower may serve stale data).
- Multi-leader — multiple nodes accept writes (e.g. multi-region). Fast local writes, but concurrent edits create conflicts that need resolution (last-write-wins, version vectors, CRDTs).
- Leaderless / quorum — clients write to
Wreplicas and read fromR; withW + R > Nthe read set overlaps the write set, giving strong-ish consistency tunable per request (Dynamo-style).
Indexing: B-Tree vs LSM tree
An index is a data structure that turns a full scan into a lookup. Two families dominate:
- B-Tree / B+Tree — balanced, sorted tree; excellent for point lookups and range scans; updates in place. The default in relational engines (PostgreSQL, MySQL). Good read performance, more random writes.
- LSM tree — writes go to an in-memory table and are flushed to immutable sorted files, then compacted in the background. Write-optimised, sequential I/O; the default in many NoSQL stores (Cassandra, RocksDB). Reads may check several levels (amplification), mitigated by bloom filters and compaction.
The trade to name: B-Trees favour read-heavy workloads with range queries; LSM trees favour write-heavy workloads. Both are examples of trading write amplification against read amplification.
Kafka: the distributed log
Kafka is an append-only, partitioned, replayable log, not a traditional queue. Understanding five concepts is enough to reason about it:
| Concept | Meaning |
|---|---|
| Topic | A named stream of records, split into partitions |
| Partition | An ordered, immutable log; the unit of parallelism |
| Offset | A consumer's position within a partition |
| Consumer group | Consumers that share partitions, one owner per partition |
| Retention | How long records are kept (time or size), regardless of reads |
The critical invariant: ordering is guaranteed within a partition, never across partitions. If ordering matters for an entity, partition by that entity's key (all events for one order go to one partition). Consumer groups give you horizontal scaling of consumption while preserving per-partition order.
Delivery semantics: at-most-once (may lose), at-least-once (may duplicate — the practical default), and exactly-once (achieved with idempotent producers/consumers and transactions, and only within Kafka's boundary). In practice, exactly-once means at-least-once delivery plus idempotent processing.
SQS vs Kafka, and event-driven fan-out
SQS is a managed queue: a message is consumed once and deleted — ideal for task distribution and decoupling. Kafka is a log: many independent consumer groups can replay the same events — ideal for event streaming, audit, and multiple downstream consumers. Choose Kafka when you need replay or several consumers; choose a queue when you need simple work distribution.
Week 4: High-Traffic & Search Systems
These three problems test the classic levers: read-through caching, atomic counters, and prefix data structures.
Problem: URL Shortener (🎯 Asked at Flipkart)
What it tests: read-heavy design, unique key generation, and separating the fast redirect path from analytics.
Requirements. Shorten a long URL, redirect on access, keep latency low, and track clicks without slowing the redirect.
Key generation. Two approaches:
- Counter + Base62 — encode a monotonically increasing integer into
[a-zA-Z0-9]. Short, collision-free, but the counter must be distributed. Use a pre-allocated ID range per node (a "range allocator") or a central generator, so nodes never collide without a lock per request. 6 Base62 chars ≈ 56 billion URLs. - Hash + collision handling — hash the long URL (MD5/SHA), truncate, Base62-encode, and check for collisions. Deterministic (same URL → same key, useful for dedup) but requires a uniqueness check.
Storage is a simple key-value mapping shortKey → longUrl, ownerId, createdAt, which is small and read-heavy — perfect for a KV store with an edge cache.
Redirect semantics: 301 (permanent) is cached by browsers and reduces load but hides repeat clicks; 302 (temporary) lets every request hit you and be counted. Prefer 302 when analytics matter.
Analytics: publish a click event to Kafka and aggregate asynchronously. The redirect must never block on a write.
Write path: client → API → ID generator → KV store → return short URL
Read path: client → CDN/edge cache → KV store → 302 redirectTrade-offs: counter vs hash (simplicity vs determinism); where the cache lives; how you prevent hot-key pressure on a single viral short link.
Problem: Distributed Rate Limiter (🎯 Asked at Razorpay)
What it tests: atomic shared state and the difference between bursty and strict limiting.
Algorithms, in increasing precision:
| Algorithm | Idea | Characteristic |
|---|---|---|
| Fixed window | Count per fixed interval | Simple; burst at window edges |
| Sliding window log | Timestamps of each request | Exact; memory heavy |
| Sliding window counter | Weighted blend of current + previous window | Approximate; cheap |
| Token bucket | Tokens refill at a rate; each request takes one | Allows controlled bursts |
| Leaky bucket | Requests drain at a fixed rate | Smooth output rate |
Token bucket allows a burst up to the bucket size while enforcing a long-run rate — the usual default. In a distributed setting, the counter lives in Redis and every check must be atomic:
allowed = INCR key
if allowed == 1: EXPIRE key window
if allowed > limit: rejectDecide the scope — per user, per API key, per IP — and return 429 with Retry-After. Scope, storage, and precision are the trade-offs.
Problem: Search Autocomplete (🎯 Asked at Google)
What it tests: prefix data structures and precomputation for read speed.
A Trie stores strings by shared prefixes: each node represents a prefix, and its children extend it. For autocomplete, each node caches the top-k completions by score, so a query is just a walk of prefix length nodes and a return of the cached list — O(prefix length), independent of corpus size.
At scale: precompute top-k at write time (updating a popular prefix must not cost a full scan), shard the trie by prefix range so a query hits one shard, cache the hottest prefixes at the edge, and blend personalisation and trending scores. This is the same Trie idea used in the DSA track for prefix matching.
Reuse the contract explorer when defining the autocomplete and query APIs:
API Contract Explorer
Compare versioning and request-response shape for the same operations.
Request
Query: ?cursor=abc&limit=20
Response
{
"data": [{ "id": "usr_1", "fullName": "Asha" }],
"nextCursor": "def"
}Week 5: Feed & Streaming Systems
Feeds test fan-out; streaming tests bandwidth and the CDN; messaging tests realtime state. All three test how you handle read-heavy, high-fan-out traffic.
Problem: Twitter / Instagram Feed (🎯 Asked at Meesho)
What it tests: the fan-out trade-off between write cost and read cost.
Fan-out is the central decision:
| Model | Behaviour | Trade-off |
|---|---|---|
| Push | Write post → fan out into every follower's feed | Fast reads; expensive for celebrities |
| Pull | Build feed at read time from followees | Cheap writes; slow reads |
| Hybrid | Push for normal users, pull for celebrities | Best of both; more complexity |
The celebrity problem is why hybrid wins: for a user with 50M followers, pushing every post into 50M feeds is untenable, so those accounts are merged at read time. Store feed ids in a cache, hydrate content lazily, rank with a scoring service, and accept eventual consistency. This is one of the few problems where you explicitly design for the tail.
Problem: YouTube / Netflix (🎯 Asked at Netflix)
What it tests: bandwidth-bound design and asynchronous processing.
Upload → chunked/segmented storage → transcoding pipeline (multiple resolutions and bitrates in parallel) → manifest files → CDN edge delivery with adaptive bitrate (HLS/DASH). The transcoding is the long pole: it is asynchronous, parallelised, and fault-tolerant (retry a failed segment, not the whole video). Streaming itself is bandwidth-bound, so the CDN does the heavy lifting and the origin rarely serves bytes directly.
Problem: WhatsApp / Messaging (🎯 Asked at Microsoft)
What it tests: realtime connections, ordering, and at-least-once delivery.
- Persistent connections (WebSocket) with a presence service and a message queue per recipient.
- Delivery receipts require message status:
sent → delivered → read, which needs per-message state and push-back to the sender. - Offline delivery: store undelivered messages and push on reconnect; order per conversation must be preserved (a per-conversation sequence number).
- Deduplication: messages can arrive twice on reconnect — carry a client message id.
Resilience patterns
When services depend on each other, failures cascade. Name these patterns:
- Circuit breaker — after repeated failures, stop calling a dependency entirely for a while, then probe; prevents piling requests onto a dead service.
- Bulkhead — isolate resource pools per dependency so one slow service cannot exhaust all threads or connections.
- Retry with exponential backoff + jitter — retry transient failures, but randomise the delay so retries do not synchronise into a thundering herd.
- Timeout + fallback — every remote call needs a deadline and a degraded response.
Week 6: Storage & Notification Systems
These problems test whether you can design storage primitives and pick the right realtime transport.
Problem: Key-Value Store (🎯 Asked at Amazon)
What it tests: distributed storage internals — partitioning, replication, and conflict handling.
A Dynamo-style store combines four ideas:
- Partitioning via consistent hashing so keys spread across nodes.
- Replication with quorum reads/writes (
W + R > N) for tunable consistency. - Conflict handling with vector clocks or version vectors; on concurrent writes, resolve with last-write-wins or application logic.
- Availability tricks: hinted handoff (a node temporarily holds writes for a down peer) and read repair (fix stale replicas during reads).
Redis is the in-memory, single-node-ish extreme; DynamoDB is the managed, durable, horizontally-scaled extreme.
Problem: Distributed Cache (🎯 Asked at Uber)
What it tests: eviction, partitioning, and the failure modes of caches.
Three axes: eviction (LRU/LFU plus TTL), partitioning (consistent hashing across cache nodes), and consistency (invalidate or update on write). The genuinely hard parts are cache stampede (many misses on one key at once — fix with single-flight and TTL jitter) and hot keys (replicate the key). A cache is not a source of truth; treat any hit as potentially stale and design the read path to tolerate it.
Problem: Notification System (🎯 Asked at Swiggy)
What it tests: pluggable channels, fan-out, and delivery guarantees.
Channels (push/SMS/email) sit behind one adapter interface; a priority queue handles fan-out; retries use backoff with a dead-letter queue for permanent failures; deduplication and user preferences prevent spam; and a delivery-status pipeline (a topic per channel) feeds observability. The design idiom is Adapter + Observer + queue.
Realtime transport: how to choose
| Transport | Direction | Mechanics | Use case |
|---|---|---|---|
| Short polling | Client pulls | Client asks repeatedly | Simple, wasteful |
| Long polling | Client holds | Server holds the request until data or timeout | Near-realtime, HTTP-only |
| Server-Sent (SSE) | Server → client | One long HTTP response streams events | Feeds, LLM token streaming |
| WebSocket | Bidirectional | Persistent full-duplex socket | Chat, collaboration, games |
The choice follows from direction and scale: one-way server pushes fit SSE; interactive two-way fits WebSocket; constrained clients fall back to long polling.
Week 7: Location & Scheduling Systems
These test spatial indexing and delayed, reliable execution — two problems with non-obvious data structures.
Problem: Uber (🎯 Asked at Zomato)
What it tests: geo-indexing and matching under realtime constraints.
- Geo-indexing: encode driver locations so "find nearby drivers" is a local query, not a global scan. Options: geohash (string prefix encodes a cell — simple, but edge cells are awkward), QuadTree (recursive split; adapts to density), H3 (near-uniform hexagonal cells). Store drivers in a spatial index keyed by cell.
- Matching: candidate drivers from nearby cells → rank by ETA (not just distance) → assign with a lock or lease so two riders cannot get the same driver.
- Surge pricing: compute demand vs supply per cell per time window and apply a multiplier.
- Tracking: drivers publish location over WebSocket to a location service; riders subscribe. Location updates are high-volume — batch and throttle.
Problem: Google Maps (🎯 Asked at Google)
What it tests: graph algorithms at scale plus heavy caching.
Routing is a weighted shortest-path problem over a road network. Plain Dijkstra is too slow on continental graphs, so production systems use A* with a heuristic and contraction hierarchies (precomputed shortcuts) to answer in milliseconds. Map tiles are pre-rendered and cached at the CDN; geo-search and reverse geocoding reuse the same spatial indexes as the ride-sharing case.
Problem: Task Scheduler (🎯 Asked at Spotify)
What it tests: priority execution, delayed jobs, and reliable delivery.
- Worker pool: a bounded set of workers pulling from a blocking queue.
- Delayed jobs: a min-heap or timing wheel keyed by
runAt; a timer promotes due jobs into the ready queue. - Reliability: leases and visibility timeouts so a job whose worker crashed is retried, without two workers running it at once.
- Cancellation: keep a
Set<jobId>of cancelled ids checked before execution. - Failure: retries with exponential backoff and a dead-letter queue after the max attempts.
Week 8: Payments & Reliability
Payments test correctness and failure handling above all — availability matters less than never losing or double-spending money.
Problem: Payment Gateway / UPI (🎯 Asked at PhonePe)
What it tests: idempotency, distributed transactions, and reconciliation.
- Idempotency keys so a retried request returns the original result instead of charging twice. The client sends a key; the server stores
key → resultand replays it. - Distributed transactions: 2PC gives strong atomicity across a few participants but blocks on coordinator failure; Saga models a long flow as a sequence of local transactions with compensating actions for rollback — better for microservices. Know when each fits.
- Ledger: double-entry, append-only, entries sum to zero. Balances are derived, never edited in place.
- Exactly-once in practice = at-least-once delivery + idempotent consumers.
Problem: Leaderboard / Distributed Counters (🎯 Asked at Netflix)
What it tests: realtime ranking at scale.
Use Redis sorted sets (ZADD/ZINCRBY) for O(log n) inserts and O(log n + k) range reads by score — exactly what a leaderboard needs. For counters that must scale beyond one node, shard and merge, or use a CRDT-style counter (G-Counter) when eventual consistency is acceptable; exact global counts require coordination.
Observability
What it tests: whether you can operate what you design. Three pillars:
- Logs — discrete events, high cardinality, useful for debugging.
- Metrics — aggregated numbers over time (counters, gauges, histograms), cheap and alertable.
- Traces — the causal path of a request across services, with span timings.
Discuss cardinality control (labels explode fast), sampling (you cannot store every trace), and alerting on symptoms (latency, error rate) rather than causes.
Week 9: AI System Design — ChatGPT & Document Q&A
AI systems reuse the building blocks of weeks 1–8 and add new constraints: streaming, token budgets, inference cost, and untrusted retrieved content.
Problem: ChatGPT-style Assistant
What it tests: streaming, stateful conversations, and inference economics.
| Concern | Design | Why it matters |
|---|---|---|
| Conversations | Store sessions + messages; paginate history | Long histories exceed context windows |
| Streaming | SSE tokens to the client; abort on disconnect | Perceived latency is time-to-first-token |
| Routing | Model/region selection by cost, latency, capability | Cheapest model that satisfies the query |
| Backpressure | Bounded queues + batching when inference is the bottleneck | GPUs are the scarce resource |
Inference economics matter: batch requests to raise GPU utilisation, measure time-to-first-token (user-perceived latency) separately from total generation time, and cache prompts/prefixes where safe. Conversation storage is a normal write-heavy table, but the context builder — summarise old turns, keep recent turns intact, never drop system instructions — is where the design pressure lives.
Problem: Document Q&A with RAG
What it tests: a retrieval pipeline and the security of untrusted context.
Ingestion: load → chunk → embed → index. Retrieval: embed query → vector search top-k → rerank → generate with citations. The hard parts are chunking quality (too big dilutes relevance, too small loses context), permission-aware retrieval, and evaluation (offline sets for recall and answer quality).
Week 10: AI Agents, Model Gateways & Mock Interviews
The final week moves from single-model apps to orchestration, and closes with practice.
Problem: AI Agent Platform
What it tests: durable, multi-step execution with tools.
Orchestration runs a loop: plan → call tool → observe → repeat, until done. The design demands:
- Tool execution with validation, timeouts, and bounded concurrency.
- Durable state — persist the run as an event log so a crash resumes instead of restarting (a Saga-shaped problem).
- Approvals — destructive actions require explicit human consent before execution.
Model the agent run as a state machine (PENDING → RUNNING → WAITING_APPROVAL → SUCCEEDED | FAILED | CANCELLED) with an append-only event log for replay and audit.
Problem: Multi-Model Gateway
What it tests: abstraction over providers with cost and reliability controls.
Route by capability, cost, or health; fall back across providers when one degrades; enforce per-tenant quotas; cache identical requests; and track cost per request/token as a first-class metric. This is the Adapter + Strategy + Circuit Breaker pattern stack applied to models.
Evaluation, tracing, and guardrails
- Offline eval sets for retrieval relevance (recall@k) and answer quality (human or model-graded).
- Tracing to attribute latency and cost across retrieval, reranking, and generation — you cannot optimise what you cannot attribute.
- Guardrails for toxicity, PII, and prompt injection.
Mock interviews
Run 2–3 mocks with company-specific framing, then a 1:1 feedback pass against the Week 1 framework. The goal is not a perfect design — it is a defensible one delivered fluently.
Interview Reflection
Which approach do you use in design rounds?
Full case studies list (19)
| # | Case Study | Focus | Asked at |
|---|---|---|---|
| 1 | URL Shortener | Hashing, redirection, analytics | Flipkart |
| 2 | Twitter / Instagram Feed | Newsfeed, fan-out | Meesho |
| 3 | YouTube / Netflix | Video upload, streaming, CDN | Netflix |
| 4 | WhatsApp / Messaging | Realtime, delivery receipts | Microsoft |
| 5 | Uber / Ride Sharing | Geo-indexing, matching, surge pricing | Zomato |
| 6 | Google Maps | Routing, geo-search, map tiles | |
| 7 | Notification System | Push, SMS, email at scale | Swiggy |
| 8 | Rate Limiter (distributed) | Token Bucket, Sliding Window | Razorpay |
| 9 | Key-Value Store | DynamoDB / Redis semantics | Amazon |
| 10 | Distributed Cache | Eviction, consistency, partitioning | Uber |
| 11 | Search Autocomplete | Trie at scale, typeahead service | |
| 12 | Payment / UPI-like Gateway | Idempotency, Saga | PhonePe |
| 13 | Distributed Task Scheduler | Priority queues, retries, delayed jobs | Spotify |
| 14 | Leaderboard / Distributed Counters | Real-time ranking at scale | Netflix |
| 15 | Logging & Monitoring | Log aggregation, alerting | Microsoft |
| 16 | ChatGPT-style AI Assistant | Streaming, conversations, model serving | — |
| 17 | Document Q&A with RAG | Ingestion, retrieval, citations | — |
| 18 | AI Agent Platform | Workflows, tools, durable execution | — |
| 19 | Multi-Model AI Gateway | Routing, quotas, fallback, cost controls | — |
Interview discipline
High level design interviews test structured reasoning, not memorized architectures. Clarify requirements, estimate scale, identify the bottleneck, then propose components with explicit trade-offs — and always name what you would change at 10×.
Interview Reflection
Which approach do you use in design rounds?
Subscribe to my newsletter
Stay up to date and get notified when I share new contents.
No spam ever, unsubscribe anytime