../

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.

LevelWhat's askedDepth expectedYou pass by
New grad / junioroften no design round; sometimes API or object design, or a small servicea working single-region design: clients, LB, stateless service, DB, cache; a sensible API and schemaclear requirements, correct basics, taking hints well; the interviewer may steer
Mid-levelthe full roundcover the whole framework end to end; one solid deep dive, possibly promptedknowing what each standard component does and when to use it; no big gaps
Seniorthe full round, less guidancemove quickly through the basics; lead 2–3 deep dives unprompted; numbers behind decisions; failure modestrade-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 ownershippushing 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

Step45 min60 minProduceWatch for
1. Requirements573–5 functional features, non-functional numbers (scale, latency, availability, consistency), what's out of scopedesigning before agreeing on scope
2. Capacity estimates35QPS (average and peak), storage per year, bandwidth, read:write ratioprecision theater; only compute what changes the design
3. API353–6 endpoints or messages with key fields, pagination, idempotencya full OpenAPI spec
4. Data model35entities, keys, the access pattern each serves, store choicepicking a DB before knowing the queries
5. High-level design810boxes and arrows for the read path and the write pathdrawing every box you know
6. Deep dives1518the 2–3 hardest parts, each with options and a decisionstaying shallow everywhere
7. Bottlenecks, trade-offs, failure modes57what breaks at 10x, single points of failure, what you'd monitorclaiming the design has no weaknesses
8. Wrap-up33a 30-second summary and what you'd do nextrunning 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

QuantityRule of thumb
Seconds per day86,400 ≈ 10^5 (dividing by the bigger number understates QPS by ~14%, which is fine)
Average QPSDAU × actions per user per day / 10^5; 1M DAU × 10 actions ≈ 100 QPS
Peak QPS2–3x average; 10x or more for launches, live events, marketing blasts
1M requests a day≈ 12 per second
Storage per yearwrites per day × bytes per write × 365 (≈ 400) × replication factor (3)
BandwidthQPS × bytes per response
Cache sizeabout 20% of a day's distinct hot data (80/20 rule)
Serverspeak QPS / QPS per server; say your per-server assumption out loud
Concurrent usersDAU × fraction online at peak (often 10–20%)
SizeValue
UUID / Snowflake ID / timestamp16 B / 8 B / 8 B
Short text post with metadata~300 B–1 KB
Compressed photo200 KB–2 MB
1080p video stream≈ 5 Mbps ≈ 2.25 GB per hour (Netflix recommends 5 Mbps for 1080p (opens in a new tab))
Powers10^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 resultMath
200M DAU, 10 feed loads and 0.5 posts per user per daygiven
Feed reads200M × 10 / 10^5 = 20,000 QPS, ~60,000 at peak
Post writes100M / 10^5 = 1,000 QPS, ~3,000 at peak
Text storage100M × 500 B = 50 GB/day → ~18 TB/year, ~55 TB with 3 replicas
Media10% of posts with a 500 KB image = 5 TB/day → ~1.8 PB/year: object storage + CDN
Fan-out on write100M posts × 200 followers = 20B feed inserts/day ≈ 200,000/s
Conclusionreads dominate requests, but fan-out dominates writes: a hybrid fan-out is needed (see Design: news feed)

Worked estimate: a chat app

Assumption or resultMath
50M DAU, 40 messages sent per user per day, 15% online at peakgiven
Messages2B/day / 10^5 = 20,000 writes/s, ~60,000 at peak
Storage2B × 200 B = 400 GB/day → ~146 TB/year, ~440 TB with 3 replicas
Open connections50M × 15% = 7.5M WebSockets at peak
Gateway serversat 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))
Conclusionstorage is append-heavy and read by conversation: a wide-column store partitioned by conversation

How to drive

Clarifying questions

AreaAsk
Users and scaleHow many daily users? Growth? Global or one region?
Core featuresWhich 3 features matter most? What can we leave out?
Read vs writeRead-heavy or write-heavy? Ratio?
LatencyWhat must be fast (p99 target)? What can be slow or async?
ConsistencyCan users see stale data? For how long? Anything that must never be wrong (money, inventory)?
AvailabilityWhat happens if it's down for a minute?
DataHow big is an item? How long is it kept? Deletions, privacy, residency?
ClientsWeb, mobile, other services? Poor networks?
Existing systemsBuild on anything existing (auth, a data warehouse)?

