Design Social & Streaming Systems
Solve complex real-world designs involving massive fan-out, feeds, and media delivery.
SPACED REPETITION Β· 17 practice questions
Make this lesson stick.
Try 3 questions now. No account needed. Sample answers aren't saved.
or sign in to practice all 17One post, two million screens
Suppose you built Snapshot, a photo-sharing app, for your university. It runs on one app server and one PostgreSQL box:
posts(id, author_id, caption, photo BYTEA, likes INT, created_at)andfollows(follower_id, followee_id)- the home feed is a JOIN: posts by everyone I follow, newest first,
LIMIT 20 - a like is
UPDATE posts SET likes = likes + 1 WHERE id = ? - the app asks
/notificationsfor news every 10 seconds
At 50,000 daily users this is a good design. It is simple, correct and cheap, and in an interview you should say so before you change anything. Then the app takes off. These are the assumptions for the rest of the lesson:
| Assumption | Value |
|---|---|
| Daily active users | 20 million |
| Feed loads per user per day | 12 |
| New photo posts per day | 2 million |
| Accounts followed, on average | 300 |
| Photos downloaded per feed load (not already on the phone) | 10, at 200 KB each |
| Users with the app open at the busiest moment | 4 million |
| Peak traffic compared with the daily average | 3Γ |
| The one viral post | 3,000 likes per second |
Predict first: at the evening peak, four things hit the box at once: feed loads, photo downloads, notification polls, and likes on the viral post. Which one breaks it first, and by how much?
Check your answer
The photo bytes. 240 million feed loads a day is 2,778 per second on average and 8,333 at peak, and each load pulls 10 photos of 200 KB:
| Load at peak | Rate | Why one box cannot take it |
|---|---|---|
| Photo bytes | 16.7 GB/s = 133 Gbps | More than five fully saturated 25 Gbps network cards, before the database even reads the bytes off disk |
| Notification polls | 4M Γ· 10 s = 400,000 requests/s | Almost every answer is "nothing new", and news still arrives 5 s late on average |
| Feed JOIN | 8,333 Γ 300 followed accounts = 2.5 million index lookups/s | Every load re-merges the same 300 lists from scratch |
| Viral like counter | 3,000 updates/s on one row | Each update holds the row lock until it commits; at about 1 ms per commit, that row tops out near 1,000 updates/s |
The photo bytes need more than five fully loaded 25 Gbps network cards, and the like counter is limited by one row's lock, not by the machine. The other two are heavy too, but how far over they are depends on hardware you haven't specified yet. Everything strains, and not in the same way. That is the point.
A bigger box can't fix the like counter at all, and it only postpones the other three. Each row is a different pressure, and each pressure has its own standard move. This lesson teaches those moves: the shared toolkit behind feeds, video platforms, chat and live features.
| # | Pressure in the requirements | Move | Full design |
|---|---|---|---|
| 1 | One write must show up for many readers | Decide where fan-out is paid: on write, on read, or both | Twitter/Instagram News Feed |
| 2 | A few keys get most of the traffic | Absorb or spread hot keys: local caches, copies, coalesced misses | this lesson |
| 3 | Items are megabytes, not bytes | Keep bytes away from your servers: object storage, CDN, async processing | YouTube / Netflix Streaming |
| 4 | The server must tell the client first | Hold connections at a gateway tier and route to them | Chat Application (WhatsApp) |
| 5 | Many writers change one number | Make the counter derived data: records, aggregation, sharding | this lesson |
At the end you design live comments for a two-million-viewer stream, which needs four of the five moves at once. The lesson assumes you know the building blocks (load balancers, caches, queues, databases) from Core Building Blocks and can turn daily volumes into requests per second as in Foundations of System Design.
Before any move: read the prompt for pressures
"Design Instagram", "design Twitch" and "design Spotify" sound alike, but they need different parts of the toolkit. Five questions tell you which parts:
| Ask the interviewer | An answer that signals a pressure | Move |
|---|---|---|
| When one user writes, how many users see it? | Hundreds to millions | 1 |
| Is popularity skewed? What are the biggest accounts, posts, streams? | A few items dominate | 2 |
| How big is one item? | Photos, audio, video: KB to GB | 3 |
| Must readers learn about new items without asking? How fast? | "Within a second, while the app is open" | 4 |
| Is there a number that many users change at once? Must it be exact? | Likes, votes, views, "watching now" | 5 |
Add one more question, because it splits real-time features into two very different designs: may a reader miss an item? A chat message must never vanish. A floating heart during a live stream can. The first needs durable storage and acknowledgements; the second can be fire-and-forget. Asking early stops you from building the expensive version for a feature that doesn't need it.
The read:write ratio depends on what you count
"Feeds are read-heavy" is true, but the number you quote depends on the unit. With Snapshot's numbers:
| What you count | Ratio |
|---|---|
| Feed loads : new posts | 240M : 2M = 120 : 1 |
| Posts displayed (20 per load) : new posts | 4.8B : 2M = 2,400 : 1 |
| With fan-out on write: timeline inserts : feed loads | 600M : 240M = 2.5 inserts per load |
The third row is the surprise. Precomputing timelines turns a read-heavy product into a timeline store that takes more writes than reads (cheap ones: each insert appends an ID to a list). Say which ratio you mean, and never quote one as "typical" without the assumptions that produced it.
π‘ The small design is a valid answer. If the interviewer says 50,000 users, the one-box Snapshot is right, and proposing Kafka, sharding and a CDN is a red flag. Every move below is triggered by a number. Say the number, then make the move.
Your turn: which pressures does each prompt contain? (a) A podcast app. (b) Chat for live streams, Twitch-style. (c) A Q&A site where answers are voted on.
Check your answer
(a) Podcast app: media (episodes are tens of MB), hot keys (a top show's new episode is downloaded by millions in its first hour), fan-out (new-episode notifications to subscribers) and counters (play counts). Persistent connections aren't needed; a mobile push is enough.
(b) Live-stream chat: real-time push, fan-out (one message to every viewer in the channel), hot keys (the biggest channels) and counters (viewer counts). Missing a line in a 100,000-viewer channel is acceptable, and that changes the delivery design. The video itself is a separate system.
(c) Q&A votes: counters (exact eventually, one vote per user) and hot keys (a question on the front page). Fan-out is small (notify the asker and the question's followers) and there is no media.
Move 1: Decide where fan-out is paid
Pressure: one write must appear in many readers' views.
When someone you follow posts, the post has to reach your home feed. There are exactly two moments to do that work.
Fan-out on write (push) Fan-out on read (pull)
author posts reader opens feed
β β
βΌ βΌ
fan-out worker βββΊ timeline(A) feed service βββΊ recent posts of followee 1
βββΊ timeline(B) βββΊ recent posts of followee 2
βββΊ ... one insert βββΊ ... one lookup
per follower per followee
β
reader: fetch timeline(me), one read merge newest first, keep 20
Push does the work at write time: it appends the post ID to each follower's precomputed timeline, so a feed load is one read. Pull does it at read time: it fetches recent post IDs from every account you follow and merges them, so a post is one write.
Where the bill lands
Push costs posts Γ followers per author. Pull costs feed loads Γ accounts followed per reader. Every follow edge has one follower and one followee, so across all users the average number of followers equals the average number followed. If authors are typical users, the comparison reduces to feed loads versus posts, which is Snapshot's 120 : 1.
| At Snapshot's peak | Operations per second |
|---|---|
| Push: 69.4 posts/s Γ 300 followers | 20,833 timeline inserts |
| Pull: 8,333 feed loads/s Γ 300 followed accounts | 2,500,000 index lookups |
Push is about 120Γ cheaper here, and each feed load becomes a single read. That is why social feeds are usually precomputed.
The celebrity problem
Averages hide the tail. A singer on Snapshot has 40 million followers, so one post from her is 40 million inserts. Size the fan-out workers generously, say 200,000 inserts per second (about 10Γ the normal peak), and her single post still keeps the whole fleet busy for 200 seconds.
Predict: who notices first?
Check your answer
Not her followers: they see her post a little late and never know it. The people who notice are everyone else who posts during those 200 seconds. Their posts wait in the same queue behind 40 million inserts, so an ordinary user's photo takes minutes to reach their few hundred followers. One account has degraded the product for everyone.
The standard fix is a hybrid: push for ordinary accounts, pull for the few accounts above a follower threshold. At read time the feed service takes your precomputed timeline, fetches recent posts from the handful of big accounts you follow, merges them by time or rank, drops duplicate IDs, and returns the top 20. Millions of feed loads read those big accounts' recent-post lists, so the lists are hot keys, and Move 2 handles them. Set the threshold by cost: push while one post's inserts take only a few seconds of fan-out capacity, and remember that every pulled account adds a lookup to each of its followers' feed loads. Many designs also skip followers who haven't opened the app for weeks and rebuild their timelines by pull when they return.
Details that decide correctness
- Store IDs, not content. A timeline entry is
(post_id, author_id, timestamp). The feed service hydrates the top 20 with one batched lookup in the post store, which is cached. Edits and deletes then take effect everywhere at once, and hydration drops IDs whose post is gone. Copy the content into 40 million timelines instead, and a delete becomes 40 million writes. - A post fans out to many timelines; a follow backfills one. When Alice follows Bob, only Alice's timeline gets Bob's recent posts. An unfollow filters Bob out of hers.
- Keep three shapes of the same data. The post store is keyed by
post_id(for hydration). The per-author post list is partitioned byauthor_idand sorted by time (for pull and backfill). The timeline is keyed by the reader'suser_id(for push). Shard posts only bypost_id, and "Bob's recent posts" becomes a query to every shard. - Put the fan-out behind a log. The post service commits the post and publishes a
post_createdevent through an outbox table or change data capture, so a crash between the two can't lose it (Move 5 explains why). The fan-out worker, notifications and the search indexer each consume it independently: in Kafka as separate consumer groups, on AWS as an SNS topic that feeds one SQS queue per consumer (one SQS queue hands each message to one consumer, not to all of them). Messaging & Queues covers delivery guarantees and ordering. - The social graph is two adjacency lists. Store "who follows X" and "whom X follows" in a sharded key-value or relational store with a cache in front. Facebook has described running its social graph this way with TAO, a graph-aware cache over sharded MySQL. Reach for a graph database when you need multi-hop queries such as friends-of-friends suggestions, and even there each hop costs work in proportion to the number of edges it touches.
Failure modes
- A worker dies halfway through a fan-out, and the event is redelivered. If an insert appends to a list, some followers get the post twice. Make the insert idempotent instead: in a Redis sorted set with
post_idas the member, adding the same member again only rewrites its score. Retries become safe. - A cache node holding timelines dies. Those users have nothing precomputed. Rebuild each timeline by pull on its next load, limit the rebuild rate so one lost node doesn't become a pull storm, and serve a simpler feed meanwhile.
- Fan-out lags. Followers seeing a post late is fine. The author not seeing their own post is not. Insert the author's post into their own timeline synchronously (or have the client add it), so they read their own write even while fan-out catches up.
Say it like this: "Feed loads outnumber posts about 120 to 1, so I'll precompute timelines with fan-out on write and store only post IDs. Accounts above a follower threshold are pulled at read time and merged in, so one celebrity post never becomes 40 million writes. I pay for that with write amplification and a second read path."
Likely follow-up: "What happens when a user crosses the threshold?" Stop pushing their new posts; their followers' timelines already hold the older ones, the read-time merge picks up the rest, and deduplicating by post ID hides the overlap.
Your turn: change the graph. A company announcements app has 100,000 employees and 5,000 channels (teams and bots). Employees follow 2,000 channels on average. The channels publish 500,000 posts a day (build bots are chatty), and employees load their feeds 300,000 times a day. Push or pull?
Check your answer
Pull. The shortcut "average followers = average followed" only holds when the same people are both authors and readers. Here 200 million follow edges (100,000 Γ 2,000) land on just 5,000 channels, so the average channel has 40,000 followers.
| Per day | Operations |
|---|---|
| Push: 500,000 posts Γ 40,000 followers | 20 billion inserts |
| Pull: 300,000 feed loads Γ 2,000 channels | 600 million lookups |
Pull is about 33Γ cheaper. Cache each channel's recent-post list (5,000 lists, read over and over) and merge at read time. You give up single-read feed loads: a load is slower for someone who follows many channels.
Move 2: Absorb or spread hot keys
Pressure: a few keys get most of the traffic.
Hash partitioning spreads keys evenly across shards. It does not spread traffic. Every read of one key goes to the one node that owns it (plus its replicas, if you read from them). And popularity in social systems is heavily skewed: one post, one account or one live stream can take a large share of all requests.
Say the winning goal of a cup final is posted on Snapshot and embedded in millions of feeds. Its record, post:42 (caption, author, counts), is read 500,000 times per second. Assume one cache node serves about 100,000 small reads per second. That is a common planning figure for an in-memory store, not a promise; use the interviewer's number if they give one.
Predict: you add 20 nodes to the cache cluster. What happens to the load on the node that owns post:42?
Check your answer
Nothing useful. The key still lands on exactly one node (perhaps a different one, if the resharding moved it), and that node still receives all 500,000 reads per second, five times what it can serve. The other nodes sit mostly idle. Scaling out helps when load is spread across many keys; one hot key needs a move of its own.
The hot-key toolbox
| Move | How | Numbers for post:42 |
Cost |
|---|---|---|---|
| Local cache | Each app server keeps hot keys in its own memory for 1 s | 300 servers β 300 reads/s reach the cache node | Up to 1 s stale; you need to know which keys are hot |
| Copies of the key | 8 copies, post:42#0 to post:42#7, on different nodes; each read picks one at random |
62,500 reads/s per copy | Every write updates 8 copies, which briefly disagree |
| Split the key | For write-hot keys: N sub-keys, summed on read | see Move 5 | Every read touches N keys |
| Coalesce misses | Only one fetch per key at a time; other requests wait for it | see below | Waiters share one fetch's latency and failure |
| Special-case the class | Celebrities pulled instead of pushed (Move 1); big live events get dedicated capacity | Two code paths to maintain |
Why it works: a local cache replaces "reads per second of a popular key", a number set by popularity, with "servers Γ· TTL", a number you choose. Popularity can grow 10Γ and the cache node won't notice.
How do you know a key is hot? Either you know the class in advance (accounts above a follower threshold, scheduled live events) or you measure it: sample requests, keep top-K counts over a few seconds, and promote the winners into the local cache or into copies.
The thundering herd
Hot keys cause a second, nastier problem. Without the local cache, suppose post:42's cache entry expires at noon and the database query that rebuilds it takes 20 ms. At 500,000 reads per second, 10,000 requests miss during those 20 ms, and every one of them sends the same query to the database. This is a cache stampede, or thundering herd. It can take down a database that was sized for the trickle of misses a healthy cache lets through.
The fix is to let one request per key do the work. Inside one server, that is request coalescing, often called single-flight:
import threading
import time
class SingleFlight:
"""Collapse concurrent loads of the same key into one call."""
def __init__(self):
self._lock = threading.Lock()
self._inflight = {} # key -> (done event, result box)
def do(self, key, load):
with self._lock:
entry = self._inflight.get(key)
leader = entry is None
if leader: # first caller for this key
entry = (threading.Event(), {})
self._inflight[key] = entry
done, box = entry
if leader:
try:
box["value"] = load(key)
except Exception as exc: # followers see the same failure
box["error"] = exc
finally:
with self._lock:
del self._inflight[key] # the next miss starts a new flight
done.set()
else:
done.wait()
if "error" in box:
raise box["error"]
return box["value"]
db_calls = 0
db_lock = threading.Lock()
def load_from_db(key):
global db_calls
with db_lock:
db_calls += 1
time.sleep(0.2) # a slow query
return f"row for {key}"
def stampede(n, get):
start = threading.Barrier(n)
def worker():
start.wait() # everyone misses at the same instant
get("post:42")
threads = [threading.Thread(target=worker) for _ in range(n)]
for t in threads:
t.start()
for t in threads:
t.join()
flight = SingleFlight()
stampede(1000, load_from_db)
print("without coalescing:", db_calls, "database calls")
db_calls = 0
stampede(1000, lambda key: flight.do(key, load_from_db))
print("with coalescing: ", db_calls, "database call")
It prints 1000 database calls without coalescing and 1 with it. In a real service, load would read the database and refill the cache, and the first line of get would still be a cache lookup.
Across 300 app servers, each server still makes its own call, so 300 queries reach the database. To get that down to one, coordinate through the cache: the first request to miss gets a short lease to refill the key, and the others wait briefly or take the stale value. Facebook has described exactly this mechanism, under that name, in Scaling Memcache at Facebook (NSDI 2013). Two cheaper habits help as well: refresh popular keys a little before they expire, and add random jitter to TTLs so that a million keys written in the same second don't all expire in the same second.
CDNs apply the same idea to media: CloudFront Origin Shield, for example, consolidates simultaneous requests for an uncached object so that as few as one reaches your origin.
Failure mode: the cache tier goes away
Suppose the cache absorbs 99% of 300,000 reads per second, so the database sees 3,000 per second and is sized for about 10,000. If the whole cache tier fails, the database receives 300,000 reads per second, 30Γ its capacity. It doesn't get slower; it falls over, and a cache incident becomes a full outage. Plan for it:
- Spread cache nodes and their replicas across availability zones, so one failure loses a slice of the cache, not all of it.
- Put a circuit breaker and load shedding in front of the database: serve stale data, hide like counts or show a simpler feed rather than letting every request through.
- Warm a cold cache gradually, and never flush the whole cache as part of a deploy.
- Data that lives only in the cache, such as precomputed timelines, has no database copy to fall back on. Rebuild it by pull at a limited rate (Move 1).
Caching strategies themselves (cache-aside, write-through, eviction, invalidation) are covered in Databases & Storage.
Say it like this: "Hashing spreads keys, not traffic, so one viral post can overload a single cache node however many nodes I add. For the few hot keys I'll keep a one-second in-process cache on the app servers and coalesce concurrent misses, so an expiry can't stampede the database. The cost is up to a second of staleness on those keys."
Likely follow-up: "How do you find the hot keys?" Sample and count per window, or mark known classes ahead of time. Then: "What if the hot key is written, not read?" Split it, which is Move 5.
Your turn: after a surprise announcement, a celebrity's profile record (bio, follower count, latest 12 posts) gets 400,000 reads per second for ten minutes. You run 400 app servers. Edits to the bio must be visible within 5 seconds. Which move, and with what setting?
Check your answer
A local cache with a short TTL, plus coalescing. With a 2 s TTL, each of the 400 servers refreshes the key once every 2 s: 400 Γ 0.5 = 200 reads/s reach the cache node, about 2,000Γ fewer than the raw traffic. A bio edit becomes visible within about 2 s plus the time the write takes to reach the cache, inside the 5 s requirement. A 5 s TTL would be cheaper (80 reads/s) but would spend the whole staleness budget and leave nothing for the write path. Coalescing makes each expiry cost one fetch per server. Copies of the key aren't needed.
Move 3: Keep bytes away from your servers
Pressure: items are megabytes, not bytes. Bandwidth and storage, not requests per second, set the size of the system.
Snapshot's photos already broke the one-box design with 133 Gbps at peak. Storage grows fast too. Each post keeps a 3 MB original and three renditions: 1080 px at 200 KB, 640 px at 80 KB and a 150 px thumbnail at 10 KB, so 3.29 MB in total.
| Per day (2 million posts) | Per year |
|---|---|
| 2M Γ 3.29 MB = 6.58 TB | 2.4 PB |
Split that total and a pattern appears. Originals are 2.19 PB a year, 91% of the bytes, and are almost never read after processing. Renditions are 0.21 PB a year, yet they carry nearly all the egress. The two want different homes.
Bytes and metadata part ways
| Data | Home | Why |
|---|---|---|
| Media record: ID, owner, status, dimensions, object keys | Database | Small, queried, updated in transactions |
| Renditions (image sizes; for video, segments plus one manifest per video) | Object storage behind a CDN | Read constantly, never modified |
| Original upload | Object storage, a colder tier after processing | Kept for reprocessing, rarely read |
The database never holds bytes, and your API servers never move them. Here is the upload path:
Client API Object storage Queue + worker
β POST /media β β β
ββββββββββββββββββββββΊβ auth; row = uploading β β
βββ presigned URL βββββ€ β β
β PUT bytes ββββββββββΌββββββββββββββββββββββββββΊβ β
β POST /media/42/complete β β
ββββββββββββββββββββββΊβ row = processing; job ββββΌββββββββββββββββββββΊβ
βββ 202 processing ββββ€ ββββ read original βββ€
β β ββββ renditions ββββββ€
β ββββββ row = ready, rendition keys ββββββββββββββ€
βββ push: ready βββββββ€ β β
- The API authenticates the user, inserts a media row with status
uploading, and returns a presigned URL: a URL that lets this client write one object for a few minutes. For large files it presigns the parts of a multipart upload instead (S3 parts are 5 MiB to 5 GiB, up to 10,000 per upload), so a dropped connection costs one part, not the whole file. The tus protocol and YouTube's resumable upload API solve the same problem. - The client sends the bytes straight to object storage.
- The client reports completion (or the storage service emits an event). The API checks that the object exists, sets the row to
processing, enqueues a processing job, and answers 202 Accepted with statusprocessing. 202 means "received, not finished yet". - A worker checks the actual bytes (not the content type the client claimed), strips location metadata (EXIF) from photos, writes renditions under new, immutable keys such as
media/42/v1/640.webp, and marks the rowready. - The client hears "ready" by push or by polling. Nobody else can see the post until then.
Viewers get a small JSON document with CDN URLs from your API and fetch the bytes from the CDN. On a miss, the edge fetches from object storage (through a shield tier if there is one) and keeps a copy for the next viewer.
Cache forever, never overwrite
Content under a key never changes, so the CDN and browsers may cache it for a year: Cache-Control: public, max-age=31536000, immutable. That is 365 days, and immutable tells clients not to revalidate while the copy is fresh. A new profile photo gets a new key (v2), and the old URL simply stops being referenced. Keep purging for takedowns and moderation. Cloudflare includes purge on every plan, while CloudFront charges per invalidation path after the first 1,000 paths each month. Either way, versioned URLs make purging the exception.
Private media works the same way with signed URLs or signed cookies that expire quickly. The CDN checks the signature at the edge, so the bytes still come from its cache.
Why the CDN pays for itself twice
Bandwidth. At a 95% edge hit ratio, object storage serves 6.7 Gbps of Snapshot's 133 Gbps peak; at 99%, 1.3 Gbps.
Latency. Light in optical fiber covers about 200,000 km per second, so every 1,000 km of distance adds at least 10 ms of round trip, and real routes are longer than straight lines. A new HTTPS connection needs three sequential round trips before the first byte arrives: the TCP handshake, the TLS 1.3 handshake and the request itself (HTTP/3 over QUIC folds the first two into one). At 150 ms per round trip to a distant origin, that is 450 ms per new connection; at 20 ms to a nearby edge, it is 60 ms.
CDNs steer users to a nearby edge in one of two ways. With GeoDNS, the DNS answer depends on where the query comes from: the user's resolver, or the user's subnet when EDNS Client Subnet is used. With anycast, every edge announces the same IP address, and internet routing (BGP) delivers packets to a nearby one. Networking Basics covers DNS and CDNs in more detail.
Storage tiers: never make a viewer wait for a restore
Tiering saves real money at petabyte scale, as long as nothing a user can open lands in a tier that needs a restore first. S3's storage classes sort cleanly:
| Class | First byte | Fits |
|---|---|---|
| Standard | milliseconds | New and popular renditions |
| Standard-IA, Glacier Instant Retrieval | milliseconds (retrieval fee; 30- and 90-day minimums) | Older renditions that must still open instantly |
| Intelligent-Tiering | milliseconds in its default tiers | Access patterns you can't predict |
| Glacier Flexible Retrieval, Glacier Deep Archive | minutes to hours, after a restore request | Originals and backups the product never serves directly |
Standard-IA and Glacier Instant Retrieval also bill any object under 128 KB as 128 KB, so a 10 KB thumbnail costs more there than in Standard: leave small renditions where they are.
β οΈ A lifecycle rule that moves year-old renditions to Deep Archive means that someone scrolling back to an old post waits hours for it to open. Originals are the right candidates: 91% of Snapshot's bytes, read only by reprocessing jobs that can wait.
Also don't multiply object-storage estimates by 3 "for replication" out of habit. S3 Standard already stores data across at least three availability zones and is designed for 99.999999999% durability, and that is included in the price. Multiply for systems you replicate yourself, and add cross-region copies only if the requirements include surviving the loss of a region.
Failure modes
- Media served through the app servers. Memory isn't the problem, because streaming code sends files through small buffers. Bandwidth and connection slots are. 133 Gbps is more than five saturated 25 Gbps network cards, and a slow phone holds a connection (and, in thread-per-request frameworks, a worker) for its whole download. When a machine's network card saturates, every API call on that machine (login, post, like) times out too, so a media problem becomes a full outage. Keeping media on separate infrastructure is a bulkhead: one workload can no longer exhaust a resource that another depends on. Users far away also pay the long round trips on every photo.
- Viral media. A new photo's first requests miss at many edges at once. Request collapsing at each edge plus a shield tier means the origin sees about one request per object, not one per viewer. That is Move 2 again, inside the CDN.
- A worker dies mid-job. Queues redeliver, so jobs run at least once. Make them idempotent: rendition keys are deterministic (
media/42/v1/640.webp), so rewriting one is harmless, and the status update is conditional (SET status = 'ready' WHERE id = 42 AND status = 'processing'). A file that crashes every worker goes to a dead-letter queue after a few attempts, and its row becomesfailed. - Presigned URLs abused. Keep them short-lived. With a presigned POST policy you can cap the size (
content-length-range) and the content type. Validate the real bytes in the worker, and keep everything private untilready.
Video has the same shape with heavier processing: a ladder of resolutions and bitrates, short segments, and one manifest per video listing every rendition so the player can switch as bandwidth changes. YouTube / Netflix Streaming builds that pipeline.
Say it like this: "Bytes never touch my API servers. Clients upload straight to object storage with a presigned URL, a queue drives processing, and viewers fetch immutable, versioned URLs from the CDN. The database holds metadata and object keys. I pay with a multi-step upload flow and a 'processing' state the UI has to show."
Likely follow-up: "What stops someone uploading a 20 GB file, or malware, through that URL?" Expiry, size and type conditions on the presigned upload; validation and scanning in the worker; and nothing public before ready.
Your turn: Snapshot adds private family albums, each shared with at most 20 people. What changes in the media design, and what does the CDN still buy you?
Check your answer
The bucket stays private and the CDN becomes its only reader. The API hands short-lived signed URLs (or cookies) only to album members, and the CDN verifies them at the edge. Object keys should be unguessable, and a leaked link should expire within minutes.
The CDN's offload shrinks: each photo is seen by at most 20 people, often in different cities, so many requests are first misses. It still buys latency (a nearby edge, reused connections) and absorbs repeat views by the same family. Album photos (apart from small thumbnails, which Standard-IA bills as 128 KB) can move to Standard-IA after 30 days, because they still open in milliseconds, but not to Deep Archive: a grandparent scrolling back a year should not wait hours.
Move 4: Hold connections and route to them
Pressure: the server has to tell the client something within about a second, without being asked.
Snapshot's polling costs 400,000 requests per second at peak, nearly all of them empty, and news still arrives 5 s late on average. Flip it around: keep one connection open per online user and send only when there is news. That is 4 million open connections instead of 400,000 requests every second.
Pick the channel by direction and by whether the app is open. Networking Basics compares the transports.
| Need | Choose |
|---|---|
| Server to browser only (notifications, live counts) | Server-Sent Events: plain HTTP, automatic reconnect, can resume from the last event ID |
| Frequent messages in both directions (chat, games) | WebSocket |
| Networks that break long-lived connections | Long polling, as a fallback |
| App closed or in the background | Mobile push (APNs, FCM) |
β οΈ Over HTTP/1.1, browsers allow only six connections per domain, shared across all tabs, so an SSE stream per tab runs out quickly. Serve SSE over HTTP/2.
The routing problem
Connections live on a tier of gateway servers. LinkedIn has described holding about 100,000 persistent connections per frontend machine for live-video likes (QCon London 2020). Take that as an assumption: 4 million connections need 40 gateways, or about 60 with 50% headroom for failures and deploys.
Gateways are stateful. Alice's connection is on gw-7 and nowhere else. When the notification service has something for Alice, how does it reach gw-7?
alice ββSSEβββΊ gw-7 ββ "alice is on gw-7" (on connect, then every 30 s) βββ
β² βΌ
β βββββββββββββββββββββββββββββββββββ
β β connection registry, TTL 60 s β
β β alice β gw-7 bob β gw-2 β
β ββββββββββββββββββ¬βββββββββββββββββ
β β 1. look up alice
ββββββββ 2. deliver to gw-7 βββββββ notification service
| Routing option | How it works | Watch out for |
|---|---|---|
| Registry plus direct send | Gateways register their users with a heartbeat and TTL; senders look up the gateway and call it (or publish to that gateway's own channel) | One lookup per message; entries go stale when a gateway dies |
| A pub/sub channel per user | gw-7 subscribes to user:alice when she connects; senders publish without knowing the gateway |
Redis Pub/Sub is at-most-once: a message published while gw-7 is reconnecting is lost |
| A log partition per gateway | Each gateway consumes its own topic or partition; senders look up the gateway and write there | Kafka is not built for a topic per user, so key by gateway, not by user |
May the reader miss one?
The question from the start of the lesson now picks the design:
- No (chat messages, notifications that change state): write the event to a durable per-user inbox with an increasing sequence number first. The push is only a hint. The client acknowledges the highest sequence number it has, and on reconnect asks for "everything after 1,042". Delivery is at-least-once, and the client drops duplicates by ID. Chat Application (WhatsApp) builds this in full.
- Yes (typing indicators, floating hearts, "watching now"): fire-and-forget pub/sub is enough. A reconnecting client fetches a fresh snapshot instead of replaying what it missed.
Presence-aware routing is a small extension: user online β send over the connection; offline and important β mobile push; offline and unimportant β a digest later. It has a race: Alice can go offline between the presence check and the send. That is why important events go to the inbox first. A send that fails, or reaches a dead gateway, is repaired from storage when she reconnects instead of being lost.
Failure modes
- Reconnect storm. A gateway with 100,000 connections crashes, and all 100,000 clients reconnect at once, hammering the load balancer, authentication and the registry. Clients must back off exponentially with random jitter: spread over 30 seconds, that is about 3,300 reconnects per second instead of 100,000 in one. Deploys should drain gateways a few at a time, never restart them all together.
- Stale registry entry. gw-7 died, but
alice β gw-7survives until its TTL expires. Sends to it fail, and durable events wait in the inbox. The heartbeat TTL bounds how long this lasts. - Slow client. A phone on a weak network can't keep up. Give each connection a bounded buffer. For fire-and-forget data, drop or merge updates and keep only the latest; for durable data, disconnect the client and let it catch up from the inbox.
Say it like this: "Online users hold a WebSocket connection to a gateway tier, and a registry maps each user to a gateway. Chat messages go to a durable inbox first and the push is just a nudge, so nothing is lost when a gateway dies; typing indicators are fire-and-forget. What I give up is a stateful tier that I have to drain carefully on every deploy."
Likely follow-up: "What happens to messages sent while the phone switches from Wi-Fi to mobile data?" They wait in the inbox. The client reconnects with its last sequence number and receives them.
Your turn: a trading app streams prices to 500,000 clients. Each client watches up to 50 tickers, and each price changes up to 10 times a second. Only the latest price matters. Design the delivery.
Check your answer
Forwarding every change could mean 500,000 Γ 50 Γ 10 = 250 million messages per second. Two moves shrink it:
- Conflate. Only the latest price matters, so each connection gets at most one message per second with the latest price of each of its tickers: 500,000 messages per second. A slow client simply gets the newest values.
- Fan out in two levels. Each gateway subscribes once to each ticker that any of its clients watches and keeps a local map from ticker to connections. A price update crosses the network once per interested gateway, not once per client.
At-most-once delivery is fine, because a missed tick is replaced a second later, so there is no inbox. What you give up: clients never see the intermediate prices within a second.
Move 5: Make the counter derived data
Pressure: many writers change one number, and everyone reads it.
Snapshot users hand out 300 million likes a day: 3,472 per second on average and 10,417 at peak. Spread over millions of posts, that is easy. The viral post's 3,000 likes per second on one row isn't. Each UPDATE posts SET likes = likes + 1 holds that row's lock until its commit is durable. If a commit takes about 1 ms, the row serializes at roughly 1,000 updates per second however large the machine is. The rest queue up holding connections, and the queue backs up into the API.
It gets worse: a counter alone can't answer the questions the product asks. Did I already like this? Can I unlike it? Was that double tap one like or two?
Records are the truth; the count is derived
- Store the like itself. Use
likes(user_id, post_id, created_at)with primary key(user_id, post_id). A double tap or a retried request hits the existing key and changes nothing, so liking is idempotent. Partition byuser_id: the viral post's 3,000 likes per second come from 3,000 different users, so they spread across every partition, and "did I like this?" is a single-partition lookup. - Emit a change only when the state changed. A new row emits +1, a delete that removed a row emits β1, and a no-op emits nothing. So that an event isn't lost when the process dies between the database write and the publish, emit events through an outbox table or change data capture (Advanced Topics & Final Prep covers the outbox pattern).
- Aggregate. Events go to a log partitioned by
post_id. A consumer sums the deltas per post for one second, then writes oneUPDATE post_counts SET likes = likes + :delta. The viral post's 3,000 events become one write per second. - Make the flush idempotent, per event. A consumer can crash after updating the counts but before the log records its position. On restart it resumes from its last saved position, which can be well behind, and its one-second flushes now cut the stream at different places than before. So guard events, not batches: store, for each log partition, the last offset already counted, in the same transaction as the counts, and skip every event at or below it. (Equivalently, on restart, seek each partition to the stored offset + 1 instead of trusting the log's saved position.)
from collections import defaultdict
class CountTable:
"""Stands in for a table of like counts plus, for each log partition, the
last offset already counted. One transaction writes both."""
def __init__(self):
self.likes = defaultdict(int)
self.applied = {} # partition -> last offset counted
self.row_writes = 0
def apply(self, partition, first_offset, events):
done = self.applied.get(partition, -1)
deltas = defaultdict(int)
counted = 0
for offset, (post_id, delta) in enumerate(events, start=first_offset):
if offset > done: # per event, not per batch
deltas[post_id] += delta
counted += 1
for post_id, delta in deltas.items(): # one transaction:
self.likes[post_id] += delta # UPDATE ... likes = likes + delta
self.row_writes += 1
last = first_offset + len(events) - 1
self.applied[partition] = max(done, last) # ...and the offset, together
return counted, len(events) - counted
# Partition 0 of the like-event log: a viral post, 40 unlikes, a long tail of 500 posts.
log = [("viral", +1)] * 3000 + [("viral", -1)] * 40 + [(f"p{i}", +1) for i in range(500)]
def consume(table, start, batch, stop=len(log)):
for first in range(start, stop, batch):
events = log[first:min(first + batch, stop)]
counted, skipped = table.apply(0, first, events)
print(f"offsets {first}-{first + len(events) - 1}: counted {counted}, skipped {skipped}")
table = CountTable()
consume(table, start=0, batch=1000, stop=3000) # three one-second flushes, then a crash
print("crash: the log's saved position is still 1000, and new flushes cut elsewhere")
consume(table, start=1000, batch=1500)
print("viral likes:", table.likes["viral"], "| row writes:", table.row_writes, "for", len(log), "events")
offsets 0-999: counted 1000, skipped 0
offsets 1000-1999: counted 1000, skipped 0
offsets 2000-2999: counted 1000, skipped 0
crash: the log's saved position is still 1000, and new flushes cut elsewhere
offsets 1000-2499: counted 0, skipped 1500
offsets 2500-3539: counted 540, skipped 500
viral likes: 2960 | row writes: 504 for 3540 events
The restart re-read 2,000 events that had already been counted, in batches that didn't line up with the originals. The flush from 2,500 to 3,539 straddled the stored offset: its first 500 events were skipped and the rest counted. The final count is exact, 3,000 likes minus 40 unlikes. A guard that compares only a batch's last offset would have applied that straddling batch in full and counted offsets 2,500β2,999 twice (3,460 instead of 2,960). The viral post's 3,040 events cost 4 row writes; the 500 long-tail posts cost one write each. Aggregation only compresses keys that are hot, which are exactly the keys that needed it.
Why it works: writes to the source of truth spread by user, and the only per-post write happens once per flush interval, a rate you choose rather than one set by popularity.
The offset guard handles replays by the consumer. It doesn't catch the producer publishing one event twice, which an outbox relay can do if it crashes after publishing but before marking the row as sent. Give each event an ID and drop duplicates within a time window, and run a periodic reconciliation job that recounts like records for all hot posts and a sample of the rest, correcting any drift. Because like records are partitioned by user, a recount per post is a batch job: an offline scan of the records, or a query on a secondary index keyed by post.
The read side
Counts are read constantly, so cache them and round them for display: "12.4K" hides a few seconds of lag. The one number a user checks closely is their own like, and that comes from their own action or their like record, not from the count. So the heart turns red instantly (the user reads their own write), and the total catches up a second later. If people viewing an open post should see the count move, push the aggregated count every second or two over the connection from Move 4, never one message per like.
Other counter shapes
| Approach | Writes at the viral post | Read cost | Freshness | Exact? |
|---|---|---|---|---|
UPDATE one row |
3,000 row locks/s, which fails | 1 read | Immediate | Yes |
| Sharded counter: N sub-counters, increment a random one | 3,000 Γ· N per key | N reads, summed | Immediate | Yes, if no increment is lost |
| Like records, log and aggregator | 1 write/s | 1 read | After the flush and consumer lag | Yes, after dedupe and reconciliation |
| HyperLogLog sketch | Constant per add | Constant | Immediate | No: about 0.81% standard error, and it counts distinct items only |
A sharded counter is the move when there is no log: keep votes:final:0 through votes:final:15 on different nodes, INCR a random one, and sum all 16 on read. β οΈ In Redis Cluster, keys that share a hash tag (the part in braces, as in votes:{final}:3) land in the same slot, and so on the same node, which defeats the purpose. Multi-key commands such as MGET only work within one slot, so the read becomes 16 separate GETs.
A HyperLogLog answers "how many different users viewed this today?" in at most 12 KB with a 0.81% standard error (Redis docs). It can't tell a user whether they are counted, and it can't subtract, so it is wrong for likes.
Likes are also an abuse target: bots inflate counts. Rate-limit per user and per device before an event is counted. Design URL Shortener & Rate Limiter covers the limiter.
Say it like this: "The source of truth is one like record per user and post, which makes likes idempotent and lets people unlike. The count is derived: each real state change emits an event, and an aggregator applies one update per post per second, storing each log partition's last counted offset in the same transaction and skipping events at or below it, so a replay can't double-count. The public count lags by a second or two; the user's own heart updates instantly from their like record."
Likely follow-up: "Can the count drift or go negative?" Only through lost or duplicated events: dedupe by event ID, and reconcile by recounting records.
Your turn: votes, not likes. A Q&A site lets users upvote, downvote, or change their vote. A hot question takes 20,000 vote actions per second, and its score must be exact eventually. What does the event carry when a user switches from an upvote to a downvote?
Check your answer
Store one vote record per (user_id, question_id), with a value of β1, 0 or +1. Each change emits delta = new β old, so switching from +1 to β1 emits β2, and withdrawing an upvote emits β1. The aggregator sums deltas per question for a second: 20,000 actions become one write per second. Sending the same vote again leaves the record unchanged and emits nothing. You give up an instantly exact public score; the voter's own arrow updates immediately from their record.
Boss level: live comments on a two-million-viewer stream
Time to combine the moves. The prompt: design live comments and reactions for a streaming platform. You ask your questions and get these answers:
| Requirement | Value |
|---|---|
| Live streams at once | 10,000; most have fewer than 100 viewers |
| Biggest stream | 2 million concurrent viewers |
| Comments on the biggest stream at peak | 5,000 per second |
| Reactions (hearts) on the biggest stream at peak | 200,000 per second |
| Comment delivery | On screen within 2 s |
| Late joiners | See the last 50 comments |
| Durability | Every comment stored (for replays and moderation); a viewer may miss some comments, for example while reconnecting |
Predict first: which number makes the naive design, "send every comment to every viewer", impossible?
Check your answer
5,000 comments per second Γ 2 million viewers = 10 billion deliveries per second on one stream. No budget covers that. Notice also that nobody can read 5,000 comments per second. That is the way out.
Step 1: shrink the problem with the product
The chat panel scrolls at a readable speed, say 10 comments per second. So the system broadcasts a selection: 0.2% of comments on the biggest stream. Everyone can still post, every comment is stored, and a stream with 30 viewers and one comment per second passes everything through. Choosing which comments to show is a product decision: always include the streamer and moderators, prefer comments from verified accounts or with many reactions, and sample the rest.
This changes a requirement, so say it out loud: "Nobody can read 5,000 comments a second. I'll broadcast about 10 per second, chosen by a selector, and store all of them. Is that acceptable?"
Step 2: fan out in two levels, in batches
At the 100,000-per-gateway figure from Move 4, 2 million viewers would fill 20 gateways, leaving nowhere for a failed gateway's viewers to go. Move 4's 50% headroom gives 30 gateways at about 67,000 viewers each. On a shared platform the load balancer doesn't group viewers by stream, so S's viewers may be spread over even more gateways; the dispatcher simply sends to every gateway that holds at least one.
viewer ββPOST commentβββΊ Comment API βββΊ comment store (write-sharded)
β
ββββΊ log, partitioned by stream_id
β
βΌ
selector for stream S: every 500 ms, pick 5 comments,
add the reaction total, build one batch
β
βΌ
dispatcher: which gateways have viewers of S? (30 of them)
ββββββββββββ¬βββββββββββββββ΄βββββββββββββ¬βββββββββββ
βΌ βΌ βΌ βΌ
gw-1 gw-2 ... gw-29 gw-30
67k viewers 67k viewers 67k viewers 67k viewers
- Level 1: the dispatcher keeps, per stream, the set of gateways that have at least one viewer of it. One batch crosses the network about 30 times, not 2 million times.
- Level 2: each gateway keeps a local map from stream to its own connections and writes the batch to each of them.
- Batching: every 500 ms the selector emits one message with 5 comments and the reaction total, so each viewer receives 2 messages per second.
| Delivery load | Value |
|---|---|
| Messages to viewers | 2M Γ 2 per second = 4 million/s, 2,500Γ fewer than naive |
| Per gateway | about 67,000 Γ 2 per second = about 133,000 messages/s |
| Per-gateway bandwidth at 1 KB per batch | about 133 MB/s = about 1.1 Gbps |
LinkedIn has described this two-level shape for live-video likes: frontend machines keep an in-memory table of their own connections, and dispatchers keep a table of which frontend machines subscribe to which video, held in a key-value store so that losing a dispatcher doesn't lose it (QCon London 2020).
Viewers only receive over the connection; they post comments with ordinary HTTP requests. So SSE is enough, though a WebSocket works too. Delivery is at-most-once by design. That is acceptable because the requirement allows a viewer to miss some comments and the broadcast is only a sample anyway. A missed batch's comments are gone for that viewer; only its reaction total is superseded by the next batch.
Step 3: reactions are a counter
200,000 hearts per second must not become 200,000 messages. Aggregate in two levels, the mirror image of the fan-out. Each of the 200 API servers that receive reactions keeps a running total per stream and flushes it every 500 ms, so 400 messages per second reach the stream's aggregator. At peak the batch carries "+100,000 hearts in the last half-second", and each client animates a handful of hearts in proportion and updates the total. This is Move 5 without per-user records: a viewer may send many hearts, and nobody takes a heart back.
Step 4: the hot keys
The big stream is a hot key everywhere it appears:
- Comment store writes: 5,000 per second, all for one stream and all landing at the newest end of its time-ordered key. If your store caps throughput per partition (DynamoDB documents 1,000 write units per second per partition), write-shard the key: 5,000 small writes per second need at least 5 shards, so use
(stream_id, shard 0β7)with time as the sort key to leave headroom. Reading every comment for a replay then merges 8 partitions, which is fine offline. - "Last 50 comments": every joining viewer reads it. Keep it as a capped list in memory per stream, and let each gateway cache it for a second: about 30 reads per second, however fast viewers join.
- The selector and aggregator for stream S: one process each, found through the log partition. They handle 5,000 comments and 400 reaction flushes per second, which one process can do. Because the selector caps its output at 10 comments per second, nothing downstream grows with the comment rate.
Step 5: the edges of the design
- Joining: subscribe first, then fetch the last 50, and drop duplicates by comment ID. Fetch first and subscribe second, and any comment that arrives in between is missed.
- Your own comment: it may never be selected for broadcast, so the client shows it locally as soon as the API accepts it.
- Moderation: cheap checks (a blocklist, a per-user rate limit such as one comment every 2 s) run before a comment is accepted. A slower classifier runs afterwards and can retract a broadcast comment with a "remove comment 9f1" event.
- Regions: viewers connect to gateways in their nearest region. Each stream's selector lives in one home region and sends every batch once to each region's dispatcher. LinkedIn described the same step: subscriptions stay local to a data center, and each published event is forwarded to the dispatchers in all the other data centers.
- "2.1M watching": each gateway reports its connection count per stream every 5 s, and the sum is the display. A viewer with two tabs counts twice, which is fine for a number shown to one decimal place.
Trace one comment
A viewer in SΓ£o Paulo comments on a stream whose home region is in Virginia. The timings are assumptions for the budget, not measurements.
| Step | Assumed time | Running total |
|---|---|---|
| POST to the nearest region, forwarded to the home region | 60 ms | 60 ms |
| Accepted: written to the comment store, appended to the log | 20 ms | 80 ms |
| Wait for the selector's 500 ms window to close (worst case) | 500 ms | 580 ms |
| Batch sent to the dispatchers in every region (worst: another continent) | 100 ms | 680 ms |
| Dispatcher to the 30 gateways | 10 ms | 690 ms |
| Gateway to viewers' phones | 50 ms | 740 ms |
The worst case is about 0.75 s, well inside 2 s. The budget also tells you which knob to turn if the requirement tightens to 500 ms: the batching window, not the servers.
Failure modes and what you give up
| Event | What happens | Design response |
|---|---|---|
| A gateway dies | About 67,000 viewers reconnect to the other gateways | Jittered backoff; on reconnect, subscribe and then fetch the last 50 |
| A dispatcher dies | Its in-flight batches are lost | The gateway table lives in a key-value store, so another dispatcher takes over; a few batches go missing |
| The selector for S dies | Comments stop appearing on S | The log hands the partition to another consumer, which resumes from the log; every comment is still stored |
| The home region fails | Comments on S stop | Fail the selector over to another region; video playback is a separate system and keeps going |
| A spam wave: 50,000 comments per second | Store and selector load jump | Per-user rate limits before acceptance; the selector's cap keeps everything downstream flat |
What you give up, said plainly: most comments on a huge stream are never broadcast; a reconnecting viewer can miss a few seconds of comments; the counts are approximate; and each stream's comments depend on one home region.
Likely follow-ups: "The streamer wants to see every comment." Give the streamer's dashboard its own subscription to the full stream: one reader of 5,000 comments per second is easy, because it is one reader, not two million. "How is this different from group chat?" Chat members must receive every message, so chat writes a durable inbox per member and delivers at least once; live comments are sampled and delivered at most once.
Final round: no label on the prompt
Real prompts don't say which move they want. For each one, name the pressures, then the moves, then what you give up.
Challenge 1: the new-episode rush
A podcast app releases a new episode of its most popular show at 06:00, and 3 million people download it in the first hour. The episode is 60 MB.
Hint
Work out the bandwidth first, then ask what is hot.
Check your answer
3M Γ 60 MB in 3,600 s = 50 GB/s, or 400 Gbps averaged over the hour, all for one object. The pressures are media and a hot key.
- Serve the file from a CDN under an immutable, versioned URL, uploaded and processed well before 06:00.
- At 06:00 the first requests miss at every edge at once. Request collapsing at the edges plus a shield tier mean the origin sees about one fetch per shield, not three million. You can also warm the edges by requesting the file through each region shortly before release. Netflix has described taking pre-positioning much further with Open Connect: its own cache appliances inside ISP networks, filled with content in nightly updates before viewers ask for it.
- The "new episode" notification to subscribers is fan-out. Send it in waves, which also spreads the download spike.
What you give up: warming and a shield cost money for one hour of traffic, and staggered notifications mean some subscribers hear about the episode minutes later.
Challenge 2: goal alerts
A sports app must alert the 30 million followers of a football club within 30 seconds of a goal.
Check your answer
30M in 30 s = 1 million notifications per second. The pressures are fan-out to a huge audience and real-time delivery to apps that are mostly closed.
- Precompute recipient lists per club, split into many chunks (say 10,000 users each), before the match starts.
- On a goal, publish one event. Users who have the app open get the alert over their live connection within seconds; that part you control.
- Everyone else needs mobile push, and the push providers set the pace. FCM's default quota is 600,000 messages per minute per project, which is 10,000 per second, a hundredth of what you need. A standard increase adds at most 25%, and larger temporary quotas for planned events need weeks of notice. At 10,000 per second, 30M alerts take 50 minutes, and only 300,000 people (1%) hear within 30 s.
- FCM topics are the cheap route, not the fast one: their fan-out runs at around 10,000 per second, is not adjustable, and is optimised for throughput rather than latency. Use a topic only if "everyone within the hour" is acceptable. To control the order, send to device tokens yourself, users who opted in to every goal first.
- Give each alert an expiry (FCM
ttl, APNsapns-expiration) and a collapse key, so a stale "Goal!" isn't shown after the next goal. Retry until it expires, then drop it.
What you give up: this is a requirement to push back on. "30 million within 30 seconds" becomes "every open app within seconds, the most eager fans first, everyone else over tens of minutes".
Challenge 3: "37 people are looking at this hotel"
A travel site shows, on each hotel page, how many people are viewing it right now. The number must update within 10 seconds and doesn't need to be exact.
Check your answer
A counter of live presence. Each open page sends a heartbeat every 5 s (or holds an SSE connection). Count viewers per hotel over a sliding window: a per-hotel sorted set of sessions scored by their last heartbeat, trimmed to the last 10 s, works for most hotels. For the few very popular pages, hold an SSE connection per viewer and let each gateway report its per-hotel connection count every few seconds, as in the live-stream design, which avoids a hot key. Because the number is approximate, a HyperLogLog per 10-second bucket also works: it counts distinct sessions in constant memory.
What you give up: exactness (a user with two tabs may count twice) and up to 10 s of lag, both allowed by the requirement.
Cheat sheet: pressure β move β cost
π§ One idea behind every move: replace a rate that popularity sets with a rate you choose. Inserts per post (hybrid threshold), cache reads per key (servers Γ· TTL), origin fetches per object (collapsing), messages per viewer (batching), writes per counter (flush interval).
| Pressure | Move | What you pay |
|---|---|---|
| One write, many readers | Push into precomputed timelines; pull big accounts at read time (hybrid) | Write amplification; a merge path; a threshold to tune |
| Few authors, many readers each | Pull, with each source's recent posts cached | Slower loads for readers who follow many sources |
| One key takes most of the reads | Local cache with a short TTL; copies of the key | Staleness up to the TTL; N writes per update |
| Many requests miss the same key at once | Coalescing (single-flight), leases, early refresh | Waiters share one fetch's latency and failure |
| Many keys expire at the same moment | Random jitter on TTLs | Slightly less predictable expiry |
| The cache tier could fail | Replicas across zones, load shedding, circuit breaker, gradual warm-up | A degraded product during recovery |
| Items are MB to GB | Presigned direct upload, async processing, immutable URLs on a CDN | A multi-step upload flow and "processing" states |
| Storage in petabytes | Originals in cold tiers; playable renditions in millisecond tiers | Restore delays for anything cold |
| The server must speak first | Gateway tier, routed through a registry or pub/sub | A stateful tier; reconnect storms to manage |
| Readers must not miss items | Durable inbox with sequence numbers; the push is a hint | Storage per message; at-least-once plus dedupe |
| Many writers, one number | Records as truth; aggregated deltas; idempotent flush | Counts lag by seconds; reconciliation jobs |
| Distinct counts, approximate is fine | HyperLogLog | About 0.81% error; no per-user answer |
Before moving on, pick one move and explain it aloud as you would to an interviewer, starting from a blank whiteboard: the number that triggers it, the move itself, one failure it must survive, and what it costs. Then build the habit that makes this stick: next time you open a social or streaming app, spend a minute guessing which of the five pressures it faces and where it pays for each one.
Next: the full designs, in any order. Twitter/Instagram News Feed takes Move 1 all the way (fan-out, ranking, pagination). YouTube / Netflix Streaming takes Move 3 to video (transcoding, adaptive bitrate, CDN delivery). Chat Application (WhatsApp) takes Move 4 to guaranteed delivery (ordering, receipts, presence).