System design interviews
How the system-design round works and how to run it: expectations by level, a minute-by-minute framework, estimation drills, how to drive the conversation, and nine worked designs in the same compact 7-step shape. The concepts behind every decision (caching, replication, quorums, queues, Bloom filters…) are in System design concepts; numbers and production tables are in System design, which also has the worked URL shortener. Coding rounds: Coding interviews.
Format & expectations
A 45–60 minute conversation around a vague prompt ("design Instagram"). There is no single right answer: you're scored on how you scope, decide, justify and adapt. Usually on a virtual whiteboard (Excalidraw, Miro, a shared doc) or a real one.
| Level | What's asked | Depth expected | You pass by |
|---|---|---|---|
| New grad / junior | often no design round; sometimes API or object design, or a small service | a working single-region design: clients, LB, stateless service, DB, cache; a sensible API and schema | clear requirements, correct basics, taking hints well; the interviewer may steer |
| Mid-level | the full round | cover the whole framework end to end; one solid deep dive, possibly prompted | knowing what each standard component does and when to use it; no big gaps |
| Senior | the full round, less guidance | move quickly through the basics; lead 2–3 deep dives unprompted; numbers behind decisions; failure modes | trade-offs stated before being asked, alternatives weighed, the hard part found early |
| Staff+ | ambiguous or open-ended prompts, sometimes "design the platform for…" | scope and prioritize; evolution over time, migrations, cost, operations, multi-team ownership | pushing back on requirements, simplifying, depth where it matters and speed elsewhere, judgment from experience |
Variants: product design (design Twitter), infrastructure design (design a rate limiter, a KV store), ML system design (design a recommendation pipeline), and architecture deep dive into a system you built. The framework below works for all of them.
The framework
| Step | 45 min | 60 min | Produce | Watch for |
|---|---|---|---|---|
| 1. Requirements | 5 | 7 | 3–5 functional features, non-functional numbers (scale, latency, availability, consistency), what's out of scope | designing before agreeing on scope |
| 2. Capacity estimates | 3 | 5 | QPS (average and peak), storage per year, bandwidth, read:write ratio | precision theater; only compute what changes the design |
| 3. API | 3 | 5 | 3–6 endpoints or messages with key fields, pagination, idempotency | a full OpenAPI spec |
| 4. Data model | 3 | 5 | entities, keys, the access pattern each serves, store choice | picking a DB before knowing the queries |
| 5. High-level design | 8 | 10 | boxes and arrows for the read path and the write path | drawing every box you know |
| 6. Deep dives | 15 | 18 | the 2–3 hardest parts, each with options and a decision | staying shallow everywhere |
| 7. Bottlenecks, trade-offs, failure modes | 5 | 7 | what breaks at 10x, single points of failure, what you'd monitor | claiming the design has no weaknesses |
| 8. Wrap-up | 3 | 3 | a 30-second summary and what you'd do next | running out of time mid-sentence |
Sums to 45 and 60. Watch the clock: if the high-level design isn't on the board by minute 20, cut estimation and API detail. The same framework with output checklists: System design: Approach framework.
Estimation drills
| Quantity | Rule of thumb |
|---|---|
| Seconds per day | 86,400 ≈ 10^5 (dividing by the bigger number understates QPS by ~14%, which is fine) |
| Average QPS | DAU × actions per user per day / 10^5; 1M DAU × 10 actions ≈ 100 QPS |
| Peak QPS | 2–3x average; 10x or more for launches, live events, marketing blasts |
| 1M requests a day | ≈ 12 per second |
| Storage per year | writes per day × bytes per write × 365 (≈ 400) × replication factor (3) |
| Bandwidth | QPS × bytes per response |
| Cache size | about 20% of a day's distinct hot data (80/20 rule) |
| Servers | peak QPS / QPS per server; say your per-server assumption out loud |
| Concurrent users | DAU × fraction online at peak (often 10–20%) |
| Size | Value |
|---|---|
| UUID / Snowflake ID / timestamp | 16 B / 8 B / 8 B |
| Short text post with metadata | ~300 B–1 KB |
| Compressed photo | 200 KB–2 MB |
| 1080p video stream | ≈ 5 Mbps ≈ 2.25 GB per hour (Netflix recommends 5 Mbps for 1080p (opens in a new tab)) |
| Powers | 10^3 K, 10^6 M, 10^9 G, 10^12 T, 10^15 P |
Latency numbers, powers of two and the full formulas: System design: Back-of-envelope numbers.
Worked estimate: a Twitter-like feed
| Assumption or result | Math |
|---|---|
| 200M DAU, 10 feed loads and 0.5 posts per user per day | given |
| Feed reads | 200M × 10 / 10^5 = 20,000 QPS, ~60,000 at peak |
| Post writes | 100M / 10^5 = 1,000 QPS, ~3,000 at peak |
| Text storage | 100M × 500 B = 50 GB/day → ~18 TB/year, ~55 TB with 3 replicas |
| Media | 10% of posts with a 500 KB image = 5 TB/day → ~1.8 PB/year: object storage + CDN |
| Fan-out on write | 100M posts × 200 followers = 20B feed inserts/day ≈ 200,000/s |
| Conclusion | reads dominate requests, but fan-out dominates writes: a hybrid fan-out is needed (see Design: news feed) |
Worked estimate: a chat app
| Assumption or result | Math |
|---|---|
| 50M DAU, 40 messages sent per user per day, 15% online at peak | given |
| Messages | 2B/day / 10^5 = 20,000 writes/s, ~60,000 at peak |
| Storage | 2B × 200 B = 400 GB/day → ~146 TB/year, ~440 TB with 3 replicas |
| Open connections | 50M × 15% = 7.5M WebSockets at peak |
| Gateway servers | at an assumed 100k connections per box: ~75, plus headroom for a zone failure (WhatsApp reported 2M+ per server in 2012, blog (opens in a new tab)) |
| Conclusion | storage is append-heavy and read by conversation: a wide-column store partitioned by conversation |
How to drive
Clarifying questions
| Area | Ask |
|---|---|
| Users and scale | How many daily users? Growth? Global or one region? |
| Core features | Which 3 features matter most? What can we leave out? |
| Read vs write | Read-heavy or write-heavy? Ratio? |
| Latency | What must be fast (p99 target)? What can be slow or async? |
| Consistency | Can users see stale data? For how long? Anything that must never be wrong (money, inventory)? |
| Availability | What happens if it's down for a minute? |
| Data | How big is an item? How long is it kept? Deletions, privacy, residency? |
| Clients | Web, mobile, other services? Poor networks? |
| Existing systems | Build on anything existing (auth, a data warehouse)? |
Behaviors that score
| Situation | Do |
|---|---|
| Thinking | narrate options and pick one: "Two options: A is simpler, B scales further. Given 20k QPS, A is enough; I'll note where B would come in." |
| Pushback ("what if the cache dies?") | treat it as a hint, not an attack: restate, answer with a mechanism, name the cost; change your design if they're right |
| You don't know a technology | reason from principles: "I haven't run Cassandra, but a leaderless store with quorums would give us…" |
| Deep vs broad | broad first (a complete, simple design by minute ~20), then deep where the problem is hard or the interviewer points |
| Interviewer goes quiet | check in: "Want me to go deeper on fan-out, or cover failure handling?" |
| Stuck | go back to requirements and numbers; simplify to one server, then scale the bottleneck |
Drawing tips
- Left to right: clients → edge (CDN, LB, gateway) → services → data stores; async parts (queues, workers) below.
- Number the arrows of the write path and the read path separately, or draw them in two colors.
- Label arrows with what flows (
post_id,event), stores with their key (messages (conv_id, seq)). - Keep a corner for requirements and numbers; point back at them when you decide something.
- Redraw rather than squeeze: a clean second diagram beats a tangle.
Design: rate limiter
Requirements. Functional: limit requests per client (user, API key or IP) per rule (100/min on POST /login);
return 429 with Retry-After; rules change without deploys. Non-functional: adds under 5 ms, available (the API
must work if the limiter doesn't), approximately accurate across many API servers. Out of scope: network-level DDoS
(CDN/WAF), billing quotas.
Estimates. 1M requests/s across the fleet, 10M active keys × ~64 B of state ≈ 640 MB: fits one Redis, but 1M operations/s needs a sharded cluster (~10 shards at an assumed 100k ops/s each).
API. Internal call from the gateway middleware; clients only see headers.
allow(key="user:42", rule="login") →
{ allowed: false, remaining: 0, retryAfterMs: 1200 }
HTTP/1.1 429 Too Many Requests
Retry-After: 2Data model. Redis hash per key and rule, rl:login:user:42 → { tokens, ts }, with a TTL equal to the time
to refill. Rules in a config store, cached in every gateway.
High-level design.
client ─► API gateway / middleware ──allowed──► service
│ ▲
Lua script│ │ allowed, remaining, retry-after
▼ │
Redis cluster (sharded by limiter key)
rules service ─► pushes rule config to every gatewayDeep dives.
| Problem | Decision |
|---|---|
| Algorithm | token bucket: allows short bursts, 2 numbers per key; sliding-window counter if bursts must be smooth (algorithms) |
| Race between gateways | the whole read-refill-take in one Lua script, atomic in Redis; use Redis's clock (TIME) so gateway clock skew doesn't matter |
| Latency | Redis in the same zone, pipelined; a local deny cache for keys already over the limit until their Retry-After |
| Redis down | fail open (allow) for normal endpoints, fail closed for abuse-prone ones (login, signup); circuit breaker around Redis |
| Very high rates | per-gateway local buckets with a share of the limit, synced every ~100 ms: cheaper, slightly inexact |
| Multi-region | a limit per region (limit / regions) or async global sync; accept small overshoot |
Trade-offs. Exact global counting costs a network hop per request; local approximate limits are free but overshoot. Failing open keeps the product up during limiter outages at the price of brief unlimited traffic.
Design: news feed
Requirements. Functional: post (text and media), follow, home feed in reverse-chronological order, paginated. Non-functional: feed loads under 500 ms p99; a post shows up in followers' feeds within seconds (eventual consistency is fine); very read-heavy; highly available. Out of scope: ads, ML ranking internals, DMs.
Estimates. From the worked estimate: 20k feed QPS (60k peak), 1k posts/s, 200k feed inserts/s if every post fans out.
API.
POST /v1/posts { text, mediaIds[] } → { postId }
GET /v1/feed?cursor=<lastPostId>&limit=20
POST /v1/users/{id}/followData model. posts (post_id Snowflake, author_id, text, media, created_at) sharded by author_id;
follows stored twice, keyed by follower and by followee; feed cache: a Redis list per user
feed:{user_id} of the latest ~800 post IDs.
High-level design.
┌─► post service ─► posts DB (by author_id)
│ │
client ─► LB┤ ▼ post.created
│ Kafka ─► fan-out workers ─► feed cache
│ │ followers (Redis list
│ ▼ from graph per user)
│ social graph DB
└─► feed service ─► feed cache + celebrity posts
│ merge, hydrate IDs
▼
post cache, user cache; media via CDNDeep dives.
| Problem | Decision |
|---|---|
| Fan-out | hybrid: push post IDs to followers' feeds, except authors above ~10k followers, whose recent posts are pulled and merged at read time (concepts) |
| Inactive users | skip fan-out to users inactive for 30 days; rebuild their feed on next login |
| Hydration | the feed holds IDs; batch-get posts and authors from caches (multi-get), fill misses from the DB |
| Pagination | cursor = last post ID; Snowflake IDs sort by time, so no offsets |
| Media | presigned upload to object storage, thumbnails by a worker, served through the CDN |
| Counts (likes) | async counters: increments go through a queue into sharded counters, cached |
| Ranking (if asked) | the feed cache becomes the candidate set; a ranking service scores and reorders |
Trade-offs. Push makes reads one lookup but multiplies writes and storage; pull is cheap to write but slow to read for people who follow many accounts. The hybrid adds a merge step and a threshold to tune.
Design: chat
Requirements. Functional: 1:1 and group chat (up to ~1,000 members), delivery receipts (sent, delivered, read), presence, offline delivery with push notifications, history synced across devices. Non-functional: delivery in a few hundred ms, no message lost, correct order within a conversation, highly available. Out of scope: voice and video; end-to-end encryption (mention: servers store ciphertext).
Estimates. From the worked estimate: 20k messages/s (60k peak), 7.5M open WebSockets, ~146 TB/year before replication.
API. WebSocket frames for live traffic, REST for history.
→ send { convId, clientMsgId, body }
← ack { clientMsgId, msgId, seq } stored durably
← msg { convId, seq, msgId, sender, body }
→ read { convId, upToSeq }
GET /v1/conversations/{id}/messages?beforeSeq=900&limit=50Data model. messages (conv_id, seq, msg_id, sender_id, body, created_at) in a wide-column store
(Cassandra/ScyllaDB): partition by conv_id, cluster by seq. members (conv_id, user_id, last_read_seq).
Session registry in Redis: user_id → gateway server.
High-level design.
client ◄══ WebSocket ══► chat gateway (stateful, holds sockets)
│ send ▲ deliver
▼ │
chat service ─► per-conversation seq
│ (Redis INCR / partition)
┌──────────────┼───────────────┐
▼ ▼ ▼
messages DB session registry push service
(conv_id, seq) user → gateway APNs / FCM for
(Redis) offline usersDeep dives.
| Problem | Decision |
|---|---|
| Ordering | a sequence number per conversation from a single owner (atomic INCR, or Kafka partitioned by conv_id); clients sort by seq, never by device clocks |
| No loss, no duplicates | the client retries with the same clientMsgId; the server stores it idempotently and acks only after the durable write |
| Delivery | look up recipients' gateways in the registry and push; at-least-once, the client dedupes by msgId |
| Offline and multi-device | each device keeps its last seq per conversation; on reconnect it asks for everything after it; push notification if no device is connected |
| Receipts | delivered = a recipient device acked; read = last_read_seq; in big groups show counts, not per-member lists |
| Presence | heartbeat every ~30 s sets a Redis key with a TTL; only send presence changes to contacts who are online and looking |
| Gateway dies | clients reconnect through the LB to another gateway, re-register, resync from their last seq |
| Large groups | fan-out to each member's gateway for small groups; for huge channels, members pull on open |
Trade-offs. WebSockets give instant two-way delivery but make gateways stateful (draining, sticky routing, a registry). A per-conversation sequencer guarantees order but is a hot spot for very busy conversations.
Design: notification system
Requirements. Functional: internal services send notifications through push (iOS, Android), SMS, email and in-app; templates with localization; user preferences and opt-outs; scheduling and quiet hours; delivery and open tracking. Non-functional: transactional messages (2FA codes) in seconds, bulk campaigns can take minutes; no loss; duplicates rare; absorbs spikes. Out of scope: the campaign authoring tool.
Estimates. 10M push + 5M email + 1M SMS a day ≈ 200/s average; a campaign of 10M in 10 minutes ≈ 17k/s: queues absorb the spike, workers drain at provider limits.
API.
POST /v1/notifications
Idempotency-Key: order-981-shipped
{ "userId": "u1", "template": "order_shipped",
"params": { "orderId": "981" },
"channels": ["push", "email"],
"priority": "transactional" }
202 Accepted { "notificationId": "n_7f3" }Data model. notifications (id, user_id, template, params, created_at), deliveries (notification_id, channel, status, attempts, provider_id), preferences (user_id, category, channel, opted_in, quiet_hours),
devices (user_id, token, platform, last_seen), versioned templates.
High-level design.
services ─► notification API ─► dedupe (idempotency key)
│
▼
preferences + per-user limits ─► render template
│
┌──────────┼──────────┬──────────┐
▼ ▼ ▼ ▼
push queue SMS queue email queue in-app queue
(per priority: transactional ≠ bulk)
│ │ │ │
▼ ▼ ▼ ▼
APNs / FCM SMS vendor email vendor inbox + WebSocket
└──── status callbacks ─► deliveries DB, retries, DLQDeep dives.
| Problem | Decision |
|---|---|
| Priority isolation | separate queues and workers for transactional and bulk, so a campaign never delays a 2FA code |
| Duplicates | idempotency key on the API; dedupe per (notification, channel) in workers; external providers make exactly-once impossible, so aim for rare duplicates |
| Provider failures | retries with backoff, then a DLQ; a secondary SMS/email vendor behind a circuit breaker |
| Provider limits | token-bucket rate limiter per provider in the workers |
| User limits | caps per user per category ("at most 3 marketing pushes a day"), quiet hours via a delay queue in the user's time zone |
| Dead device tokens | remove tokens the provider reports as invalid |
| Tracking | provider webhooks for delivered and bounced; tracked links for opens and clicks; events to analytics |
Trade-offs. Checking preferences at send time (not enqueue time) respects late opt-outs but adds a lookup per message. More channels and vendors improve reach and resilience at the cost of more integration code.
Design: typeahead
Requirements. Functional: top 5–10 completions for the typed prefix, ranked by popularity; refresh at least daily, with trending terms sooner. Non-functional: results visible under ~100 ms while typing (server budget ~10–20 ms); very high read QPS; slight staleness fine. Out of scope: spell correction, the search results page.
Estimates. 100M DAU × 10 searches × ~6 requests per search (debounced) = 6B/day ≈ 60k QPS, ~150k at peak. 10M popular queries; the prefix → top-k map is a few GB to tens of GB: fits in RAM, so replicate rather than shard.
API.
GET /v1/suggest?q=new%20yo&limit=8
200 ["new york times", "new york weather", …]
Cache-Control: public, max-age=300Data model. Offline: query_counts (query, day, count) built from search logs. Serving: a trie where each
node stores its top-k completions, or a flat map prefix → [top-k queries].
High-level design.
client: debounce ~100 ms, cancel stale requests,
cache results per prefix
│
▼
CDN / edge cache (popular prefixes, TTL minutes)
│ miss
▼
suggest service ─► in-memory prefix → top-k (replicas)
▲ load new snapshot, swap atomically
search logs ─► Kafka ─► aggregation job (counts, time decay)
─► build job ─► snapshot in object storeDeep dives.
| Problem | Decision |
|---|---|
| Query speed | precompute top-k at every prefix node, so a lookup is O(length of prefix), no subtree walk |
| Freshness | rebuild the snapshot hourly or daily with time-decayed counts; a small streaming layer (count-min sketch per window) adds trending queries, merged at serve time |
| Size | cap prefix length (~20 chars) and store query IDs instead of strings |
| Scaling | replicate the whole index per region; if it outgrows RAM, shard by hash of the prefix (first-letter sharding is uneven: many queries start with "s") |
| Client load | debounce, reuse the results of a shorter prefix when fewer than k items would match, cache at the edge |
| Safety | filter offensive and legally blocked terms at build time; a kill-switch list checked at serve time |
| Personalization | blend the user's recent searches (client-side or a small per-user list) with global results |
Trade-offs. Precomputing top-k makes reads trivial but rebuilds cost compute and add staleness; a streaming layer buys freshness with complexity. Replication over sharding wastes RAM but avoids fan-out per keystroke.
Design: web crawler
Requirements. Functional: from seed URLs, fetch HTML pages, extract links, store pages for an indexer,
recrawl pages as they change. Non-functional: 1B pages a month; polite (obey robots.txt, limit requests per
host); robust to traps and broken HTML; no duplicate work. Out of scope: indexing and ranking, JavaScript
rendering (a headless-browser pool is an extension).
Estimates. 1B / 2.6M s ≈ 400 pages/s, ~1,000 at peak. At 100 KB per page: 40 MB/s inbound and 100 TB a month before compression.
API. Internal: seed(urls); output is pages in the content store plus a page.fetched event for the
indexer.
Data model. URL frontier (priority queues feeding per-host queues); seen-URL set (Bloom filter in front of a
store keyed by URL hash); pages (url_hash, url, fetched_at, status, content_hash, next_crawl_at); page bodies in
object storage by content hash.
High-level design.
seeds ─► URL frontier ───────────────► fetcher workers
front: priority queues │ ▲ DNS cache
back: one queue per host │ │ robots.txt cache
▲ ▼
│ content store (S3)
│ + pages metadata DB
│ │
seen? (Bloom + store) ▼
normalize, filter ◄── link extractor ◄── parser, dedupe
(content hash, SimHash)Deep dives.
| Problem | Decision |
|---|---|
| Politeness | each host maps to one back queue, served by one worker, with a delay between requests (≥ 1 s or the site's crawl delay); cache and obey robots.txt (RFC 9309 (opens in a new tab)) |
| Priority | front queues by importance (link-based rank, domain quality) and expected change rate; the Mercator layout (paper (opens in a new tab)) |
| Duplicate URLs | normalize (lowercase host, drop fragments, sort query parameters), then a Bloom filter: a false positive skips one URL, acceptable |
| Duplicate content | exact: content hash; near-duplicates: SimHash with a small Hamming distance |
| Traps | caps on depth, URL length and pages per host; detect calendar-like infinite URL patterns |
| DNS | a local caching resolver; DNS lookups are often the hidden bottleneck |
| Recrawl | next_crawl_at from the observed change rate; conditional GET (If-Modified-Since, ETag) |
| Distribution and recovery | partition the frontier by host hash across nodes; checkpoint queues so a crash resumes |
Trade-offs. Politeness caps throughput per host, so parallelism must come from many hosts at once. A Bloom filter saves memory but occasionally skips a new URL; a full seen-set is exact but costs a lookup per link.
Design: proximity search
Requirements. Functional: find businesses near a point within a radius, filter by category, sort by distance or rating; owners add and edit businesses; view details. Non-functional: search under 200 ms p99; very read-heavy; edits can take minutes to appear; highly available. Out of scope: reviews, map rendering, directions.
Estimates. 200M businesses × 1 KB = 200 GB of details. Geo index entries (ID + geohash + lat/lng ≈ 32 B) × 200M ≈ 6.4 GB: fits in memory. 100M DAU × 5 searches = 500M/day ≈ 5k QPS, ~15k peak.
API.
GET /v1/search?lat=40.73&lng=-73.99&radius=2000
&category=coffee&cursor=…&limit=20
GET /v1/businesses/{id}
POST /v1/businesses (owners, rate limited)Data model. businesses (id, name, lat, lng, category, rating, …) in Postgres as the source of truth. Geo
index: (geohash, business_id) with a B-tree on geohash, or PostGIS with a GiST index, or an in-memory
quadtree service.
High-level design.
client ─► LB ─► search service ─► geo index service
│ (in memory, replicated;
│ geohash or quadtree)
│ → candidate IDs
▼
business cache (Redis) ─► businesses DB
│ CDC / events
owners ─► business service ─► DB ──────────┴─► index updaterDeep dives.
| Problem | Decision |
|---|---|
| Index choice | geohash: simple, works in any B-tree or Redis GEOSEARCH; quadtree: adapts to density; PostGIS ST_DWithin is enough at modest scale (structures) |
| Query | pick the geohash length from the radius (~1 km → 6 chars, ~5 km → 5), search the cell plus its 8 neighbors, then filter by exact distance and sort |
| Sparse areas | too few results → shorter prefix (bigger cells) until enough |
| Quadtree | split any node holding more than ~100 businesses; build in memory at startup, apply updates incrementally |
| Scaling | the index fits in RAM, so replicate it per region; shard by region only if it stops fitting (and watch dense cities) |
| Caching | cache results per (geohash, category) for popular areas; business details in Redis |
| Moving objects (ride-hailing variant) | drivers send locations every few seconds into an in-memory geo store (Redis GEO, H3 cells); no durable writes per ping |
Trade-offs. Geohash is easy to store and shard but has cell-edge effects (hence neighbors); quadtrees fit density better but live in memory and need rebuilding. Minutes of index lag keep writes cheap.
Design: video streaming
Requirements. Functional: upload videos, transcode into several resolutions, stream with adaptive bitrate on any device, video metadata. Non-functional: playback starts in about 2 s with rare rebuffering; global viewers; uploads never lost; egress cost under control. Out of scope: recommendations, comments, DRM details, live streaming (differences noted below).
Estimates. 500k uploads/day × 5 min; ~500 MB per video across all renditions → 250 TB/day of new storage.
Viewing: 100M views/day × 5 min × 5 Mbps ≈ 1.5 × 10^11 Mb/day ≈ 1.7 Tbps average egress: only a CDN can
serve that.
API.
POST /v1/videos → { videoId, uploadUrls[] }
PUT <presigned part URL> (multipart, resumable)
POST /v1/videos/{id}/complete
GET /v1/videos/{id} → { title, status, manifestUrl }
GET cdn/v/{id}/master.m3u8 (HLS manifest)Data model. videos (id, owner_id, title, status: uploading → processing → ready, duration), renditions (video_id, height, bitrate, codec, playlist_path). Object storage: raw/{id} and
hls/{id}/{rendition}/seg_00001.m4s.
High-level design.
creator ─► upload API ─► multipart PUTs ─► object store (raw)
│ upload event
▼
transcode DAG: split into chunks →
encode renditions in parallel →
package HLS/DASH, thumbnails
│
▼
viewer ◄── CDN / ISP caches ◄──── object store (segments)
│ manifest, then segments
└─► video API ─► metadata DB (status, renditions)Deep dives.
| Problem | Decision |
|---|---|
| Transcoding | a DAG of jobs: split at keyframes, encode chunks in parallel (360p … 1080p, 4K), reassemble, package; workers pull from a queue and are retried per chunk |
| Adaptive bitrate | HLS (RFC 8216 (opens in a new tab)) or DASH: a manifest lists renditions; the player picks a rendition per segment (2–6 s) from measured bandwidth and buffer |
| Delivery | segments are immutable: cache forever at the CDN; push popular titles to edge and ISP caches ahead of demand, pull the long tail from origin |
| Cost | encode expensive codecs (AV1) only for videos that get views; cold storage tiers for old raw files |
| Uploads | multipart with per-part retries, resumable; content hash to dedupe |
| Status | transcode events update status; the creator is notified when ready |
| Live variant | ingest over RTMP/SRT, transcode in real time, short segments (1–2 s) or low-latency HLS |
Trade-offs. More renditions mean smoother playback on bad networks but more storage and encode time. Pre-positioning at the edge cuts egress and latency for hits but wastes cache on videos nobody watches.
Design: key-value store
A Dynamo-style store, as in the Dynamo paper (opens in a new tab) and Cassandra.
Requirements. Functional: get(key), put(key, value), delete(key); values up to ~1 MB.
Non-functional: always writable, tunable consistency, horizontal scale to hundreds of nodes, single-digit ms p99
in a region, durable. Out of scope: transactions, secondary indexes, range scans.
Estimates. 10 TB × 3 replicas = 30 TB; at ~2 TB usable per node → ~15 nodes for space; 1M ops/s at an assumed 50k ops/s per node → ~20 nodes. Plan ~25–30 for headroom.
API.
get(key, consistency=QUORUM) → [(value, version)]
more than one version = siblings to merge
put(key, value, context, consistency=QUORUM)
context = the version from the previous get
delete(key) → writes a tombstoneData model. Each node runs an LSM engine: commit log → memtable → SSTables with Bloom filters, compacted in the background. Cluster state: ring membership and token ownership, spread by gossip.
High-level design.
client ─► any node = coordinator for this request
│ hash(key) on ring → preference list: next N=3
┌────────┼────────┐ distinct physical nodes
▼ ▼ ▼
node B node C node D write: wait for W=2 acks
(LSM) (LSM) (LSM) read: ask R=2, newest wins,
▲ ▲ ▲ repair stale replicas
└── gossip membership, Merkle-tree anti-entropy
D down → node E keeps D's writes as hints (hinted handoff)Deep dives.
| Problem | Decision |
|---|---|
| Partitioning | consistent hashing with virtual nodes; replicas on distinct machines and zones (concepts) |
| Consistency | N=3; per-request R and W; R + W > N for overlap, W=1 for fastest writes |
| Conflicts | version vectors return concurrent versions as siblings for the client to merge (Dynamo), or last-write-wins by timestamp (Cassandra: simpler, silently drops concurrent writes) |
| Temporary failures | sloppy quorum + hinted handoff: another node accepts writes and hands them back later |
| Permanent divergence | read repair on reads; background anti-entropy compares Merkle trees per range |
| Membership | gossip every second; a phi-accrual failure detector marks nodes down |
| Deletes | tombstones kept for a grace period (longer than the repair interval) so deleted data doesn't come back |
| Adding nodes | the new node takes virtual nodes and streams those ranges from their current owners |
Trade-offs. Choosing availability (AP) means clients see conflicts or stale reads; a consensus-based CP store (etcd, Spanner) avoids that but rejects writes without a majority. LWW is simple but loses data under concurrency; version vectors keep it but push merging to clients.
Common mistakes
| Mistake | Instead |
|---|---|
| Jumping straight to boxes | spend 5 minutes on requirements and write them down |
| Estimating everything to 3 significant figures | orders of magnitude, and only numbers that change a decision |
| Naming technologies instead of reasons ("use Kafka") | the property you need, then a tool that has it ("a replayable log, e.g. Kafka") |
| Microservices for a 1k QPS problem | the simplest design that meets the numbers, then scale the bottleneck |
| No read path vs write path distinction | trace one request of each through the diagram |
| Ignoring failure | for each box: what if it's slow, down, or the data is stale? |
| Single points of failure left in the final design | replicate, or say explicitly why it's acceptable |
| Hand-waving the hard part (fan-out, ordering, hot keys) | find it early and spend the deep-dive time there |
| Defending a design against a good hint | take it, adjust, explain the new trade-off |
| Silent thinking | narrate; the interviewer scores reasoning they can hear |
| Running out of time before a summary | keep the last 3 minutes for wrap-up |
Scoring rubric
| Dimension | Weak | Strong |
|---|---|---|
| Problem navigation | vague scope, missing non-functional requirements | crisp scope with numbers; finds the hard part early |
| High-level design | missing components or data flow; wrong store for the access pattern | complete read and write paths; each component justified |
| Technical depth | buzzwords, can't go one level down | mechanisms explained (how the cache is invalidated, how order is kept) |
| Trade-offs | one option presented as the only one | alternatives compared against requirements; costs admitted |
| Failure and operations | assumes everything works | failure modes, monitoring, rollout, cost (expected more at senior+) |
| Communication | long silences, messy diagram, ignores hints | structured, checks in, adapts, clean diagram |
Interviewers write down specific evidence for each dimension. Make evidence easy to write: say the requirement, the options, the decision and the cost, in that order.
Prep plan
| Week | Focus | Do |
|---|---|---|
| 1 | concepts | read System design concepts and the Primer; DDIA chapters 5–9 (replication, partitioning, transactions, consistency) |
| 2 | framework + classics | run the framework on a URL shortener, rate limiter and news feed, timed at 45 minutes, out loud |
| 3 | breadth | chat, notifications, typeahead, crawler, proximity, video, KV store; Alex Xu vol. 1–2 for more designs |
| 4 | mocks | 3–5 mock interviews with peers or a service; review each against the rubric; re-do the weakest design |
Resources: System Design Primer (opens in a new tab), Designing Data-Intensive Applications (opens in a new tab), Alex Xu's System Design Interview vol. 1 and 2, Hello Interview (opens in a new tab), ByteByteGo (opens in a new tab). With a week or less: the framework, the estimation table, and three designs practiced out loud.
Recipes
The first 5 minutes
A script that turns a one-line prompt into a scoped problem.
- "Let me make sure I understand what we're building. The core features are A, B and C; is that right? Anything you want to add or drop?"
- "Who uses it and how much? Roughly how many daily users, and is it global?"
- "What matters most: latency, availability, consistency? Can a user see stale data for a few seconds?"
- "I'll leave out D and E unless you want them."
- Write the list in a corner: features, numbers, non-functional targets, out of scope. "I'll do a quick estimate next."
Estimation template
Fill in, round hard, circle the numbers that change the design.
DAU = ______
actions/user/day = ______ (reads ____, writes ____)
avg QPS = DAU × actions / 10^5 = ______ peak ×3 = ______
read:write ratio = ______
bytes per write = ______
storage/year = writes/day × bytes × 400 × 3 replicas = ______
bandwidth = peak QPS × bytes per response = ______
cache = 20% of daily hot data = ______
servers = peak QPS / (QPS per server) = ______
→ implications: ______________________________________Deep-dive question bank
Questions to ask yourself about each box before the interviewer does.
| Component | Ask |
|---|---|
| Load balancer | L4 or L7? Health checks? Sticky sessions needed (WebSockets)? What if it dies? |
| Service | stateless? how does it scale? timeouts and retries to dependencies? |
| Cache | what's cached, TTL, invalidation on write, hit rate, stampede, hot keys, cache down? |
| Database | why this store? access patterns and keys, indexes, replication, sharding key, hot partitions? |
| Queue | delivery semantics, idempotent consumers, ordering, DLQ, consumer lag, backpressure? |
| Search index | how is it fed (CDC), how stale, how to rebuild? |
| Object store + CDN | presigned uploads, cache keys, invalidation, egress cost? |
| Anything | single point of failure? what breaks at 10x? what do we monitor? |
Wrap-up in 3 minutes
- One-sentence recap of the design and the two key decisions.
- The main bottleneck at 10x and how you'd address it.
- The biggest risk or weakest point you know about.
- What you'd build next (monitoring, multi-region, a feature you left out).
When you're stuck
- Re-read the requirements and numbers; the answer is usually a constraint you forgot.
- Solve it for one server, then ask "what breaks first when this gets 100x the load?"
- Name two options and compare them on the requirement that matters most.
- Ask a targeted question: "Is it acceptable if a user sees their own post a second late?"
Practice loop for one design
- Set a 45-minute timer; talk out loud; draw on a real whiteboard tool.
- Afterward, compare with a reference design (Primer, Alex Xu, Hello Interview) and note what you missed.
- Score yourself on the rubric; write the 3 gaps.
- Redo the same design three days later in 30 minutes.
References
- Amazon Dynamo paper (opens in a new tab): the model for the key-value store design
- Mercator: a scalable, extensible web crawler (opens in a new tab): URL frontier with priority and politeness queues
- RFC 9309: Robots Exclusion Protocol (opens in a new tab) and RFC 8216: HLS (opens in a new tab): crawler and streaming standards
- Twitter: Timelines at Scale (opens in a new tab): hybrid fan-out in production
- Discord: how Discord stores trillions of messages (opens in a new tab): chat storage on Cassandra and ScyllaDB
- The System Design Primer (opens in a new tab): free topic index and practice designs
- Designing Data-Intensive Applications (opens in a new tab): the concepts in depth
- Alex Xu, System Design Interview vol. 1 and 2 (ByteByteGo (opens in a new tab)): worked designs in interview format
- Hello Interview: system design in a hurry (opens in a new tab): delivery framework and level expectations
- System design concepts: the building blocks behind every decision here
- System design (production reference): numbers, tables and the worked URL shortener