Kazi Rahamatullah
Kazi Rahamatullah
AboutProjectsBlogContact
UsesBooks
ResumeView CV
Kazi Rahamatullah

© Copyright 2026 Kazi Rahamatullah

AboutProjectsBlogBooksUses
Twitter/XGitHubProduct HuntCodeSandbox
Back to Blog
High Level DesignHLDSystem DesignScalabilityDistributed SystemsCachingKafkaInterview Preparation

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.

Oct 2, 202643 min read

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:

  1. Clarify — functional and non-functional requirements.
  2. Estimate — QPS, storage, bandwidth from first principles.
  3. Model — API and data model.
  4. Design — components, each with a reason to exist.
  5. 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

#TopicDescription
1Foundations & Mental ModelsFramework, capacity estimation, latency numbers, CAP, consistency, storage.
2Core Infrastructure Building BlocksLoad balancers, consistent hashing, CDN/DNS, caching, Redis.
3Database & Messaging SystemsTransactions, sharding, replication, indexing, Kafka and event-driven fan-out.
4High-Traffic & Search SystemsURL Shortener, distributed rate limiter, search autocomplete.
5Feed & Streaming SystemsTwitter/Instagram, YouTube/Netflix, WhatsApp, resilience patterns.
6Storage & Notification SystemsKey-value store, distributed cache, notifications, realtime transport.
7Location & Scheduling SystemsUber, Google Maps, task scheduler, geo-indexing.
8Payments & ReliabilityPayment gateway, leaderboards, observability, distributed transactions.
9AI System Design: ChatGPT & Q&AStreaming assistants, RAG, inference economics, tenant isolation.
10AI Agents, Gateways & Mock InterviewsAgent platforms, multi-model gateways, evaluation, mocks.
11Full 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:

PhaseTimeWhat you produce
Requirements5 minFunctional scope + non-functional targets (latency, scale)
Estimation5 minQPS, storage, bandwidth, cache size
High-level design20 minComponents, APIs, data model, data flow
Deep dive15 minOne or two hard parts: consistency, hotspots, failure
Wrap-up5 minBottlenecks, 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-functionalQuestion to askArchitectural consequence
Latencyp99 read under 100 ms?Cache at the edge + read replicas
Availability99.99% or 99.9%?Redundancy vs a single leader
ConsistencyCan a feed be 5 seconds stale?Eventual consistency is allowed
DurabilityCan we lose a payment?Synchronous replication + ledger
Scale100 QPS or 100,000 QPS?One DB vs sharding + a queue

Note

Write the non-functional targets down and refer back to them. When you later propose a cache, you can say "this protects the p99 read target" — that is the sentence interviewers reward.

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:

  1. Users — total users, then daily active users (DAU).
  2. Actions — requests per user per day.
  3. QPS — (DAU × actions) / 86,400, then multiply by 2–3 for peak.
  4. Storage — writes per day × record size × retention.
  5. Bandwidth — record size × QPS (and the read:write ratio).
  6. 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:

text
DAU           = 200M × 10%          = 20M
Requests/day  = 20M × 20            = 400M
Average QPS   = 400M / 86,400       ≈ 4,600 QPS
Peak QPS      = 2–3× average        ≈ 12,000 QPS

Storage 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:

QuantityValue
Seconds per day86,400 (≈ 100,000 for mental math)
1 million requests/day≈ 12 QPS
1 billion requests/day≈ 12,000 QPS
Powers of two2^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

Performance

These latency numbers explain most architecture decisions. Disk (~10 ms) is ~100,000× slower than memory (~100 ns), which is why caching works. A cross-continent hop (~150 ms) is slower than an entire in-memory computation, which is why CDNs exist.

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.
ChoiceUnder partitionExamples
CPReject writes/reads it cannot make consistentHBase, ZooKeeper, Spanner-like
APKeep serving, accept divergence, reconcile laterCassandra, DynamoDB, DNS

Note