Behaviors that score

SituationDo
Thinkingnarrate 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 technologyreason from principles: "I haven't run Cassandra, but a leaderless store with quorums would give us…"
Deep vs broadbroad first (a complete, simple design by minute ~20), then deep where the problem is hard or the interviewer points
Interviewer goes quietcheck in: "Want me to go deeper on fan-out, or cover failure handling?"
Stuckgo 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: 2

Data 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 gateway

Deep dives.

ProblemDecision
Algorithmtoken bucket: allows short bursts, 2 numbers per key; sliding-window counter if bursts must be smooth (algorithms)
Race between gatewaysthe whole read-refill-take in one Lua script, atomic in Redis; use Redis's clock (TIME) so gateway clock skew doesn't matter
LatencyRedis in the same zone, pipelined; a local deny cache for keys already over the limit until their Retry-After
Redis downfail open (allow) for normal endpoints, fail closed for abuse-prone ones (login, signup); circuit breaker around Redis
Very high ratesper-gateway local buckets with a share of the limit, synced every ~100 ms: cheaper, slightly inexact
Multi-regiona 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}/follow

Data 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 CDN

Deep dives.

ProblemDecision
Fan-outhybrid: push post IDs to followers' feeds, except authors above ~10k followers, whose recent posts are pulled and merged at read time (concepts)
Inactive usersskip fan-out to users inactive for 30 days; rebuild their feed on next login
Hydrationthe feed holds IDs; batch-get posts and authors from caches (multi-get), fill misses from the DB
Paginationcursor = last post ID; Snowflake IDs sort by time, so no offsets
Mediapresigned 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=50

Data 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 users

Deep dives.

ProblemDecision
Orderinga 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 duplicatesthe client retries with the same clientMsgId; the server stores it idempotently and acks only after the durable write
Deliverylook up recipients' gateways in the registry and push; at-least-once, the client dedupes by msgId
Offline and multi-deviceeach device keeps its last seq per conversation; on reconnect it asks for everything after it; push notification if no device is connected
Receiptsdelivered = a recipient device acked; read = last_read_seq; in big groups show counts, not per-member lists
Presenceheartbeat every ~30 s sets a Redis key with a TTL; only send presence changes to contacts who are online and looking
Gateway diesclients reconnect through the LB to another gateway, re-register, resync from their last seq
Large groupsfan-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, DLQ

Deep dives.

ProblemDecision
Priority isolationseparate queues and workers for transactional and bulk, so a campaign never delays a 2FA code
Duplicatesidempotency key on the API; dedupe per (notification, channel) in workers; external providers make exactly-once impossible, so aim for rare duplicates
Provider failuresretries with backoff, then a DLQ; a secondary SMS/email vendor behind a circuit breaker
Provider limitstoken-bucket rate limiter per provider in the workers
User limitscaps 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 tokensremove tokens the provider reports as invalid
Trackingprovider 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=300

Data 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 store

Deep dives.

ProblemDecision
Query speedprecompute top-k at every prefix node, so a lookup is O(length of prefix), no subtree walk
Freshnessrebuild 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
Sizecap prefix length (~20 chars) and store query IDs instead of strings
Scalingreplicate 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 loaddebounce, reuse the results of a shorter prefix when fewer than k items would match, cache at the edge
Safetyfilter offensive and legally blocked terms at build time; a kill-switch list checked at serve time
Personalizationblend 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.

ProblemDecision
Politenesseach 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))
Priorityfront queues by importance (link-based rank, domain quality) and expected change rate; the Mercator layout (paper (opens in a new tab))
Duplicate URLsnormalize (lowercase host, drop fragments, sort query parameters), then a Bloom filter: a false positive skips one URL, acceptable
Duplicate contentexact: content hash; near-duplicates: SimHash with a small Hamming distance
Trapscaps on depth, URL length and pages per host; detect calendar-like infinite URL patterns
DNSa local caching resolver; DNS lookups are often the hidden bottleneck
Recrawlnext_crawl_at from the observed change rate; conditional GET (If-Modified-Since, ETag)
Distribution and recoverypartition 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.

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 updater

Deep dives.

