How to Approach a Design
Every system-design problem follows the same skeleton: (1) clarify functional and non-functional requirements, (2) do back-of-the-envelope capacity math, (3) sketch a high-level architecture, then (4) deep-dive the hard parts and their trade-offs. This page works four canonical systems end to end using that skeleton.
Numbers to memorize
1 day ≈ 86,400 s ≈ 10⁵. So 1M writes/day ≈ 12/s; 1B/day ≈ 12K/s. Read:write is often 100:1. One SSD does ~10–50K IOPS; one commodity box ~10K QPS of light work.
Keep the classic latency ladder in mind — it tells you which hops dominate a request budget:
| Operation | Approx latency |
|---|---|
| L1 / main-memory reference | ~1 ns / ~100 ns |
| Read 1 MB from RAM | ~10 µs |
| SSD random read | ~100 µs |
| Round trip within a datacenter | ~0.5 ms |
| Round trip across continents | ~150 ms |
1 · URL Shortener
Requirements
Functional: shorten a long URL to a short one; redirect the short URL to the original; optional custom alias and expiry. Non-functional: very high read availability, low-latency redirects (< 100 ms), short codes must be unique and unguessable enough, no collisions.
Capacity Estimation
Assume 100M new URLs/day → ~1,160 writes/s. Read:write 100:1 → ~116K reads/s. Store URLs 5 years: 100M × 365 × 5 ≈ 180B rows. At ~500 bytes/row → ~90 TB. A base62 code of length 7 gives 62⁷ ≈ 3.5 trillion combinations — plenty.
High-Level Design
client
| POST /shorten GET /{code}
v |
[API / LB] --------------------- [Redirect svc]
| |
[Key Gen] [Cache: Redis] --hit--> 301
| | miss
v v
[DB: code -> longUrl (NoSQL / KV, e.g. DynamoDB/Cassandra)]
Deep Dives
Base62 encoding: map a unique integer ID to [0-9a-zA-Z]. This is compact and collision-free because each ID is unique. Key generation: two options — (a) a distributed counter / ID generator (e.g. a range-allocator or Snowflake-style ID) then base62-encode it; or (b) hash the URL (MD5) and take the first 7 chars, handling rare collisions by re-hashing. Option (a) avoids collisions entirely and is preferred at scale.
Redirect flow: the hot path is a KV lookup; cache aggressively (redirects are read-heavy and immutable). Return 301 for permanent (browser caches, fewer hits) or 302 if you need per-click analytics. DB choice: a partitioned KV store (DynamoDB, Cassandra) keyed by code scales horizontally and needs no joins.
2 · News Feed
Requirements
Functional: users post; a user's home feed shows recent posts from people they follow, ranked; support likes/comments. Non-functional: feed loads fast (< 200 ms), highly available, eventual consistency acceptable, handle celebrities with tens of millions of followers.
Capacity Estimation
Assume 300M DAU, each opens feed ~10×/day → 3B feed reads/day ≈ 35K reads/s. 100M posts/day ≈ 1,160 writes/s. Average 200 followers → fan-out on write does ~200 × 1,160 ≈ 230K feed inserts/s — big, but bounded.
High-Level Design
post --> [Post svc] --> [Post DB]
|
v
[Fanout svc] --(for each follower)--> [Feed cache: Redis lists]
^
read: [Feed svc] --> merge precomputed feed + pull celebrity posts --> rank --> client
Deep Dives
| Model | Fan-out on write (push) | Fan-out on read (pull) |
|---|---|---|
| When | On post, write to all follower feeds | On read, gather from followees |
| Read cost | Cheap (precomputed) | Expensive (query many) |
| Write cost | Huge for celebrities | Cheap |
Celebrity problem: pushing a celebrity post to 50M feeds is wasteful. Use a hybrid: fan-out on write for normal users, but for celebrities skip fan-out and pull their recent posts at read time, merging them into the precomputed feed. Ranking: beyond recency, score by affinity (interaction history), post type, and predicted engagement (an ML model). Storage: posts in a sharded store; feeds as capped Redis lists of post IDs (store IDs, hydrate content on read).
Key insight
Store post IDs in the feed, not full posts — one edit updates one place, and feeds stay tiny. Hydrate content from cache at read time.
3 · Distributed Rate Limiter
Requirements
Functional: allow N requests per window per key (user/IP/API-key); return 429 with Retry-After when exceeded. Non-functional: low latency (< 5 ms added), accurate across many app servers, fault-tolerant, minimal memory.
Capacity Estimation
If the API does 100K QPS and every request checks the limiter, the limiter store must sustain > 100K ops/s with sub-ms latency — a natural fit for in-memory Redis, which handles > 100K ops/s per node.
Algorithms
| Algorithm | Behavior | Trade-off |
|---|---|---|
| Token bucket | Refill tokens at rate R; spend 1/req | Allows bursts, simple |
| Leaky bucket | Queue drains at fixed rate | Smooths, no bursts |
| Fixed window | Counter per clock window | 2× burst at window edge |
| Sliding window log | Timestamps in a sorted set | Accurate, more memory |
| Sliding window counter | Weighted prev+current window | Accurate + cheap |
Deep Dives
Distributed counters: keep counters in Redis so all app servers share state. To avoid race conditions between the read and the increment, run the whole check-and-increment as an atomic Lua script (or use INCR + EXPIRE).
# token bucket in Redis (per key), atomic Lua:
count = INCR key
if count == 1: EXPIRE key window # first hit starts the window
if count > limit: return 429
else: allow
Latency vs accuracy: a strictly-shared Redis is accurate but adds a network hop. For extreme throughput, use local token buckets per node with a fraction of the global quota, periodically reconciled — trading a little accuracy for zero-hop checks. Handle Redis failure by failing open (allow traffic) rather than blocking all requests.
4 · Chat System (WhatsApp / Messenger)
Requirements
Functional: 1:1 and group messaging, delivery receipts (sent/delivered/read), online presence, message history, push notifications when offline. Non-functional: low latency, ordered delivery within a conversation, high availability, exactly-once perceived delivery.
Capacity Estimation
500M DAU × 40 messages/day = 20B messages/day ≈ 230K messages/s. Persistent connections: at peak maybe 100M concurrent WebSockets — spread over servers holding ~65K connections each → ~1,500+ connection servers.
High-Level Design
A <=WebSocket=> [Chat/WS server A] --> [Message svc] --> [Message store]
| |
[Presence] [Kafka queue]
| v
B <=WebSocket=> [Chat/WS server B] <-- routed via [Session registry: user -> server]
(if B offline) --> [Push: APNs/FCM]
Deep Dives
Connection & routing: clients hold a WebSocket to a stateful chat server. A session registry (Redis) maps each user to the server holding their connection so the message service can route a message to the right box. Delivery: persist first, then push; the client ACKs receipt, and the server retries unacked messages. Track sent → delivered → read as separate signals.
Ordering: assign each message a per-conversation sequence number (or a time-sortable ID); the client sorts by it so messages never appear out of order even if the network reorders them. Group chat: for small groups, fan-out the message to each member's connection; for huge groups, a pull model. Storage: a write-heavy store partitioned by conversation ID (Cassandra / HBase) fits the append-only, time-ordered access pattern. Presence: heartbeats update a TTL key; absence of heartbeat → offline; debounce to avoid flicker. Offline: if the recipient has no live connection, queue and send a push notification via APNs/FCM.
Why WebSocket
HTTP polling wastes connections and adds latency. A persistent bidirectional WebSocket lets the server push instantly and keeps one connection per client — the standard for real-time chat.
Storage Choice by System
Notice how the access pattern — not fashion — dictates the datastore in each design:
| System | Store | Why |
|---|---|---|
| URL Shortener | KV (DynamoDB / Cassandra) | Point lookups by code, no joins, huge scale |
| News Feed | Redis lists + sharded post DB | Precomputed feeds in memory, content elsewhere |
| Rate Limiter | Redis counters | Sub-ms atomic ops, TTL expiry built in |
| Chat | Cassandra / HBase | Write-heavy, append-only, partition by conversation |
Cross-Cutting Themes
Across all four: cache the read path, shard by the natural key (code, conversation, user), precompute where reads dominate, prefer eventual consistency where users tolerate it, and always ask "what breaks at 100× scale and how does it degrade?" A good design states its trade-offs explicitly rather than claiming one true answer.
Practice Exercises
- Design Twitter — timelines, tweets, follow graph, search, trending topics.
- Design a web crawler — URL frontier, politeness/robots.txt, dedup, distributed fetching.
- Design YouTube — video upload, transcoding pipeline, CDN delivery, view counts.
- Design Uber — rider/driver matching, geospatial indexing (quadtree/geohash), ETA, surge.
- Design a distributed cache — consistent hashing, replication, eviction (LRU), hot keys.
- Design Google Drive — file storage, chunking, dedup, sync, sharing & permissions.