Say "CP or AP under partition". Outside a partition, a system can be both consistent and available. Treating CAP as a permanent three-way switch is the most common way to lose points.

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.

ModelGuaranteeTypical use
StrongReads always see the latest committed writePayments, inventory
LinearizableStrong + operations appear in real-time global orderLeader election, distributed locks
CausalWrites that are causally related are seen in orderComment threads, collaboration
Read-your-writesA user sees their own writesProfiles, posting then refreshing
Monotonic readsA user never sees time go backwardsFeeds, timelines
EventualReplicas converge given no new writesLikes, counters, view counts

Interview Answer

When asked "what consistency do you need?", answer with the weakest model that satisfies the use case, and name the cost. "Feeds can be eventually consistent within a second; balances cannot." Weaker consistency buys availability and latency — that trade is the whole game.

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.

AspectSQLNoSQL
SchemaFixed, migration-drivenFlexible / schema-on-read
ScalingVertical, then read replicas/shardingHorizontal by design
TransactionsStrong ACID across tablesOften single-item or eventual
JoinsNative and efficientUsually denormalised / application-side
Best forRelationships, ledger, reportingHigh 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

users → orders (1 : N)

Back to index


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:

LayerOperates onSeesStrengthsTrade-off
L4TCP/UDP connectionsIP, portVery fast, protocol-agnosticCannot route by URL/header
L7HTTP requestsPath, headers, cookiesSmart routing, TLS, retries, stickyMore CPU; terminates requests

Common algorithms and when each fits:

AlgorithmHow it worksUse when
Round robinRotate evenlyHomogeneous servers
Least connectionsPick the least busyVariable request cost
IP/consistent hashSame key → same serverCaching, session affinity
WeightedBigger servers get more trafficMixed 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:

  1. The client opens a TCP connection to the balancer's virtual IP.
  2. The balancer picks a backend (e.g. least connections) and forwards packets — often by NAT, without terminating the connection.
  3. Because the connection is pinned to that backend for its lifetime, every packet of that connection goes to the same server.
  4. 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:

  1. It terminates the client's TCP and TLS, decrypts, and reads the full HTTP request.
  2. It routes by path, host, header, or cookie — and can send each request to a different backend.
  3. It can retry idempotent requests, rewrite headers, add authentication, compress, and pool connections to backends.
  4. The cost is CPU: parsing HTTP and doing TLS handshakes is far more expensive than forwarding packets.
text
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)
PathL4L7
Terminates TLSPass-through (cannot see content)Yes (can inspect and modify requests)
Routing unitConnectionHTTP request
Can retryNo (cannot tell a request from a packet)Yes, for idempotent requests
Typical useRaw TCP, databases, extreme throughputHTTP 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:

  1. 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.
  2. 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.
  3. Hit. The edge serves the object directly — single-digit milliseconds, with no origin involvement.
  4. 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.
  5. Revalidation. When a TTL expires, the edge asks the origin If-None-Match/If-Modified-Since; a 304 Not Modified refreshes the TTL without re-downloading the body.
  6. Invalidation. An explicit purge removes an object by URL or surrogate key; versioned filenames make old entries harmless.
text
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 → HIT

Pull 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:

http
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 edge

The 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.

Performance

An edge hit is ~5–20 ms; a miss that crosses an ocean to origin adds ~150 ms RTT plus origin compute. A high CDN hit ratio is one of the cheapest latency wins in system design — the same object is paid for once per edge instead of once per user.

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.

StrategyBehaviourConsistencyTrade-off
Cache-asideApp reads cache, loads DB on miss, fills cacheEventualSimple; stale until invalidated
Read-throughCache itself loads from DB on missEventualCleaner app code; cache must know DB
Write-throughWrite cache and DB togetherStrong-ishHigher write latency
Write-backWrite cache now, flush to DB laterEventual, risk of lossFast writes; durability risk
Write-aroundWrite DB only, bypass cacheEventualAvoids 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

Performance