ProblemDecision
Index choicegeohash: simple, works in any B-tree or Redis GEOSEARCH; quadtree: adapts to density; PostGIS ST_DWithin is enough at modest scale (structures)
Querypick 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 areastoo few results → shorter prefix (bigger cells) until enough
Quadtreesplit any node holding more than ~100 businesses; build in memory at startup, apply updates incrementally
Scalingthe index fits in RAM, so replicate it per region; shard by region only if it stops fitting (and watch dense cities)
Cachingcache 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.

ProblemDecision
Transcodinga 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 bitrateHLS (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
Deliverysegments are immutable: cache forever at the CDN; push popular titles to edge and ISP caches ahead of demand, pull the long tail from origin
Costencode expensive codecs (AV1) only for videos that get views; cold storage tiers for old raw files
Uploadsmultipart with per-part retries, resumable; content hash to dedupe
Statustranscode events update status; the creator is notified when ready
Live variantingest 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 tombstone

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

ProblemDecision
Partitioningconsistent hashing with virtual nodes; replicas on distinct machines and zones (concepts)
ConsistencyN=3; per-request R and W; R + W > N for overlap, W=1 for fastest writes
Conflictsversion 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 failuressloppy quorum + hinted handoff: another node accepts writes and hands them back later
Permanent divergenceread repair on reads; background anti-entropy compares Merkle trees per range
Membershipgossip every second; a phi-accrual failure detector marks nodes down
Deletestombstones kept for a grace period (longer than the repair interval) so deleted data doesn't come back
Adding nodesthe 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

MistakeInstead
Jumping straight to boxesspend 5 minutes on requirements and write them down
Estimating everything to 3 significant figuresorders 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 problemthe simplest design that meets the numbers, then scale the bottleneck
No read path vs write path distinctiontrace one request of each through the diagram
Ignoring failurefor each box: what if it's slow, down, or the data is stale?
Single points of failure left in the final designreplicate, 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 hinttake it, adjust, explain the new trade-off
Silent thinkingnarrate; the interviewer scores reasoning they can hear
Running out of time before a summarykeep the last 3 minutes for wrap-up

Scoring rubric

DimensionWeakStrong
Problem navigationvague scope, missing non-functional requirementscrisp scope with numbers; finds the hard part early
High-level designmissing components or data flow; wrong store for the access patterncomplete read and write paths; each component justified
Technical depthbuzzwords, can't go one level downmechanisms explained (how the cache is invalidated, how order is kept)
Trade-offsone option presented as the only onealternatives compared against requirements; costs admitted
Failure and operationsassumes everything worksfailure modes, monitoring, rollout, cost (expected more at senior+)
Communicationlong silences, messy diagram, ignores hintsstructured, 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

WeekFocusDo
1conceptsread System design concepts and the Primer; DDIA chapters 5–9 (replication, partitioning, transactions, consistency)
2framework + classicsrun the framework on a URL shortener, rate limiter and news feed, timed at 45 minutes, out loud
3breadthchat, notifications, typeahead, crawler, proximity, video, KV store; Alex Xu vol. 1–2 for more designs
4mocks3–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.

  1. "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?"
  2. "Who uses it and how much? Roughly how many daily users, and is it global?"
  3. "What matters most: latency, availability, consistency? Can a user see stale data for a few seconds?"
  4. "I'll leave out D and E unless you want them."
  5. 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.

ComponentAsk
Load balancerL4 or L7? Health checks? Sticky sessions needed (WebSockets)? What if it dies?
Servicestateless? how does it scale? timeouts and retries to dependencies?
Cachewhat's cached, TTL, invalidation on write, hit rate, stampede, hot keys, cache down?
Databasewhy this store? access patterns and keys, indexes, replication, sharding key, hot partitions?
Queuedelivery semantics, idempotent consumers, ordering, DLQ, consumer lag, backpressure?
Search indexhow is it fed (CDC), how stale, how to rebuild?
Object store + CDNpresigned uploads, cache keys, invalidation, egress cost?
Anythingsingle 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

  1. Set a 45-minute timer; talk out loud; draw on a real whiteboard tool.
  2. Afterward, compare with a reference design (Primer, Alex Xu, Hello Interview) and note what you missed.
  3. Score yourself on the rubric; write the 3 gaps.
  4. Redo the same design three days later in 30 minutes.

References