System design practice

News Feed

medium fan-out queues caching read-heavy

Fan out on write through a queue into feed caches.

Solve it in your browser Read the lesson first

Design the home timeline of a social network: people publish short posts, and everyone who follows them sees those posts at the top of their feed, newest first.

Functional requirements

Use these use case names exactly: the traffic, requirements and tests in problem.proschi refer to them.

Scale

Constraints

What is given

problem.proschi declares the user and the existing graph service that knows who follows whom. Listing the followers of an account pages through up to 5,000 ids and takes about 50 ms. The file also holds the traffic, requirements and tests. Add the components, the connections and the two use cases.

How your design is checked

You write the design as text in Proschi. Tests run in your browser: a simulation of the traffic above checks latency, availability, cost and what happens when a machine fails. How the simulation works.

Start designing

Lesson · 12 min read

Learn it: News Feed #

News Feed: build the feed before anyone asks for it #

A home timeline looks like a query: "the 50 newest posts by people I follow, sorted by time". Run that query every time someone opens the app, for millions of users, and you join a social graph with a posts table ten thousand times a second. This lesson turns the query inside out. Instead of building a feed when it is read, you deliver each post into its readers' feeds when it is written: through a queue, into a cache.

What you'll learn #

The problem, explained #

Who uses it. 20 million daily active users of a social network. They publish short posts and open the app to see what the people they follow posted.

Functional requirements.

Non-functional requirements. p99 (the latency 99% of requests beat) under 50 ms for reading and under 90 ms for publishing. Reading available 99.9%. A post is never lost after 201, and no feed shows a post that was not stored. Publishing never waits for the social graph or the feeds. Reading never touches a database or the social graph. Any single machine can fail. At most $4,000 a month, including the social graph.

What is given, and why. given.proschi declares the user and an existing graph service (four replicas) that knows who follows whom. A capacity line sets its latency to 50 ms per call, because listing up to 5,000 follower ids is slow. That number is the reason the design looks the way it does.

What the tests check.

Back-of-the-envelope #

QuantityArithmeticResult
Read feedgiven10k rps
Publish postgiven500 rps
Read : write requests10k ÷ 50020 : 1
Feed inserts from fan-out500 posts/s × 200 followers100k writes/s
Feed cache operations100k inserts + 10k reads110k ops/s
Feed API requests10k reads + 500 publishes10.5k rps
Graph callsone per post500 rps
Feed size per user500 ids × 8 bytesabout 4 KB raw
All feeds (20M users)20M × 4 KBabout 80 GB raw, several times that with Redis overhead
Posts per day, upper bound500/s × 86,400 s × ~1 KBabout 43 GB/day

The ratio flips. Requests are 20:1 reads to writes, but on the feed cache the fan-out makes it 10:1 writes to reads. If you size the cache as "read-heavy", you will be surprised. In Proschi, writing the feed step as x200 ZADD … tells the model to count 200 cache writes per post; its latency is counted once, as if the writes were batched or pipelined.

Sizing the feed cache. A Redis replica in the model takes 100k operations a second, and in the model caches serve writes on every replica. So 110k ops/s on two replicas is 55% busy. That looks fine until survive any node failure re-runs the analysis with one replica fewer: 110k on one replica is 110%, saturated. Count replicas for the degraded case, not just the healthy one.

Sizing the Feed API. 10.5k requests a second at 2k per service replica is 5.25 replicas at 100%. Divide by a 70% target and check it with one replica lost.