The single most effective lever is cache hit ratio at the right layer: CDN for static assets, Redis for computed/hot reads. A 99% hit ratio means the database sees 1% of traffic — that is the entire point.

Back to index


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:

PropertyGuaranteeHow it is implemented
AtomicityAll statements succeed, or none doUndo log / rollback of partial changes
ConsistencyConstraints hold before and afterForeign keys, unique, check constraints
IsolationConcurrent transactions do not corrupt each otherLocks + MVCC snapshots
DurabilityA committed write survives a crashWrite-ahead log flushed to disk (fsync)

The canonical example — moving money must never lose or create it:

sql
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 levelDirty readNon-repeatable readPhantom read
Read uncommittedYesYesYes
Read committedNoYesYes
Repeatable readNoNoYes*
SerializableNoNoNo

*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 SERIALIZABLE transactions can fail and must be retried.

The classic lost update, and three fixes:

sql
-- 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;

Note

A single-node transaction gives you ACID. Once work spans services or databases, that guarantee is gone — reach for 2PC (strong, blocking, few participants) or Saga (compensating actions, many participants). Payments reliability is covered in Week 8.

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.

MethodHow it worksStrengthWeakness
RangeContiguous key ranges per shardGreat for range scansHot spots on sequential keys
Hashhash(key) → shardEven spreadRange queries scatter; resharding costly
DirectoryLookup service maps key → shardFlexible, easy rebalanceExtra 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/N of keys move, not all.
typescript
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.

Performance

Resharding without downtime follows the migration pattern: dual-write to the old and new mapping, backfill historical rows, verify, then cut reads over and retire the old mapping. Plan for it before you need it — resharding is a data migration, not a config change.

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 W replicas and read from R; with W + R > N the read set overlaps the write set, giving strong-ish consistency tunable per request (Dynamo-style).

Note

Synchronous replication guarantees the write survives the failure of the leader, at the cost of write latency. Asynchronous replication is fast but can lose the most recent writes on failover. State which one you are choosing and why — the answer differs for a likes table and a payments ledger.

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:

ConceptMeaning
TopicA named stream of records, split into partitions
PartitionAn ordered, immutable log; the unit of parallelism
OffsetA consumer's position within a partition
Consumer groupConsumers that share partitions, one owner per partition
RetentionHow 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.

Note

Async fan-out decouples producers from consumers: an order event can feed inventory, notifications, analytics, and fraud detection independently. The cost is eventual consistency and the need for idempotent consumers — every consumer must tolerate seeing the same message twice.

Back to index


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.

text
Write path:  client → API → ID generator → KV store → return short URL
Read path:   client → CDN/edge cache → KV store → 302 redirect

Trade-offs: counter vs hash (simplicity vs determinism); where the cache lives; how you prevent hot-key pressure on a single viral short link.

Interview Answer

URL shortener interview flow: estimate read:write ratio (heavily read-skewed) → Base62 with a counter → KV store → cache hot keys at the edge → 302 when you need analytics → state why the redirect path must never block on writes.

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:

AlgorithmIdeaCharacteristic
Fixed windowCount per fixed intervalSimple; burst at window edges
Sliding window logTimestamps of each requestExact; memory heavy
Sliding window counterWeighted blend of current + previous windowApproximate; cheap
Token bucketTokens refill at a rate; each request takes oneAllows controlled bursts
Leaky bucketRequests drain at a fixed rateSmooth 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:

text
allowed = INCR key
if allowed == 1: EXPIRE key window
if allowed > limit: reject

Warning

INCR then EXPIRE are two commands — a crash between them can leave a key with no TTL, leaking memory. Use a Lua script (or SET key value EX ttl NX plus INCR) to make the check-and-expire atomic, and be aware that a strict global limit needs a single Redis instance or a consensus approach.

Decide 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"
}

Back to index


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:

ModelBehaviourTrade-off
PushWrite post → fan out into every follower's feedFast reads; expensive for celebrities
PullBuild feed at read time from followeesCheap writes; slow reads
HybridPush for normal users, pull for celebritiesBest 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.

