Fan-out on write.
A feed read that has to ask 380 accounts "what's new?" is slow, and it repeats the same work every time someone scrolls. Fan-out on write moves that work to posting time: a queue hands each new post to workers that append its ID to the feed of every active follower. Posting stays fast, feeds are ready before anyone asks, and a handful of huge accounts are left out of the push on purpose.
Builds on Delivery semantics.
Requirements.
The write path of Chorus's home feed: from tapping Post to the post sitting in followers' feeds. Reading, merging and ranking those feeds are the next two topics.
Functional requirements
Non-functional requirements
Capacity estimates
What the push costs
- New posts per day
- 50MAssumption for Chorus (400M monthly and 200M daily users).
- Mean followers of a pushed post
- 250Median 90. Excludes pull-mode accounts over 100K followers.
- Followers active in the last 30 days
- 60%Assumption. Active means opened the feed.
- Evening peak over the daily mean
- ×3Assumption.
- Logical bytes per feed entry
- 20 Bpost_id 8 B + author_id 8 B + type and flags 4 B: what the entry means, not what Redis spends.
- In-memory bytes per entry
- 100 BAssumption. Over 128 members a sorted set leaves listpack for a skiplist plus a hash table: a text member of up to ~40 B, an 8 B score, a skiplist node and a hash entry.
- Entries kept per feed
- 500About 25 pages of 20; older posts come from author timelines.
- Feeds kept
- 400MEveryone active in the last 30 days.
- Copies of each feed
- 3
- Data per Redis node
- 64 GBAssumption.
- Kai Ren's followers
- 30MA footballer; our example of a very large account.
- Appends per postfollowers × active = 250 × 0.6150from Mean followers of a pushed post and Followers active in the last 30 days
- Appends per dayposts × per-post = 50M × 1507.5B/dayfrom New posts per day and Appends per post
- Average append ratedaily ÷ 86,400 s~86,800/sfrom Appends per day
- Peak append rateavg × peak = 86,806 × 3~260K/sfrom Average append rate and Evening peak over the daily mean
- One feedkeep × entry = 500 × 20 B10 KBfrom Entries kept per feed and Logical bytes per feed entry · Logical size. All feeds would be 400M × 10 KB = 4 TB of raw references.
- One feed in Rediskeep × mem = 500 × 100 B50 KBfrom Entries kept per feed and In-memory bytes per entry
- All feeds, one copyfeeds × per-feed-mem = 400M × 50 KB20 TBfrom Feeds kept and One feed in Redis
- Redis shards (primaries)store ÷ node = 20,000 GB ÷ 64 GB = 312.5~313from All feeds, one copy and Data per Redis node
- All feeds, three copiesstore × replicas = 20 TB × 3; 313 shards × 360 TB on ~939 nodesfrom All feeds, one copy, Copies of each feed and Redis shards (primaries) · Raising zset-max-listpack-entries above 500 keeps each feed in the compact encoding at a fraction of this, at the cost of O(n) inserts; benchmark before counting on it.
- Kai Ren's post, if pushedkai-followers × active = 30M × 0.6 = 18M appends; 18M ÷ 260K/s69 s of the whole fleetfrom Kai Ren's followers, Followers active in the last 30 days and Peak append rate
- Pushing is cheap for the typical author: 150 small appends.
- Keep references, not posts: about 100 B per entry in Redis, against ~1 KB for a post, is what lets 400M feeds live in memory.
- One very large account would take over the whole fleet for a minute, so it can't be pushed like everyone else.
Appends needed for one post
Data
| Author | appends per post |
|---|---|
| Median author (90) | 54 |
| Mean author (250) | 150 |
| Dara (1,840) | 1,104 |
| Cut-off (100K) | 60,000 |
| Kai Ren (30M) | 18,000,000 |
High-level design.
Posting and fan-out are split by a queue. The author waits for one database write; the copies are made afterwards, in parallel, by workers that can fall behind without anyone noticing.
Fan-out on write
Components
| Component | Responsibility | Owns |
|---|---|---|
| API gateway | Authenticates, applies per-user posting limits, routes. | rate-limit counters |
| Post service | Validates the post, assigns a time-ordered ID, writes it, publishes post.created. | idempotency records |
| Post store | Every post once, plus each author's posts newest first. | posts, author_posts |
| Post events | Buffers post and follow events; keeps each author's events in order. | post.events, graph.events |
| Fan-out workers | Page followers, drop inactive ones, append the reference to each feed. | nothing (stateless) |
| Social graph | Answers who follows whom in pages of 1,000. | followers, following |
| Activity store | Last feed open per user, read in batches of 1,000. | last_active |
| Feed store | One capped, ordered set of post references per active user. | feed:{user_id} |
Dara posts a new loaf
- Dara's app → API gateway: POST /v1/posts (Idempotency-Key 7c1e…)
- API gateway → Post service: forward, authenticated
- Post service → Post store: insert post + author_posts row
- Post store → Post service (reply): ok
- Post service → Post events: post.created, key = Dara's author_id
- Post service → Dara's app (reply): 201 {post_id}
- Note over Dara's app: Dara is done: about 60 ms
- Post events → Fan-out workers: deliver post.created
- Fan-out workers → Social graph: followers of Dara, page 1
- Social graph → Fan-out workers (reply): 1,000 IDs; page 2 has 840
- Note over Fan-out workers: 1,104 of 1,840 active: 1,104 appends
- Fan-out workers → Feed store: 313 shard pipelines in parallel, 3-4 each
- Feed store → Fan-out workers (reply): ok
- Note over Fan-out workers: Commit offset after all batches succeed
What one append does
# One pipeline per shard; followers are
# already grouped by shard.
pipe = shard.pipeline(transaction=False)
for uid in followers_on_this_shard:
key = f"feed:{uid}"
# LPUSHX pushes only if the feed exists:
# never create a stub feed.
pipe.lpushx(key, f"{post_id}:{author_id}")
pipe.ltrim(key, 0, 499) # newest 500
pipe.execute()
# LTRIM right after LPUSH removes one
# element: O(1) in practice. But a
# redelivered batch pushes the post twice.
# Plain ZADD on a missing key creates a
# one-entry feed the read path would trust.
# So append only if the key exists.
APPEND = shard.register_script("""
if redis.call('EXISTS', KEYS[1]) == 0
then return 0 end
-- same member twice = one entry
redis.call('ZADD', KEYS[1],
ARGV[1], ARGV[2])
-- keep the newest 500
redis.call('ZREMRANGEBYRANK',
KEYS[1], 0, -501)
return 1""")
# The ID's 41-bit ms timestamp + epoch.
score = created_ms(post_id)
member = f"{post_id}:{author_id}"
pipe = shard.pipeline(transaction=False)
for uid in followers_on_this_shard:
APPEND(keys=[f"feed:{uid}"],
args=[score, member], client=pipe)
pipe.execute()
# Scores are doubles, exact only to 2^53:
# a full 64-bit post_id would lose its low
# bits, so score by the timestamp inside it.
The list is the smallest and fastest structure, and Twitter's home timeline used one. It has two weaknesses in a system with at-least-once delivery: a redelivered batch adds the post again, and entries sit in arrival order rather than post order, so a post from a lagging partition lands above newer ones. A sorted set keyed by the post fixes both at the cost of O(log n) inserts and more memory per entry. Either way the append must not create a missing feed: a key holding one entry looks like a live feed to the read path, which then skips the rebuild. So the list uses LPUSHX and the sorted set a four-line script that checks EXISTS first. Step through the strip to see the difference.
Ines's feed as a list and as a sorted set
- Tomas, 1 slot, Existing entry, value 06:01:40
- Dara, 1 slot, Existing entry, value 05:58:02
- Rowing club, 1 slot, Existing entry, value 05:51:30
- Tomas, 1 slot, Existing entry, value 05:40:15
- 496 more, 3 slots, Older entries (to 500)
- Tomas, 1 slot, Existing entry, value 06:01:40
- Dara, 1 slot, Existing entry, value 05:58:02
- Rowing club, 1 slot, Existing entry, value 05:51:30
- Tomas, 1 slot, Existing entry, value 05:40:15
- 496 more, 3 slots, Older entries (to 500)
- Just appended
- Existing entry
- Duplicate
- Older entries (to 500)
As it starts. 4 steps follow.
Data model.
Posts live once, in the post store. Feeds hold only references, and can be thrown away and rebuilt.
Written once, read by ID from everywhere; the author timeline is the one range query the system needs.
| Column | Type | Key | Note |
|---|---|---|---|
| post_id | bigint | primary | Time-ordered: 41-bit ms timestamp, then shard and sequence bits, so ID order is time order. |
| author_id | bigint | ||
| body | text | ||
| media | json | Blob keys only; bytes live in the blob store. | |
| visibility | enum | public | followers | close_friends | |
| created_at | timestamp | ||
| deleted_at | timestamp nullable |
| Column | Type | Key | Note |
|---|---|---|---|
| author_id | bigint | partition | |
| post_id | bigint | clustering desc |
Two directions of the same edge, each read as a paged list.
| Column | Type | Key | Note |
|---|---|---|---|
| user_id | bigint | partition | |
| follower_id | bigint | clustering | |
| created_at | timestamp |
| Column | Type | Key | Note |
|---|---|---|---|
| user_id | bigint | partition | |
| followee_id | bigint | clustering | |
| muted | bool | ||
| created_at | timestamp |
One small value per user, read 1,000 at a time by fan-out.
| Column | Type | Key | Note |
|---|---|---|---|
| user_id | bigint | primary | |
| last_active_at | timestamp | Last feed open, the same event that refreshes the feed key's TTL. Written at most once an hour, so it can lag the TTL by an hour; the conditional append covers that gap. |
Every feed read starts here; it fits in RAM (20 TB a copy) because entries are references, not posts.
| Column | Type | Key | Note |
|---|---|---|---|
| member | text | post_id:author_id | |
| score | double | Created time in ms, from the ID's top 41 bits. |
| Query | Uses | How |
|---|---|---|
| Who follows Dara? | followers | partition user_id, paged by follower_id, 1,000 per page |
| Append to Ines's feed | feed:{user_id} | script: if EXISTS, ZADD then ZREMRANGEBYRANK 0 -501 |
| Dara's last 20 posts | author_posts | partition author_id, first 20 rows |
| Is this follower active? | last_active | multi-get of 1,000 user IDs |
Keyed by author so a delete can never overtake its create. Every feed-service node also tails this topic for large-account posts and deletes (feed-reads).
follow backfills the followee's last 20 posts into the follower's feed; unfollow, mute and block scan the follower's 500 entries and remove that author's, then repeat 60 s later. This topic is keyed by follower and post.events by author, so an append in flight can land after the cleanup: the read path also filters muted and blocked authors (feed-reads), and a stray post from someone just unfollowed can survive until the repeat pass or the trim.
A big author's post split into 1,000-follower batches so many workers share it (see Optimizations).
The life of one reader's feed key
States of9Feed store
Fan-out never creates a feed key; only the rebuild does. So a reader whose key is missing, for any reason, has no feed at all and gets a fresh one on her next open, never a feed that silently missed posts.
| From → To | Event | Guard | Action | Actor |
|---|---|---|---|---|
| No feed key → No feed key | post.created | nothing (filtered, or no key) | Fan-out workers | |
| No feed key → Rebuilding | app opened | key + placeholder, TTL 60 s | Feed service (feed-reads) | |
| Rebuilding → Rebuilding | post.created | ZADD lands | Fan-out workers | |
| Rebuilding → Live and pushed to | rebuild done | 500 refs in, placeholder out | Feed service (feed-reads) | |
| Live and pushed to → Live and pushed to | post.created | key exists | ZADD + trim | Fan-out workers |
| Live and pushed to → Live and pushed to | feed read | refresh TTL to 30 d | Feed service (feed-reads) | |
| Live and pushed to → No feed key | 30 d unread | key expires | Redis | |
| Live and pushed to → No feed key | shard lost | Redis cluster |
- No feed keystart
- New user, 30 days without a feed read, or shard lost
- Rebuilding
- The read path refills it from author timelines (feed-reads)
Interface.
One public call to post, one to delete, two for following. Everything else is internal: events and a paged follower list.
Stores the post and publishes post.created. Returns before any follower's feed is touched.
Endpoints
| Method | Path | Does | Returns |
|---|---|---|---|
| POST | /v1/posts | Create a post; fan-out follows asynchronously. | 201 { post_id, created_at } |
| DELETE | /v1/posts/{id} | Mark the post deleted and emit post.deleted. | 204 |
| POST | /v1/follows | Follow { followee_id }; emits follow, which backfills recent posts. | 201 |
| DELETE | /v1/follows/{followee_id} | Unfollow; emits unfollow. | 204 |
| GET | /internal/users/{id}/followers?cursor&limit=1000 | Internal; the page fan-out reads. | { ids[], next_cursor } |
Errors
| HTTP | Type | Body code | Client behaviour |
|---|---|---|---|
| 400 | error | invalid_post | Empty, over 2,000 characters, or more than 4 photos. Show the message; keep the draft. |
| 403 | error | not_author | Deleting someone else's post. Nothing to retry. |
| 409 | error | idempotency_mismatch | The key was used with a different body. Generate a new key only if this is really a new post. |
| 413 | error | media_too_large | Reject before upload; tell the author the limit. |
| 429 | retry | rate_limited | Wait for Retry-After. Posting limits stop a spam flood from turning into a fan-out flood. |
Optimizations.
What keeps fan-out lag under ten seconds when posting spikes.
A big post on Dara's partition
Scenario 1 of 2: As described.
Timeline as a list
A big post on Dara's partition: 3 lanes, from 0 s to 4.5 s.
- 0–3.6 s · Partition 17 · Harbour FC: 90 pages + 54K appends
- 0.3 s · Partition 17 · Dara posts
- 0.3–3.68 s · all lanes · window: Dara waits 3.4 s
- 3.6–3.68 s · Partition 17 · Dara: 1,104
- 3.68 s · Feeds · Dara delivered (delayed)
- 01Launch
- Load
- 50 posts/s
- Bottleneck
- None
- Change
- Feeds computed on read with one SQL query over following JOIN posts; no feed store.
- Adds
- 4Post store7Social graph
- 0210× launch
- Load
- 500 posts/s
- Bottleneck
- Reads slow down as follow counts grow
- Change
- Add the feed store and push synchronously from the post service.
- Adds
- 9Feed store
- 03Today's peak
- Load
- ~1,750 posts/s
- Bottleneck
- Posting latency tied to the author's follower count
- Change
- Move the push behind a queue with workers; skip inactive followers.
- Adds
- 5Post events6Fan-out workers8Activity store
- 04Today, with very large accounts
- Load
- ~1,750 posts/s, 260K appends/s
- Bottleneck
- One post stalls the fleet or its partition
- Change
- Pull-mode accounts over 100K followers; two-stage batched fan-out.
Trade-offs.
The chosen option is first; the others stay visible so the reasoning can be checked.
- Pro:A read is one lookup plus a few merges
- Pro:The cost of any single post is bounded (60K appends at most)
- Con:Two code paths to build, test and keep consistent
- Con:Readers who follow many large accounts pay extra merge latency
One large account's post is millions of writes and minutes of lag for everyone else; Writes to feeds nobody opens
Every page asks hundreds of author timelines and merges them; reads outnumber posts about 50:1 here (2.4B page loads vs 50M posts a day); Works at scale only with heavy in-memory infrastructure, as Facebook's Multifeed does: aggregators query leaf servers holding recent actions at read time
Decide per pair, not per author
The cost of pushing Tomas's posts to Ines is Tomas's posting rate: one write per post. The cost of pulling them is Ines's reading rate: one timeline query per feed read. Push the pair when the author posts less often than the reader reads, weighted by the relative cost of a write and a read-time query. Silberstein et al. (Feeding Frenzy, SIGMOD 2010) show that making this choice per producer-consumer pair from the ratio of rates minimises total cost.
Worked, with a write and a query costing the same: Ines opens her feed 8 times a day. Tomas posts once a day, so push: one write serves 8 reads. @metro_alerts posts 150 times a day, so pull: 150 writes a day, most scrolled past, which would also push Tomas's posts out of her 500 entries. Follower count, which Chorus uses, is only the easy proxy for this rule.
One author and one reader, per day
- Push
- Pull
Data
| Author's posts per day | Push | Pull |
|---|---|---|
| 0.1 | 0.1 | 8 |
| 1 | 1 | no value |
| 8 | 8 | no value |
| 150 | 150 | no value |
| 1,000 | 1,000 | 8 |
- 8 posts/day: Author's posts per day = 8
- At 1: Tomas: push
- At 150: @metro_alerts: pull
- Pro:Idempotent; a redelivered append overwrites itself
- Pro:Ordered by post time even when partitions lag
- Pro:Small entries (a reference, ~100 B in Redis)
- Con:Needs hydration on read
- Con:O(log n) inserts and more memory per entry than a list
- Con:History beyond 500 comes from author timelines
Duplicates on redelivery; Arrival order, not post order; Removing one entry is O(n)
About 50× the raw bytes (1 KB post vs 20 B reference); Edits and deletes must touch every copy
- Pro:Instant for readers: hydration sees deleted_at, and a small set of recently deleted IDs is checked
- Pro:No second fan-out storm
- Con:Dead references occupy slots until trimmed
As expensive as the post itself; Still races with appends in flight, so the read filter is needed anyway
| Failure | Impact | Detection | Mitigation | Meanwhile |
|---|---|---|---|---|
| Workers fall behind (lag 2 min)6Fan-out workers | Posts appear late in feeds | Consumer lag per partition against the 10 s p95 target | Autoscale workers; split big posts into batches sooner | Posts are still on the author's profile at once |
| A shard and its replica lost9Feed store | Those users' feeds are empty | Shard health checks | Appends skip the missing keys; rebuild on read from author timelines (feed-reads) | A slower first page for about 0.3% of users (one shard of 313) |
| Duplicate delivery after a worker crash5Post events | None visible | Not needed | ZADD is idempotent per post | No change |
| Follower lookups slow7Social graph | Fan-out lag grows for every author | p99 of follower-page latency | Retry with backoff; posts wait safely in the queue for up to 3 days | Feeds lag; posting is unaffected |
| Activity store unavailable8Activity store | Fan-out can't tell active from inactive followers | Error rate on the batch lookup | Fail open and push to all followers of small accounts; delay large ones | About 1.7× the append load until it recovers |