Design a news feed

Fan-out on write vs on read - the timeline trade-off that decides whether a celebrity post melts your database.

9 min read

A feed is not a list of posts. It is a list of posts assembled for each person, which is a very different problem. The work is not storing a post; it is deciding where copies of it need to end up, and when.

Just select the posts from everyone you follow

The first design anyone writes is one query: find who I follow, fetch their recent posts, sort by time, take fifty. It is correct and readable, and it falls over long before you have a real product.

Each square is a shard queried to build one feed page; a red one is having a slow moment
Shards per feed open
200
Shard queries a second
2M
Feed opens hitting a slow shard
87%
200
10K
Follow 200 accounts and every feed open asks 200 shards. If each is slow 1% of the time, 87% of feed opens are slow.

One query is cheap. But the feed is the busiest screen in the product, so this runs constantly, and because posts are stored by author, the 200 accounts you follow live on 200 different shards. Every feed open becomes a 200-way scatter-gather, and the page is as slow as the slowest shard. Nothing about the query is wrong; the timing is. The work happens at read time, and reads outnumber writes by orders of magnitude.

Do the work when they post, or when they look

The fix is to move the work to the rarer event. Fan-out on write pushes each new post into every follower's feed cache as it is published, so opening the app is a single cache read. Fan-out on read writes nothing extra and merges at open time, as before. Same feed, opposite costs.

One post is copied into every follower's feed cache right now
+2M more caches
Opening the app later reads one prebuilt list
Work per post
2M cache writes
Work per feed open
1 cache read
2M
A two-million-follower account under push: one post, two million cache writes. Switch to pull and every reader pays a 200-way merge instead.

Push makes reads instant, which is perfect for the vast majority of accounts. But pushing means literally writing a reference into every follower's cache, so a post from someone with two million followers is two million writes. Doing it in the background changes when you pay, not how much. It is the hot-key problem in its purest form, and it is why no large system is purely push.

Push for almost everyone, pull for the handful who would break it

Follower counts follow a power law: a long tail of small accounts and a tiny head of enormous ones. Fan-out cost is followers times posts, so that tiny head generates most of the total.

Casual700K accounts, 150 followers
push, 2%
Active250K accounts, 1.2K followers
push, 5%
Popular45K accounts, 25K followers
push, 19%
Influencer4.8K accounts, 400K followers
pull, 33%
Celebrity200 accounts, 12M followers
pull, 41%
Bars: each group's share of all fan-out writes
Moving 5K of a million accounts to pull removes 74% of all fan-out writes.
316K
Five thousand accounts, half a percent of them, generate most of the fan-out. Put them on pull and almost everyone else keeps instant feeds.

That asymmetry is why the hybrid is not a grubby compromise. Ordinary accounts fan out on write and their followers read one cache key. Celebrity posts are not fanned out at all; they are merged in at read time from the short list of big accounts each person follows. Two code paths and a threshold, and every large feed system converges on some version of it.

What you actually tune is the cutoff. It moves with your follower distribution, your cache budget and how much read latency you will spend to protect write throughput. Getting it wrong in either direction is recoverable. Having only one path is not.

The short version

  • Building feeds at read time means a scatter-gather across a shard per followed account.
  • Push copies each post into followers' caches so reads are one lookup.
  • Push breaks on celebrities: one post, millions of writes.
  • Push for most accounts, pull for the few huge ones, and tune the cutoff.