Short answers to what readers ask most about this topic.
01What is the difference between fan-out on write and fan-out on read?
Fan-out on write copies a new post into every follower's precomputed feed when it is created, so opening a feed is a single lookup. Fan-out on read stores the post once under its author and gathers the posts of every followed account when the feed is opened. Write pays per follower and read pays per account followed, so the cheaper one depends on how often feeds are opened versus how often people post.
02What is the celebrity problem in news feed design?
It is the case where one account has so many followers that a single post triggers a huge burst of fan-out writes. With 500,000 followers, one post is 500,000 inserts, which in the worked example equals 48 minutes of average write load. The usual fix is a hybrid: skip the push above a follower threshold and merge that account's posts into feeds at read time.
03How do you store a news feed in Redis?
Use one sorted set per active user, with the post id as the member and the same unique id as the score. Add with ZADD, trim the oldest entries with ZREMRANGEBYRANK to keep a fixed cap such as 200, and read newest-first with ZRANGE using BYSCORE and REV. Keep the post bodies in Postgres and treat the Redis feed as a rebuildable cache.
04Why use cursor pagination instead of page numbers for a feed?
A feed changes while someone reads it, so a new post shifts every item down and page 2 can repeat an item from page 1. A cursor, here the id of the last item returned, asks for entries strictly older than that point and stays stable. It also avoids the cost Redis documents for large LIMIT offsets, which must walk past that many elements.
05Should a small app use fan-out on write?
Usually not at first. With a few thousand users, one indexed SQL query that joins follows to posts is correct and easy to measure. Move to fan-out on write only when you can show that reads per post are high and that the query is the bottleneck, and keep Postgres as the source of truth when you do.
News Feed System Design: Fan-Out on Write vs Fan-Out on Read
How to design a news feed: fan-out on write versus fan-out on read, the celebrity problem, Redis sorted set feeds, cursor pagination and ranking, with the arithmetic shown.
To design a news feed, use fan-out on write for ordinary accounts: when someone posts, push the post id into each follower's Redis sorted set, so a feed read is one lookup. Skip the push for accounts with huge followings and merge their posts at read time. Page with a cursor and rank only the newest few hundred candidates.
A feed looks like one query: show me the latest posts from the people I follow. On a small table that is exactly what it is, a join and an ORDER BY, and it is fine. Then you write down how often feeds are opened compared with how often anyone posts, and that join becomes the most expensive thing in the system.
This is a worked design, not a war story. I have not run a social network, so every number below comes from stated assumptions with the arithmetic shown, and the Redis behaviour comes from its documentation. The stack is the one I would reach for on a single VPS: Postgres as the source of truth, Redis for the feeds and a NestJS service in front.
What do fan-out on write and fan-out on read mean?
Fan-out is one event spreading to many destinations. Wikipedia describes it in messaging as delivering a message to one or multiple destinations, possibly in parallel. In a feed, the event is a new post and the destinations are the timelines of everyone who follows the author. The only real decision is when you do that spreading.
Fan-out on write (push): when a post is created, copy its id into every follower's precomputed feed. Opening a feed is one lookup.
Fan-out on read (pull): store the post once, under its author. When a user opens the feed, fetch the newest posts of every account they follow and merge them.
Hybrid: push for ordinary authors, pull for authors with a very large following, and merge both when the feed is read.
Every feed design is a choice about where to pay: at write time, once per follower, or at read time, once per account followed. Which is cheaper depends on how those two counts compare against how often each event happens, so the next section puts numbers on it.
How much work does each model actually do?
Write the assumptions down first and derive everything from them. These are inputs for a design exercise on a single VPS, not measurements from a real product, and you should replace them with your own logs.
Assumed inputs (a design exercise, not a measurement):
registered users = 1,000,000
daily active users (DAU) = 200,000
feed opens per DAU per day = 10
posts per DAU per day = 0.5
average followers = 150 (so the average account follows 150 too)
peak vs average = 3x
seconds per day = 86,400
Reads and posts
feed opens per day = 200,000 x 10 = 2,000,000
average opens / s = 2,000,000 / 86,400 = 23.1 peak = 69.4
posts per day = 200,000 x 0.5 = 100,000
average posts / s = 100,000 / 86,400 = 1.16 peak = 3.5
reads per post = 2,000,000 / 100,000 = 20
Fan-out on write (push)
feed inserts per day = 100,000 x 150 = 15,000,000
average inserts / s = 15,000,000 / 86,400 = 173.6 peak = 520.8
per feed open = 1 range query (23.1 / s average)
Fan-out on read (pull)
inserts per day = 100,000 (1 per post)
outbox lookups per day = 2,000,000 x 150 = 300,000,000
average lookups / s = 300,000,000 / 86,400 = 3,472.2 peak = 10,416.7
Pull does 300,000,000 / 15,000,000 = 20x the lookups that push does inserts.
That 20 is just reads per post.
Pull does 300 million outbox lookups a day against push's 15 million feed inserts, a factor of 20, and that factor is simply reads per post: 2,000,000 feed opens divided by 100,000 posts. Whenever a feed is opened more often than anyone posts, which is nearly always, push is the cheaper side to pay. Its price is that the work happens whether or not the follower ever opens the app.
Model
Work per post
Work per feed open
Main weakness
Pick it when
Fan-out on write (push)
150 inserts on average
1 range query
A huge account turns one post into a burst of writes; inactive followers still get the work
Feeds are opened far more often than people post and follower counts are bounded
Fan-out on read (pull)
1 insert
150 outbox lookups plus a merge
Read latency grows with the number of accounts followed
The graph is small, or most users rarely open the feed
Hybrid
150 inserts for ordinary authors, 1 for large ones
1 range query plus about 2 outbox lookups
Two code paths and a threshold to tune
A few accounts hold a large share of all follow edges
The table shows why neither model is free. Push moves the cost to writes and storage. Pull keeps storage tiny but makes every read slower and less predictable, because its latency scales with how many accounts the reader follows.
What is the celebrity problem and how does a hybrid fix it?
Averages hide the account that breaks push. Assume one author has 500,000 followers, half the user base. A single post from that account is 500,000 inserts, which is a very different event from the average post at 150.
One author with 500,000 followers (half the user base) posts once:
push inserts = 500,000
average push load = 173.6 inserts / s
500,000 / 173.6 = 2,880 s = 48 minutes of AVERAGE load, from one post
if workers sustain 20,000 inserts / s (assumed, load-test it):
500,000 / 20,000 = 25 s until the last follower's feed has the post
the same account posting 3 times in an hour = 1,500,000 inserts
Hybrid with a 10,000-follower threshold. Assume 20 such accounts average 100,000
followers and post 2 times a day (these posts are part of the 15,000,000 above):
inserts avoided = 20 x 2 x 100,000 = 4,000,000 per day (26.7% of 15,000,000)
push inserts remaining = 15,000,000 - 4,000,000 = 11,000,000 per day
average inserts / s = 11,000,000 / 86,400 = 127.3
Read side, assuming a reader follows 2 of those large accounts:
calls per feed open = 1 feed ZRANGE + 2 outbox ZRANGEs = 3 (pull would be 150)
calls per day = 2,000,000 x 3 = 6,000,000 = 69.4 / s average
The hybrid rule is a threshold on follower count. Below it, push as usual. At or above it, do not fan out at all: write the post only to the author's own outbox and merge those outboxes into the reader's feed when it is read. Readers follow only a handful of such accounts, so the merge costs a few extra lookups instead of hundreds.
const PUSH_FOLLOWER_LIMIT = 10_000; // a starting guess; tune it from fan-out queue lag
const OUTBOX_CAP = 200;
export async function onPostCreated(post: { id: number; authorId: number }) {
// 1. Always write the author's outbox: it is the pull path AND the rebuild source.
const outbox = "outbox:" + post.authorId;
await redis.zadd(outbox, post.id, String(post.id));
await redis.zremrangebyrank(outbox, 0, -(OUTBOX_CAP + 1));
// 2. Large accounts are pulled at read time, so there is nothing to fan out.
const followers = await followerCount(post.authorId); // users.follower_count
if (followers >= PUSH_FOLLOWER_LIMIT) return;
// 3. Everyone else is pushed by a worker, never inside the request.
await fanOutQueue.enqueue({ postId: post.id, authorId: post.authorId });
}
The threshold is a tuning knob, not a constant of nature. I would start at 10,000 followers, a number I picked rather than measured, then watch fan-out queue lag and move it. Moving it is safe because every feed can be rebuilt from Postgres.
Put only post ids in the feed, never bodies. A feed entry then costs the same however long the post is, an edit needs no fan-out, and a deleted post is dropped when you hydrate ids from Postgres.
How do you store a feed in Redis sorted sets?
Postgres stays the source of truth for posts and follows. Redis holds one sorted set per active user, with the post id as the member and the same id as the score, so newest-first reading is a range query. The Redis docs give ZRANGE a cost of O(log(N)+M) for N elements and M returned, and ZREMRANGEBYRANK the same shape, which is why trimming after every insert is cheap.
-- Source of truth. Redis feeds are rebuildable from these two tables.
CREATE TABLE posts (
id bigint GENERATED ALWAYS AS IDENTITY PRIMARY KEY, -- also the Redis score
author_id bigint NOT NULL,
body text NOT NULL,
created_at timestamptz NOT NULL DEFAULT now(),
deleted_at timestamptz -- soft delete
);
CREATE INDEX posts_author_idx ON posts (author_id, id DESC); -- the outbox / pull query
CREATE TABLE follows (
follower_id bigint NOT NULL,
followee_id bigint NOT NULL,
created_at timestamptz NOT NULL DEFAULT now(),
PRIMARY KEY (followee_id, follower_id) -- fan-out: who follows X?
);
CREATE INDEX follows_follower_idx ON follows (follower_id); -- pull: whom do I follow?
-- users.follower_count is a maintained counter, so the hybrid check is one primary-key read.
-- The pure pull query. Correct, easy to measure, and fine until it is not:
SELECT p.id
FROM posts p
WHERE p.author_id IN (SELECT followee_id FROM follows WHERE follower_id = $1)
AND p.deleted_at IS NULL
ORDER BY p.id DESC
LIMIT 20;
A fan-out worker does the writing, outside the request that created the post: insert the post, enqueue a job, return. The worker walks the follower list in keyset pages and writes to Redis in pipelines.
const FEED_CAP = 200; // newest 200 entries per feed
const PAGE_SIZE = 1000;
export async function fanOut(postId: number, authorId: number) {
let afterId = 0;
for (;;) {
// Keyset page over the primary key (followee_id, follower_id). Never OFFSET.
const { rows } = await pg.query(
"SELECT follower_id FROM follows " +
"WHERE followee_id = $1 AND follower_id > $2 " +
"ORDER BY follower_id LIMIT $3",
[authorId, afterId, PAGE_SIZE],
);
if (rows.length === 0) return;
const pipe = redis.pipeline();
for (const { follower_id } of rows) {
const key = "feed:" + follower_id;
// Same member + same score twice is a no-op update, so a retried job is safe.
pipe.zadd(key, postId, String(postId));
// Ranks are 0-based from the LOWEST score; -1 is the newest. Removing 0 .. -201
// leaves exactly the newest 200.
pipe.zremrangebyrank(key, 0, -(FEED_CAP + 1));
}
await pipe.exec();
afterId = rows[rows.length - 1].follower_id;
// Production: skip followers inactive for 30 days; rebuild their feed when they return.
}
}
Trimming bounds memory. 200 entries per active feed times 200,000 active feeds is 40,000,000 entries. At an assumed 50 bytes each that is 2.0 GB, and that 50 is a placeholder: replace it with what MEMORY USAGE reports on your own data before trusting the total. A user who scrolls past 200 entries falls through to a Postgres query.
A Redis feed is a cache, not a record. If it is flushed, fan-out lagged, or a follower was skipped, the answer is to rebuild from Postgres, never to treat Redis as the only copy. Design the rebuild path on day one, because you will use it.
How do you paginate a feed with a cursor?
A feed changes while people read it, so page numbers drift: a new post pushes everything down and page 2 repeats an item from page 1. I covered the general case in the cursor versus offset pagination post. Here the cursor is simply the id of the last item returned.
# Score and member are the same unique id.
ZADD feed:42 1001 1001 1002 1002 1003 1003 1004 1004 1005 1005
# Page 1: newest first, 3 entries, no cursor yet.
ZRANGE feed:42 +inf -inf BYSCORE REV LIMIT 0 3
# 1) "1005" 2) "1004" 3) "1003" -> cursor = 1003
# Page 2: strictly below the cursor. A new post (1006) would not shift this page.
ZRANGE feed:42 (1003 -inf BYSCORE REV LIMIT 0 3
# 1) "1002" 2) "1001"
# Wrong: ZRANGE feed:42 +inf -inf BYSCORE REV LIMIT 3 3
# An offset drifts when 1006 arrives, and the server must walk past 3 entries first.
With BYSCORE and REV, ZRANGE takes the highest score first, and a leading parenthesis makes a bound exclusive, so the next page starts strictly below the last id. This only works when no two entries share a score, because Redis orders equal scores lexicographically and a tie at a page boundary would be skipped. A unique id as the score avoids that. The docs also warn that a large LIMIT offset forces Redis to walk past that many elements, which is the other reason to use a cursor and not an offset. Doubles hold integers exactly up to 2 to the power 53, about 9.0e15, far above any realistic row count.
const PAGE = 20;
export async function readFeed(userId: number, cursor: string | null) {
// "(" makes the bound exclusive: strictly older than the last id already shown.
const max = cursor ? "(" + cursor : "+inf";
const largeAccounts = await largeAccountsFollowed(userId); // usually 0 to a handful
const keys = ["feed:" + userId, ...largeAccounts.map((a) => "outbox:" + a)];
// With BYSCORE + REV the FIRST bound is the highest score, the second the lowest.
const lists = await Promise.all(
keys.map((key) =>
redis.zrange(key, max, "-inf", "BYSCORE", "REV", "LIMIT", 0, PAGE + 1),
),
);
// Merge: ids are numbers below 2^53, so a numeric sort is exact. Fetch one extra to know
// whether another page exists.
const ids = [...new Set(lists.flat())]
.sort((a, b) => Number(b) - Number(a))
.slice(0, PAGE + 1);
const hasMore = ids.length > PAGE;
const pageIds = ids.slice(0, PAGE);
// Hydrate from Postgres. Deleted posts simply do not come back, so a page may be short.
const { rows } = await pg.query(
"SELECT id, author_id, body, created_at FROM posts " +
"WHERE id = ANY($1::bigint[]) AND deleted_at IS NULL",
[pageIds],
);
const byId = new Map(rows.map((r) => [String(r.id), r]));
const items = pageIds.map((id) => byId.get(id)).filter(Boolean);
return { items, nextCursor: hasMore ? pageIds[pageIds.length - 1] : null };
}
Where does ranking fit in?
Ranking is a second stage on top of the feed, not a different storage model. Get candidates cheaply first: the newest ids from the feed plus the large-account outboxes. Then score and sort only those. Twitter's open-sourced recommendation repository describes the same shape, with candidate sources, a light ranker, a heavy neural ranker, visibility filters and a mixer, and says about half of the posts in the For You timeline come from the in-network source.
For a first version, a transparent formula beats a model you cannot debug. This one is my own illustration, not a recommendation from any source, so tune the weights against what your users actually do.
Illustrative score, my own toy formula, not taken from any source:
rank = (1 + 2 x likes + 3 x comments) / (age_hours + 2) ^ 1.5
Post A: 10 likes, 2 comments, 1 hour old
(1 + 20 + 6) / 3 ^ 1.5 = 27 / 5.196 = 5.20
Post B: 100 likes, 10 comments, 12 hours old
(1 + 200 + 30) / 14 ^ 1.5 = 231 / 52.38 = 4.41
Post C: 0 likes, 0 comments, just posted (0 hours)
1 / 2 ^ 1.5 = 1 / 2.828 = 0.35
Chronological order: C, A, B. Ranked order: A, B, C.
Cost of ranking 200 candidates per feed open:
2,000,000 opens x 200 = 400,000,000 score evaluations per day = 4,629.6 / s average
Ranking breaks the pure score cursor, because the returned order is no longer the stored order. Either snapshot the ranked id list for a few minutes and page by position inside the snapshot, or keep chronological paging and rank within each page. Decide which you ship and say so in the product, because users notice when a feed reshuffles under them.
Which model should you choose, and what breaks first?
Run this checklist before building anything.
Compute reads per post from your own logs. If feeds are opened far more often than anyone posts, pay at write time.
Plot the follower distribution. If a few accounts hold a large share of all follow edges, add a threshold and the hybrid path.
If you have a few thousand users, ship the single indexed SQL query first. It is correct, and you can measure before you add Redis.
Keep Postgres as the source of truth and make the feed rebuild a tested job, not a runbook paragraph.
Make every fan-out job idempotent. Adding the same member with the same score again changes nothing, so retries are safe.
What breaks first is rarely the arithmetic. Deleted posts and unfollows leave ghost ids in feeds, so hydrate from Postgres and drop what no longer exists. A flushed Redis serves empty feeds until you rebuild lazily on the first read from each followed account's recent posts. And fan-out lag during a burst makes new posts appear late, so alert on queue depth, not on CPU.
Pay for fan-out where the traffic is cheap to predict. Reads outnumber posts, so push ordinary authors into a capped Redis sorted set per user, treat large accounts as pull, page with a cursor, and rank only a few hundred candidates. Keep Postgres as the truth and every other layer disposable.