Back to index


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

TransportDirectionMechanicsUse case
Short pollingClient pullsClient asks repeatedlySimple, wasteful
Long pollingClient holdsServer holds the request until data or timeoutNear-realtime, HTTP-only
Server-Sent (SSE)Server → clientOne long HTTP response streams eventsFeeds, LLM token streaming
WebSocketBidirectionalPersistent full-duplex socketChat, 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.

Back to index


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.

Note

Geohashing trades precision for simplicity; QuadTrees adapt to density; H3 gives near-uniform hexagonal cells. Name the trade-off rather than picking one blindly — the right choice depends on whether your load is urban-dense or spread out.

Back to index


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 → result and 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.

Warning

Never hold a database transaction open across an external provider call. Persist a pending intent, call the provider, then reconcile asynchronously and rely on idempotency for retries. Long-held locks across a network call are how systems deadlock under load — this is the reliability answer interviewers probe for.

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.

Back to index


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.

ConcernDesignWhy it matters
ConversationsStore sessions + messages; paginate historyLong histories exceed context windows
StreamingSSE tokens to the client; abort on disconnectPerceived latency is time-to-first-token
RoutingModel/region selection by cost, latency, capabilityCheapest model that satisfies the query
BackpressureBounded queues + batching when inference is the bottleneckGPUs 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).

Warning

Retrieved content is untrusted input and a direct prompt-injection vector — a malicious document can contain "ignore previous instructions". Defend by separating instructions from data, validating tool calls, and enforcing tenant isolation at the retrieval layer so a user can never retrieve another tenant's chunks. Filter by permission before the vector search, not after.

Back to index


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?

Interview Answer

The HLD signal I look for: you state the requirement that makes a component necessary, and the failure mode that makes it risky. "I add a cache here for read QPS, and I accept eventual staleness with a 60s TTL" is a senior sentence. "I add a cache because caches are faster" is not.

Back to index


Full case studies list (19)

#Case StudyFocusAsked at
1URL ShortenerHashing, redirection, analyticsFlipkart
2Twitter / Instagram FeedNewsfeed, fan-outMeesho
3YouTube / NetflixVideo upload, streaming, CDNNetflix
4WhatsApp / MessagingRealtime, delivery receiptsMicrosoft
5Uber / Ride SharingGeo-indexing, matching, surge pricingZomato
6Google MapsRouting, geo-search, map tilesGoogle
7Notification SystemPush, SMS, email at scaleSwiggy
8Rate Limiter (distributed)Token Bucket, Sliding WindowRazorpay
9Key-Value StoreDynamoDB / Redis semanticsAmazon
10Distributed CacheEviction, consistency, partitioningUber
11Search AutocompleteTrie at scale, typeahead serviceGoogle
12Payment / UPI-like GatewayIdempotency, SagaPhonePe
13Distributed Task SchedulerPriority queues, retries, delayed jobsSpotify
14Leaderboard / Distributed CountersReal-time ranking at scaleNetflix
15Logging & MonitoringLog aggregation, alertingMicrosoft
16ChatGPT-style AI AssistantStreaming, conversations, model serving—
17Document Q&A with RAGIngestion, retrieval, citations—
18AI Agent PlatformWorkflows, tools, durable execution—
19Multi-Model AI GatewayRouting, quotas, fallback, cost controls—

Back to index


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?

Interview Answer

What caught me in my last design interview: estimation is not decoration. The moment you write down peak QPS and storage growth, the correct architecture narrows itself — you stop guessing and start defending choices. Candidates who skip the numbers end up debating components instead of constraints.

Back to index


Share this article

XLinkedInFacebook
Kazi Rahamatullah

Written by

Kazi Rahamatullah

FullStack Developer

X / TwitterGitHubLinkedIn

Subscribe to my newsletter

Stay up to date and get notified when I share new contents.

No spam ever, unsubscribe anytime