Latency. Read feed is load balancer → API → feed cache → post cache: four short hops, two of them to a 1 ms cache. Publish post is load balancer → API → posts database → queue, then 201. Everything after the 201 (the graph's 50 ms, the 200 writes) is asynchronous and does not count toward publish latency. If you put the graph call before the 201, that 50 ms hop and its tail alone threaten the 90 ms limit.

Database choice. 500 inserts a second fits on one PostgreSQL primary (5k writes/s), so this is not forced. A partitioned store like Cassandra fits append-only posts keyed by author well, and it grows without a resharding project. Either way, the database is off the read path.

Cost. Graph $400, service replicas $100, a Redis replica $150, a Kafka broker $200, a Cassandra replica $500, a load balancer $50. The budget fits a careful design, but not an extra replica on every tier.

Concepts #

Fan-out on write versus fan-out on read #

There are two ways to produce a timeline:

Pick by the ratio of reads to writes, weighted by the cost of each. Here reads are 20 times more frequent, the read latency budget is tight (50 ms), and follower counts are bounded at 5,000, so push wins clearly. Twitter's timeline service has long been described this way: a fan-out step inserts tweet ids into Redis-backed timelines of each follower.

When not to push: when some authors have millions of followers. One post becomes millions of writes and takes minutes to land. The fix is a hybrid: push for ordinary accounts, pull for the few huge ones, and merge the two at read time. The "Feeding Frenzy" paper makes this choice for each producer/consumer pair, based on how often one posts and the other reads.

title "Push into inboxes"

author  "Author"       [Actor]
api     "Inbox API"    [REST API] x2
log     "Updates"      [Kafka]    x2
pusher  "Pusher"       [Worker]   x2
members "Member Index" [REST API] x2
inboxes "Inboxes"      [Redis]    x2

author -> api     : HTTPS
api    -> log     : produce
log    -> pusher  : consume
pusher -> members : list
pusher -> inboxes : write

usecase "Announce" {
  author   -> api     : POST /announcements
  api      -> log     : PRODUCE Announced a_1
  log     --> api     : ack
  api     --> author  : 202
  log     ->> pusher  : Announced a_1
  pusher   -> members : GET /groups/9/members
  members --> pusher  : 200 [~50 ids]
  pusher   -> inboxes : x50 LPUSH inbox:{member} a_1
}

Store first, then fan out (through a queue) #

Two ordering rules make push safe.

Store first. Write the post to durable storage before anything else references it. A feed must never point at a post that does not exist. If the fan-out fails halfway, it can be retried from the stored post; if the store failed, the user saw an error and nothing leaked into feeds.

Fan out behind a queue. After the post is stored, put a PostCreated event on a durable queue and answer 201. A fan-out worker consumes the event, calls the graph, and writes the feeds. Now the author's latency does not depend on the follower count or the graph's speed. You also get retries for free: if the worker dies, the event is consumed again. Because that can happen, feed inserts should be idempotent (safe to repeat). A sorted-set insert keyed by post id is, because adding the same member twice changes nothing.

The trade-off is eventual consistency: a follower may open the app a second before the post lands in their feed. The requirement ("within a few seconds") allows it. The author's own feed is a different story: clients usually show the author's own post immediately.

Caching ids and objects separately #

Store feeds as lists of post ids (a Redis sorted set per user, scored by time, trimmed to the newest 500), and store the posts themselves in a separate cache keyed by id. A read is then two cache calls: fetch 50 ids, then fetch those 50 posts in one batched get.

Why split them? A post is written once into the post cache but referenced by 200 feeds. Copying the full text into every feed would multiply memory by 200, and an edit or deletion would need 200 updates. Ids are tiny and immutable.

title "Ids, then objects"

reader "Reader"      [Actor]
api    "Board API"   [REST API] x2
lists  "Board Lists" [Redis]    x2
items  "Card Cache"  [Redis]    x2

reader -> api   : HTTPS
api    -> lists : ZREVRANGE
api    -> items : MGET

usecase "Open board" {
  reader -> api    : GET /boards/3
  api    -> lists  : ZREVRANGE board:3 0 19
  lists --> api    : 20 card ids
  api    -> items  : MGET card:… (20 keys)
  items --> api    : 20 cards
  api   --> reader : 200
}

Trimming matters too: a feed only needs its newest 500 ids, so memory per user is bounded no matter how long they have been following people. Inactive users' feeds can expire and be rebuilt on their next visit, a common optimisation that this problem leaves out.

Designing it step by step #

1. Scope. Ask: how many users and how active (20M DAU), what the feed contains (posts by followees, newest first, no ranking), follower distribution (average 200, max 5,000: no celebrities), freshness (a few seconds), and the latency targets. Explicitly park ranking, media and celebrities as follow-ups.

2. High-level design. Sketch two flows.

Say why not fan-out on read: every read would call a 50 ms graph and merge posts from up to 5,000 accounts out of the database. The test Reading the feed never queries a database rules it out, and so would the 50 ms p99.

3. Deep dive.

4. Wrap-up. Walk the requirements: reads touch only caches; publish answers after the store and the queue; the fan-out retries from the queue; nothing is single. Then extend: the hybrid for celebrities, ranking (fetch more candidates than you show, score them), feed rebuild for users whose feed expired, and deletion (remove the post from the post cache; feeds that still list its id skip it on read).

Common mistakes #

Fanning out inside the request (wrong/fan-out-in-request). The API stores the post, queues the event, then also calls the graph and writes the 200 feeds before answering. The queue is decoration: the author waits for the 50 ms graph and the fan-out, and an author with 5,000 followers would wait much longer. In the model the publish p99 climbs past 200 ms. It fails Fan-out runs behind a queue, and the p99 limit for Publish post as well.

Two feed cache nodes (wrong/two-feed-cache-nodes). When both are up, two Redis replicas handle 110k ops/s at about 55%. Lose one and the survivor needs 110% of its capacity. In production, that is the moment feeds stop updating. Or worse, the cache falls over and takes reads with it. It fails survive any node failure.

Other classic mistakes.

In the interview #

Start from the ratio: "Reads outnumber writes 20 to 1 and the read budget is 50 ms, so we should do the work at write time." Then draw the two flows and mark the async boundary clearly. Do the 500 × 200 = 100k multiplication on the whiteboard; it is the number the rest of the design depends on.

Likely follow-ups:

Further reading #

Now design it

More system design problems