Fan-out: on write vs on read
Choose whether distribution happens on write, on read, or with a measured hybrid boundary.
When to reach for this
Reach for this when…
- Social graph with asymmetric follow counts (power-law follower distribution)
- Newsfeed / activity stream / home timeline design
- One-to-many notification delivery at scale
- Any design where a single write must be visible to M other users
- "Design Twitter" / "Design Instagram feed" / "Design LinkedIn feed"
Not really this pattern when…
- Bounded small groups only (team chat, <50 members) — just materialise the group inbox directly
- No follow relationship — user pulls their own data, not other users' data
- Symmetric messaging (1:1 DMs) — that's a queue, not a fan-out
- Pub/sub where consumers are services, not users — that's producer-consumer
Good vs bad answer
Interviewer probe
“How would you design a Twitter-like home timeline?”
Weak answer
"I'd use a database to store tweets and when a user wants their feed, I'd query all the tweets from people they follow, sorted by time. For scale I'd add caching. If there are a lot of followers I'd use a message queue to handle the writes asynchronously."
This answer doesn't identify the core tradeoff (push vs pull), doesn't mention the celebrity problem, proposes no concrete data structures, and conflates "caching" with architectural decisions. The candidate hasn't shown they understand why the naive approach breaks.
Strong answer
"This is a fan-out problem — the question is whether to pay distribution cost at write time or read time. Pure push (fan-out on write into Redis ZSET inboxes) gives cheap reads but fails when one high-follower author creates a huge write storm. Pure pull gives cheap writes but fails on power users following thousands of accounts. The production answer is a hybrid: push when active-follower fan-out fits the queue and freshness SLO, pull when read-time merge is cheaper.
Write path: post service → fan-out router → cost model using active followers, post rate, queue depth, and freshness SLO → push into bounded inboxes or store in author timeline only. Read path: ZREVRANGE the inbox, pull-fetch selected author timelines with a per-reader cap, merge-K with a heap, apply policy filters and dedupe, rank, and return a page."
This answer names the pattern, identifies the structural tradeoff, avoids a magic threshold, proposes concrete data structures (Redis ZSET), and sketches both paths of the hybrid. It demonstrates principal-engineer-level reasoning.
Why it wins: The strong answer identifies the structural tradeoff, proposes the hybrid with a cost-based boundary, and sketches concrete data structures — showing the candidate can reason about cost distributions, not just draw boxes.
Cheat sheet
- •Push: O(1) read, O(followers) write. Pull: O(following) read, O(1) write.
- •Hybrid with a measured push/pull boundary is the production answer for skewed graphs.
- •Inbox = Redis ZSET or equivalent, bounded via add + trim.
- •Merge-K at read: heap of inbox + K pull-path author timelines, O(N log K).
- •Spread fan-out: chunk follower list and drain steadily for mid-tier high-follower authors.
- •Follow → optional backfill. Unfollow → lazy filter + sweep. Deletes/privacy/moderation → tombstone or read-time policy filter.
- •Ranking pipeline: candidates → features → ML score → diversity filter.
- •Inactive user inboxes: 7-day TTL, reconstruct on next login.
- •Multi-region: regional fan-out workers + regional inbox Redis clusters.
- •Dynamic boundary with hysteresis prevents push/pull oscillation.
- •Fan-out workers are idempotent — ZADD same member+score is a no-op, safe to retry.
- •High-follower author timelines are often cache-hot, but measure hit rate before promising latency.
Core concept
Fan-out is the one-to-many distribution problem behind social feeds, activity feeds, notifications, repo watchers, and channels. When a user publishes an item, who pays the distribution cost: the writer, each reader, or a carefully chosen mix?
Fan-out on write (push) pays at write time. The moment a user publishes, fan-out workers copy the post identifier into follower inboxes, often Redis sorted sets keyed by user id and scored by timestamp. From the reader's perspective the feed is already materialised: open the app, read the inbox, hydrate the posts. Push gives O(1) reads. Its cost is write amplification: one author write becomes roughly active followers × eligible posts inbox writes. A high-follower account can turn one post into tens of millions of inbox writes.
Fan-out on read (pull) pays at read time. Posts are stored once in the author's timeline. When a reader opens their feed, the feed service fetches recent posts from followed authors, merges them with a heap (O(N log K), where K is the number of sources), ranks, and returns a page. Writes are O(1). Reads now scale with the reader's follow count and cache hit rate. A reader following thousands of accounts can turn one refresh into thousands of timeline lookups.
Neither shape survives all workloads alone. Attention networks are usually power-law: most authors have modest reach, a small set has very high reach, and a few readers follow huge numbers of authors. The system needs to keep the median path cheap without letting either tail dominate capacity.
The canonical production architecture: push for normal users, pull for celebrities, merge at read.
The hybrid approach partitions authors by cost. Push ordinary authors into follower inboxes. Store very high-follower or high-post-rate authors in their own timelines and pull them at read time for the readers who follow them. At feed-read time, merge the reader's pre-materialised inbox with a bounded set of pulled author timelines, filter, rank, and return the feed.
The boundary is not "1M followers" by magic. Model it as a capacity and freshness decision:
| Signal | Push path cost | Pull path cost | What moves the boundary |
|---|---|---|---|
| Active followers | active followers × posts | follower reads × merge cost | More inactive followers favour pull less than raw follower count suggests |
| Author post rate | more fan-out work per author | more fresh items to merge | High post rate lowers the safe push threshold |
| Feed-read rate | amortises pushed writes | increases pull reads | Frequent readers favour push for authors they follow |
| Fan-out capacity | queue seconds and write IOPS | little effect | Tight capacity lowers the push threshold |
| Freshness SLA | must finish fan-out before SLA | fresh on next read | Stricter freshness can require push or prioritised pull |
Correctness rules ride alongside the cost model: deletes usually write a tombstone and filter at read time before asynchronous cleanup; blocks, mutes, privacy changes, and moderation decisions must be applied at read time even for pre-materialised inbox entries; follow creates optional backfill; unfollow uses lazy filtering plus later sweep; dedupe protects retries and merge overlap. Ranking consumes candidates from both paths, so push is a candidate-generation optimisation, not the final order.
Twitter's public 2012 QCon talk, "Timelines at Scale", is the classic historical reference: it framed timeline delivery as thousands of events per second and explained why high-follower accounts stress write-time fan-out. Use it as a pattern source, not as a universal threshold. Modern systems should tune from their own active-follower ratios, author post rates, follower feed-read rates, fan-out capacity, and freshness SLOs.
Interview walkthrough
Worked example: Design a Twitter-like home timeline
Step 1: Recognise the pattern. The prompt says "home timeline" or "news feed" — this is fan-out. The key question is: when a user posts, how does the post reach their followers' feeds? Immediately say: "This is a fan-out problem. The core tradeoff is push vs pull, and at scale we'll need a hybrid."
The canonical production architecture: push for normal users, pull for celebrities, merge at read.
Step 2: Estimate the volumes. Use round numbers and keep the units honest. Example: 200M DAU × 5 feed opens/day = 1B feed reads/day ≈ 11.6K reads/sec. 500M posts/day ≈ 5.8K author writes/sec. The user-level read:write request ratio is about 2:1, not thousands-to-one. The push-path cost comes from write amplification: if an eligible post fans out to 300 active followers on average, 500M posts/day would imply 150B inbox writes/day ≈ 1.7M inbox writes/sec before filtering, chunking, or excluding pull-path authors.
Step 3: Choose the hybrid. Pure push fails for high-follower authors: one post can expand into tens of millions of inbox writes. Pure pull fails for high-fan-in readers: a user following thousands of accounts can trigger thousands of timeline reads per refresh. The hybrid pushes where write amplification fits capacity and pulls where read-time merge is cheaper. State the boundary as a cost model, not a magic follower count.
Step 4: Sketch the write path. User posts → post service persists to posts DB → enqueue fan-out task → router evaluates active followers, post rate, queue depth, fan-out capacity, and freshness SLO. Push path: paginate through active followers, pipeline ZADD operations to bounded inboxes. Pull path: store in the author's timeline only. Always store in the author's own timeline for profile reads, pull reads, and deep scroll.
Step 5: Sketch the read path. User opens app → GET /feed → read bounded inbox → pull recent posts from followed pull-path authors, with a cap per reader → merge-K with heap → apply deletes, tombstones, blocks, mutes, privacy, moderation, and dedupe → rank candidates → return page → cache briefly with invalidation on new inbox writes.
Step 6: Address the high-follower account problem. The interviewer will ask what happens when a very large account posts. Answer: "That author is on the pull path, so the post does not enqueue millions of inbox writes. It goes into the author's timeline. Followers see it when their feed read merges pull-path authors, subject to freshness and ranking. If a mid-tier author is still on push, we chunk fan-out and drain at a controlled rate."
Step 7: Address ranking. "We don't return a chronological feed — we rank. The inbox and pulled timelines provide candidates. We enrich each candidate with features from the feature store (engagement rate, recency, author affinity) and score with a lightweight model within the read-path latency budget."
Step 8: Address follow/unfollow churn and policy changes. "On follow: optionally backfill the new followee's recent posts into the follower's inbox. On unfollow: filter at read time and sweep later. Deletes use tombstones so old inbox entries disappear immediately. Blocks, mutes, privacy changes, and moderation decisions are enforced at read time even if stale entries remain in the inbox."
This walkthrough hits every signal an interviewer looks for: pattern recognition, honest volume math, principled tradeoff reasoning, concrete data structure choices (bounded inbox, author timeline, heap merge), failure mode awareness, and operational concerns (churn, policy filters, caching, ranking).
Interview playbook
When it comes up
- "Design a news feed" or "Design Twitter home timeline"
- "Design an activity feed" or "Design Instagram feed"
- "How would you deliver a post to all followers?"
- Any mention of asymmetric follow graphs at scale
- "Design a notification system" (bounded fan-out variant)
Order of reveal
- 11. Name the pattern. This is a fan-out problem — the core question is whether to pay the distribution cost at write time (push) or read time (pull).
- 22. Frame the tradeoff. Push gives O(1) reads but O(followers) writes. Pull gives O(1) writes but O(following) reads. The follower distribution is power-law, which means neither alone works.
- 33. Propose the hybrid. The production answer is a hybrid: push authors whose active-follower fan-out fits capacity; pull authors whose write amplification would dominate the queue. Merge at read time.
- 44. Sketch the write path. Post → fan-out router → cost model → push into bounded Redis ZSET inboxes or store in author timeline only.
- 55. Sketch the read path. Read inbox + pull selected author timelines → merge-K with min-heap → policy filter + dedupe → rank → return page.
- 66. Address failure modes. High-follower storms are handled by pull or spread fan-out. Stale follows are handled by lazy filter + sweep. Deletes/privacy/moderation are enforced at read time.
- 77. Discuss scaling levers. Multi-region inboxes for global latency. Ranking model complexity vs latency budget. Dynamic cost boundary tuning. Inbox TTL for inactive users.
Signature phrases
- “Push and pull are two regimes of the same cost function; the boundary is capacity plus freshness, not folklore.” — Shows you understand the structural reason for the hybrid, not just the implementation.
- “The high-follower problem is not an edge case in a power-law graph; it is the constraint that shapes the architecture.” — Reframes what sounds like a corner case as a central design driver.
- “Inbox bounded to 800 entries via ZADD + trim — old content falls back to author timelines.” — Demonstrates concrete Redis mechanics and the bounded-storage invariant.
- “Spread fan-out: chunk the follower list and drain at a steady rate to avoid thundering-herd writes.” — Shows awareness of the mid-tier celebrity problem beyond just the threshold.
- “Lazy filter on unfollow, nightly sweep for space reclamation — correctness now, efficiency later.” — Demonstrates pragmatic engineering: immediate consistency via cheap filtering, deferred cleanup via batch.
- “The merge step is O(N log K) where K is the number of celebrity pulls — typically under 20.” — Precise complexity analysis shows you've thought about the read-path cost.
Likely follow-ups
?“What if a normal user goes viral and suddenly has 5M followers?”Reveal
The routing model should notice that active-follower fan-out, post rate, or queue cost has crossed the safe boundary. Future posts route through pull or spread fan-out. In-flight fan-out for the current post can continue idempotently, but new chunks should respect queue backpressure. Hysteresis and cooldowns prevent oscillation if the author's measured cost hovers near the boundary.
?“How do you handle unfollow?”Reveal
Lazy filter at read time: when assembling the feed, check each candidate post's author against the reader's current follow set (cached in a bloom filter or set lookup). Posts from unfollowed authors are filtered out before ranking. A nightly batch sweep scans inboxes and removes entries from authors the user no longer follows, reclaiming Redis memory. This avoids expensive real-time deletion while ensuring immediate correctness.
?“How does ranking work in this system?”Reveal
The merge step produces 100-200 candidates (inbox + celebrity pulls). Each candidate is enriched with features from the feature store: post recency, author engagement rate, reader-author interaction history, content type, early engagement signals (likes in first 5 minutes). A lightweight ML model (gradient-boosted tree or small neural net) scores each candidate by P(engagement). A diversity filter ensures no author dominates and content types are mixed. Total ranking latency: <20ms for the full candidate set.
?“How would you handle multi-region deployment?”Reveal
Partition followers by region. When a post arrives, the region router dispatches fan-out tasks to regional workers. Each region has its own Redis inbox cluster, so feed reads are local. Celebrity timeline pulls may cross regions — cache them in each consuming region with a 30-second TTL. The follow graph is globally consistent (replicated) but fan-out is regional. Cross-region replication lag for celebrity content is acceptable because users don't expect sub-second freshness for celebrity posts.
?“Celebrity handling — why not just use pull for everyone?”Reveal
Pure pull means every feed read triggers O(following) timeline lookups. For a user following 500 accounts, that's 500 reads per feed refresh. Even with caching, the p99 latency is unacceptable — a single cache miss on one of 500 timelines adds 20-50ms. Push pre-materialises the feed for the common case (normal-follower-count authors), reducing feed reads to O(1) for ~95% of the content. The hybrid pays the push cost for the 95% where it's cheap and uses pull only for the 5% where push would be catastrophically expensive.
?“Real-time vs batch ranking — which one?”Reveal
Real-time ranking at feed-read time for the primary feed. The candidate set is small enough (100-200 items) that scoring completes in <20ms. Batch ranking (pre-computing ranked feeds offline) doesn't work because the candidate set changes with every new post and follow/unfollow event — the precomputed feed would be stale by the time it's read. However, some features used by the ranking model ARE batch-computed: author engagement rates, topic embeddings, and user interest profiles are updated by offline pipelines every few hours and served from the feature store.
Canonical examples
- →Twitter / X home timeline
- →Instagram feed
- →LinkedIn activity feed
- →Activity notifications (mentions, likes, retweets)
- →Push notification fan-out to mobile devices
- →GitHub activity / notification feed
Variants
Fan-out on write (push)
Copy the post id into every follower's inbox at write time.
Fan-out on write is conceptually the simplest feed architecture: when a user publishes a post, a background worker iterates over their follower list and inserts the post identifier into each follower's inbox. The inbox is often a Redis sorted set (ZSET) keyed by the follower's user id, with the post timestamp or monotonic id as the score. With bounded inboxes, insertion and trim costs stay predictable.
The write path looks like this: the post service persists the post to the posts database, then enqueues a fan-out task. The fan-out worker reads the author's follower list (paginated, from a follow-graph service), pipelines batches of ZADD commands to the inbox Redis cluster, and trims the inbox to its maximum size. Pipelining keeps network round trips low; the real limit is aggregate write IOPS and queue latency, not the single-command API.
When a user posts, the fan-out worker copies the post id into every follower inbox. Reads are O(1).
For the median author, this is excellent. A post to a few hundred active followers is cheap to copy, and feed reads become one bounded inbox read plus hydration. For a high-follower account, one post can expand into tens of millions of inbox writes. That queue work competes with every other author's delivery and can make unrelated feeds stale. This is the high-follower post storm: the defining failure mode of pure push.
There's a subtler problem too: wasted storage. In any social network, many followers are inactive. Pure push writes into their inboxes anyway unless you filter by active follower or expire inactive inboxes. That is why production push paths usually combine active-follower targeting, bounded inboxes, TTLs, idempotent writes, and read-time policy filters.
Despite these limitations, pure push remains the right answer for bounded fan-out scenarios: team activity feeds, GitHub repo watchers, notification delivery to device tokens, and any system where the maximum recipient count is architecturally bounded or the inactive-recipient waste is acceptable.
Pros
- +O(1) reads — the feed is pre-materialised in the inbox
- +Read latency is predictable and low (~5ms for a ZREVRANGE)
- +Ranking and filtering can be baked in at write time
- +Simple read path — single data source, no merge logic
Cons
- −Celebrity posts are catastrophic — O(followers) writes per post
- −Write amplification = followers × posts per day
- −Inactive followers waste storage and write bandwidth
- −Fan-out queue backlog delays all users during celebrity storms
Choose this variant when
- Most users have <100K followers (bounded fan-out)
- Feed reads need very low latency and most authors have bounded active followers
- You control the follow graph shape (e.g., team-based, not social)
- The system has no celebrity / power-law problem
Fan-out on read (pull)
Store posts once in the author's timeline; gather and merge at read time.
Fan-out on read inverts the cost structure: writes are O(1) — the author appends a post to their own timeline — and reads pay the aggregation cost. When a reader opens their feed, the feed service looks up the reader's follow list, fetches the last N posts from each followed account's timeline, merges them into a single sorted stream, and returns the top page.
The merge operation is a classic K-way merge. You have K sorted streams (one per followed account), each producing posts in reverse-chronological order. A min-heap of size K lets you extract the global top-N in O(N log K) time. In practice, K is the number of accounts the reader follows — for a typical user following 200 accounts, this means a heap of 200 entries and ~50 extractMin operations to fill a page.
Posts live in author timelines. At read time the feed service gathers, merges, and ranks.
This was Facebook's early approach (before the shift to algorithmic ranking): at news feed render time, pull the latest posts from each friend's wall, merge, and display. For the median user with ~150 friends, this was manageable. But the p99 user follows 2,000+ accounts, and the feed service must issue 2,000 parallel timeline reads, wait for the slowest one, then merge. Even with aggressive caching (each timeline cached in memcached with a 60-second TTL), the tail latency was painful.
The deeper problem is cache efficiency. In a push model, the inbox is a single hot key per user — extremely cache-friendly. In a pull model, each followed account's timeline is a separate cache entry, and the hit rate depends on how many other readers are also following that account. Celebrity timelines are hot (millions of readers pull them) and stay cached. But the long tail of low-follower accounts has poor cache hit rates, and a cache miss means a database read for a timeline that might be read once and evicted.
Pull does have genuine advantages. Writes are trivially fast — one append, no fan-out queue, no worker fleet. There's no wasted work on inactive readers. And when you add algorithmic ranking (which requires scoring each candidate against the reader's interest model), pull naturally produces the candidate set you need to rank. Push pre-materialises a chronological inbox that you then have to re-rank anyway, which partly defeats the purpose.
Pull is the right answer when the fan-in is low (users follow <100 accounts), when celebrity content dominates the feed, or in read-light workloads where most users check their feed infrequently. It's also the simpler starting point — no fan-out infrastructure, no inbox storage, just timelines and a merge service.
Pros
- +O(1) writes — no fan-out queue, no worker fleet
- +No wasted work on inactive readers
- +Naturally produces candidate set for algorithmic ranking
- +No inbox storage amplification
Cons
- −Read latency scales with fan-in (number of followed accounts)
- −Cache hit rate per-followed-account varies wildly
- −K-way merge adds CPU cost at read time
- −Tail latency from slowest timeline read dominates p99
Choose this variant when
- Users follow <100 accounts on average
- Celebrity / one-to-many broadcast dominates the content mix
- Algorithmic ranking is already required (pull provides candidates naturally)
- Write volume is high relative to reads (rare in social, common in logging)
Hybrid (push + pull threshold)
Push for normal users, pull for celebrities, merge at read time.
The hybrid is not a compromise. It is a cost model with two execution paths. Push is cheap for authors whose active-follower fan-out can finish inside the freshness SLO. Pull is cheaper for authors whose write-time fan-out would dominate the queue but whose followers can tolerate a bounded merge at read time.
A practical router evaluates a few signals, not just raw follower count: active followers, post rate, follower feed-read rate, fan-out worker capacity, storage budget, and freshness SLO. A high-follower author who rarely posts to mostly inactive followers might still be safe to push in chunks; a lower-follower author with high post rate and highly active followers may need pull or spread fan-out.
At feed-read time, the service reads the user's pre-materialised inbox and also pulls recent posts from the small set of followed authors assigned to the pull path. The merged result is filtered for policy and relationship changes, deduped, and passed to ranking. Keep the pull set bounded per reader; if it grows, use cached author timelines, candidate limits, or move more authors back to push.
Twitter's Raffi Krikorian described the push/pull timeline split in the 2012 QCon "Timelines at Scale" talk. The durable lesson is the shape of the tradeoff, not a fixed public threshold. The exact boundary should be measured continuously and changed with hysteresis so authors do not oscillate between paths.
The complexity cost is real: two code paths, two storage layers, a merge step at read time, read-time policy filters, and operational tuning for path assignment. The benefit is that normal authors keep O(1)-style feed reads while high-follower authors do not monopolise the fan-out queue.
Pros
- +High-follower accounts do not destroy the write path — no unbounded write storms
- +Normal users still get O(1) reads for most of their feed content
- +Bounded number of celebrity pulls at read time (typically <20)
- +Matches the publicly described Twitter timeline split and the general cost model used by large feeds
Cons
- −Two code paths with different failure modes
- −Path-assignment tuning is ongoing (static thresholds break during viral moments)
- −Merge + rank at read time adds latency and complexity
- −Monitoring for threshold crossings adds operational overhead
Choose this variant when
- Any feed where follower or subscriber distribution is strongly skewed
- Push queue latency and pull merge latency must both stay bounded
- The product needs fast feed reads and support for high-follower authors
Activity stream (non-social fan-out)
System-event fan-out where the source set is bounded — push always works.
Not every fan-out problem has a celebrity problem. Activity streams — GitHub's notification feed, Slack's channel activity, JIRA's project activity log — differ from social feeds in one structural way: the source set is bounded by system design, not by user behaviour.
On GitHub, a repository might have 10,000 watchers. When someone pushes a commit, the activity service fans out a notification to those 10,000 watchers' activity feeds. The fan-out is bounded because repository watchers are bounded (and self-selected). There is no "celebrity repository" with 100 million watchers that breaks the push path. Slack channels have a hard member limit. JIRA projects have bounded team sizes.
In these systems, pure fan-out on write is the correct answer and the hybrid complexity is unnecessary. The architecture is straightforward: event source → message queue → fan-out worker → per-user activity inbox (often a bounded list in Redis or a DynamoDB partition). The fan-out is measured in thousands, not millions, and completes in milliseconds.
The key difference from social fan-out is that the architect controls the maximum fan-out by controlling the maximum group size. If you can enforce that no single source fans out to more than N recipients, and N is bounded by product rules (not by user growth), pure push is simpler, cheaper, and more predictable than any hybrid.
Scaling path
V1: Single database, pull at read time
Ship a working feed with zero infrastructure beyond the application database.
Posts and follows tables, feed built at read time with a join.
Store posts in a posts table, follows in a follows table. When a user requests their feed, run a join: SELECT p.* FROM posts p JOIN follows f ON p.author_id = f.followed_id WHERE f.follower_id = ? ORDER BY p.created_at DESC LIMIT 50. This is pure pull — no fan-out, no inbox, no workers.
This works for the first 10K users. The join is fast when the follows table fits in memory and the posts table is indexed on (author_id, created_at). Read latency is 10-50ms depending on the number of followed accounts.
What triggers the next iteration
- The join becomes expensive as the follows table grows — each feed read scans O(following) index entries
- No caching layer — every feed read hits the database
- Cannot support real-time delivery — users must refresh to see new posts
- Ranking requires sorting the merged result on every read
V2: Materialised inbox with Redis ZSET (push)
Eliminate the read-time join by pre-materialising each user's feed into a Redis sorted set.
Fan-out worker writes into Redis ZSET inboxes; feed reads are O(1).
Add a fan-out worker: when a user posts, iterate their follower list and ZADD the post id (scored by timestamp) into each follower's inbox ZSET. Feed reads become a single ZREVRANGE call — O(1) regardless of follow count.
Trim inboxes to 800 entries with ZREMRANGEBYRANK after each ZADD batch. Older content is still available via the author's timeline (fallback for deep scroll). This handles millions of users with sub-10ms feed reads.
What triggers the next iteration
- Celebrity posts generate millions of ZADD operations, backing up the fan-out queue
- Inactive users waste Redis memory — their inboxes are written but never read
- Fan-out latency for high-follower authors can exceed 60 seconds
- No ranking beyond chronological order
V3: Hybrid fan-out + ranking service
Survive high-follower posts and add algorithmic ranking.
Push for normal users, pull for celebrities, ranking service merges and scores.
Introduce a measured push/pull boundary. Cheap authors use push (fan-out on write into inboxes). Expensive authors skip fan-out and store in their author timeline only. At read time: merge the pre-materialised inbox with pull-fetched author timelines, then pass candidates through a ranking service that scores by recency, engagement prediction, and affinity.
The ranking service receives a bounded candidate set from the inbox and pulled timelines, enriches candidates with features from the feature store, runs an ML scoring model, applies diversity filtering, and returns a ranked page.
What triggers the next iteration
- Single-region deployment means cross-region read latency for global users
- Ranking model cold-start for new users with sparse interaction history
- Fan-out still takes time for mid-tier high-follower authors
- Follow/unfollow churn creates stale entries in inboxes
V4: Multi-region fan-out with regional inboxes
Serve global users with <50ms feed reads from the nearest region.
Posts originate in one region and are replicated to regional fan-out workers with regional inboxes.
Partition followers by region. When a post arrives, the region router splits the follower list into regional chunks and dispatches fan-out tasks to regional fan-out workers. Each region maintains its own inbox Redis cluster. Feed reads are served entirely from the local region — no cross-region calls for the inbox read path.
Pull-path author timelines may still cross regions if the author timeline lives elsewhere, so cache hot timelines in each consuming region or replicate them asynchronously. Make the freshness tradeoff explicit instead of promising sub-second global delivery for every high-follower post.
What triggers the next iteration
- Cross-region pull-path timeline cache invalidation adds complexity
- Follow-graph partitioning must track user region changes
- Regional inbox divergence during network partitions
- Operational cost of managing Redis clusters in 3+ regions
Deep dives
Choosing and managing the push/pull boundary
Path assignment is the most important tuning parameter in a hybrid fan-out system. Set the push boundary too low and feed reads pull too many author timelines. Set it too high and high-follower posts storm the fan-out queue. Treat the boundary as a cost model:
| Term | Meaning | Push-favouring direction | Pull-favouring direction |
|---|---|---|---|
| active_followers | followers likely to read soon | low | high |
| post_rate | posts per author per hour | low | high |
| feed_read_rate | reads by those followers per hour | high | low |
| fanout_capacity | safe inbox writes per second | high | low |
| freshness_slo | max delay before item should appear | strict with capacity | relaxed or read-triggered |
Approximate the decision as: push cost ≈ active_followers × post_rate, bounded by fan-out capacity and freshness SLO. Pull cost ≈ follower_feed_reads × merge_cost, bounded by the number of pull-path authors each reader follows. Push while the write work fits the queue budget and saves more read work than it creates; pull when write amplification would dominate capacity.
Four tiers of users and their fan-out strategy. The boundary is never arbitrary — it is an engineering partition.
Static thresholds break during viral moments. A user may gain followers quickly, post more often, or attract more active readers. If the threshold check happens only at post time, the next post can suddenly create a much larger fan-out than the queue budget expected. The reverse is also a problem: an author oscillates around the boundary, and some posts are pushed while others are pulled, creating inconsistent delivery and operational noise.
Use hysteresis and cooldowns so authors do not bounce between paths. Promote to pull only after the model stays above the boundary for several windows, and demote only after sustained lower cost. Evaluate at post time for routing, but update model inputs continuously from follower activity, queue depth, post rate, and read latency.
What happens when a 100M-follower celebrity posts into a pure-push system.
The storm diagram shows what happens without a boundary: one high-follower post expands into a huge number of inbox writes, saturates the queue, and delays unrelated delivery. With a measured boundary, that post goes to the author's timeline and is pulled by followers when they read, while ordinary authors keep the cheap pre-materialised path.
Inbox storage: Redis ZSET mechanics and memory budget
The inbox is the core data structure of fan-out on write. Each user's inbox is a Redis sorted set (ZSET) where members are post ids and scores are timestamps (or snowflake ids for ordering). The ZSET provides three critical operations:
- 1ZADD — insert a post id with a score. O(log N) where N is the set size. With bounded inboxes (N ≤ 800), this is effectively O(1).
- 2ZREVRANGE — read the top K entries by descending score. O(K + log N). For a feed page of 50 items: sub-millisecond.
- 3ZREMRANGEBYRANK — trim the set to keep only the top N entries. Called after ZADD to enforce the size bound.
Fan-out workers ZADD into bounded sorted sets; feed reads via ZREVRANGE; deep scroll falls back to author timelines.
The fan-out worker pipelines these operations for efficiency. A typical batch: ZADD inbox:{user_id} {timestamp} {post_id} for each of 50 follower inboxes in a single Redis pipeline, followed by ZREMRANGEBYRANK inbox:{user_id} 0 -(max_size+1) to trim. The pipeline completes in under 1ms for 50 operations.
Memory budget: each ZSET entry costs ~80 bytes (64-byte member + 8-byte score + 8-byte overhead). An inbox of 800 entries costs ~64KB. For 200M users: 200M × 64KB = 12.8TB. This is the primary cost driver and the reason inboxes must be bounded. Without the bound, memory grows linearly with post volume × follower count, which is unsustainable.
Deep scroll (page 2+) works by continuing the ZREVRANGE with an offset. Beyond the inbox's 800 entries, the feed service falls back to pulling from author timelines — switching from push to pull for historical content. This is acceptable because deep scroll is rare (<5% of feed reads go past page 1) and latency expectations are relaxed.
TTL is an additional lever: set a 7-day TTL on inbox keys so that users who haven't opened the app in a week don't consume memory. When they return, the feed service detects the empty/expired inbox and does a full pull-based reconstruction, then resumes normal push delivery.
Merge-K streams at read time
In the hybrid model, feed reads merge two data sources: the pre-materialised inbox (from push) and K celebrity timelines (from pull). The merge is a K-way merge of sorted streams, where K = 1 (inbox) + number of celebrities the user follows.
The algorithm uses a min-heap of size K+1. Initialise the heap with the most recent entry from each stream. To produce the next feed item, extract the max from the heap, then insert the next entry from that item's source stream. Repeat until you have a full page (typically 50 items).
The reader's inbox plus K celebrity timelines are merged via a min-heap, then ranked and filtered.
Time complexity: O(N log K) where N is the page size and K is the number of streams. For a typical user following 10 celebrities: O(50 × log 11) ≈ 170 comparisons — trivial.
The expensive part isn't the merge — it's fetching the celebrity timelines. Each timeline read is a ZREVRANGE against a Redis key or a database query. With 10 celebrity pulls, and each taking 2-5ms (cached) or 20-50ms (cache miss), the total pull latency is 2-50ms depending on cache hit rates. Celebrity timelines are hot (millions of readers access them) so cache hit rates exceed 99% in practice.
Caching the merged result is tempting but tricky. The merged feed is personalised (different users follow different celebrities), so you can't share it across users. You can cache it per-user with a short TTL (30-60 seconds) to handle rapid pull-to-refresh. The cache key includes the user id and a version number that increments on each inbox write, ensuring the user sees new posts within one refresh cycle.
For the ranking pipeline, the merge step produces a candidate set of ~100-200 items (50 from inbox + 10-15 × 10 from celebrity timelines). This candidate set is then scored and re-ranked by the ranking model, which may reorder items based on engagement prediction, recency decay, and interest-graph affinity. The merge-K step is therefore a pre-filter, not the final ordering.
Spread fan-out: taming mid-tier high-follower authors
The hybrid boundary routes the most expensive authors to pull. But what about mid-tier high-follower authors who are still on the push path? A single post from a large account can generate hundreds of thousands of inbox writes. If several such authors post simultaneously, the fan-out queue can still spike.
Spread fan-out addresses this by chunking the follower list and scheduling chunks over time. Instead of enqueuing a single "fan out to 500K followers" task, the fan-out scheduler splits the follower list into chunks of 50K and schedules them with staggered delays: chunk 1 at t=0, chunk 2 at t=10s, chunk 3 at t=20s, and so on. The fan-out workers drain at a steady rate, never spiking above their provisioned throughput.
Chunk the follower list and drain at a steady rate to avoid thundering-herd writes.
The trade-off is fan-out latency. A large post now reaches followers over a controlled window instead of all at once. For many feeds this is acceptable; users rarely notice whether a post appears a few seconds later. The most engaged or recently active followers can be placed in the first chunks if freshness matters.
Implementation details: the fan-out scheduler stores chunk metadata in a delayed message queue (e.g., SQS with delay seconds, or a Redis sorted set used as a delay queue). Each chunk message contains the post id, the follower list offset, and the chunk size. Workers process chunks idempotently — if a chunk is retried, the ZADD operations are naturally idempotent (same member + same score = no-op).
Rate limiting fan-out workers is the complementary control. Each worker has a configurable writes-per-second limit (e.g., 50K ZADD/sec). If the queue depth exceeds a threshold, the scheduler automatically increases the inter-chunk delay for new posts, providing backpressure without dropping fan-out tasks. The system self-regulates: during quiet periods, fan-out is near-instant; during peak load, it spreads gracefully.
Handling follow/unfollow churn
The follow graph is not static. Users follow and unfollow accounts constantly. Each event has implications for the fan-out system that are easy to overlook in a design interview.
On follow: When user A follows user B, A's inbox should contain B's recent posts. Two approaches: (1) Backfill — a background worker fetches B's last N posts and ZADDs them into A's inbox. This provides immediate gratification but adds write load proportional to the follow rate. (2) Lazy — don't backfill; A will start seeing B's posts on the next fan-out. This is simpler but means A's feed looks sparse right after following someone. Most production systems use backfill with a small N (e.g., 10 posts).
On follow: backfill inbox. On unfollow: lazy filter at read, nightly sweep cleans stale entries.
On unfollow: When user A unfollows user B, B's posts should disappear from A's feed. Two approaches: (1) Eager deletion — scan A's inbox and remove all post ids authored by B. This requires a secondary index (post_id → author_id) and is expensive for large inboxes. (2) Lazy filter — leave B's posts in A's inbox but filter them out at read time by checking the follow graph. This is cheaper but means stale entries consume inbox space.
The production answer is lazy filter + nightly sweep. At read time, the feed service checks each candidate post's author against the reader's current follow set (cached) and filters out unfollowed authors. A nightly batch job scans inboxes and removes entries from unfollowed authors, reclaiming space. The lazy filter ensures correctness immediately; the nightly sweep ensures space efficiency eventually.
Edge case: A follows B, B posts, A unfollows B, A re-follows B. With lazy filter, B's posts are filtered out during the unfollow window and reappear on re-follow (because the filter checks current follow state). With eager deletion, the posts are gone and would need backfill on re-follow. The lazy approach handles this naturally.
Another subtlety: if B is on the pull path, follow/unfollow doesn't affect A's inbox at all (B's posts were never pushed). The follow graph update only affects the pull list used at read time — the feed service checks which pull-path authors to read from based on A's current follow set.
The ranking pipeline: from candidates to feed
A modern feed is not chronological — it is ranked. The ranking pipeline transforms a raw candidate set into a personalised, engaging feed. In the hybrid fan-out model, the pipeline sits between the merge step and the final response.
Candidate generation from inbox + celebrity pull feeds into feature enrichment, ML scoring, and diversity filtering.
Stage 1: Candidate generation. The merge-K step produces 100-200 candidates from the inbox (push) and celebrity timelines (pull). This is the initial candidate pool.
Stage 2: Feature enrichment. Each candidate is annotated with features from the feature store: post age, author engagement rate, media type, number of likes/retweets in the first 5 minutes (early engagement signal), the reader's historical interaction rate with this author, topic embeddings, and recency decay. The feature store is typically a low-latency key-value store (Redis or a purpose-built feature service) pre-computed by offline pipelines.
Stage 3: Scoring. A lightweight ML model (typically a gradient-boosted tree or a small neural network) scores each candidate. The model predicts P(engagement) — the probability the reader will like, comment, or share the post. Training data comes from historical engagement logs. Inference must complete in <10ms for the full candidate set, which constrains model complexity.
Stage 4: Diversity filtering. A ranked list of the top 50 by score goes through a diversity filter that ensures no single author dominates the feed (e.g., max 3 posts per author in a single page), mixes content types (text, image, video), and injects exploration candidates (posts from authors the user hasn't interacted with recently) to avoid filter bubbles.
Stage 5: Materialisation. The final ranked page is cached per-user with a 60-second TTL. Subsequent feed reads within the TTL window return the cached result. On pull-to-refresh, the cache is invalidated and the pipeline runs again.
The ranking pipeline adds 20-50ms to feed read latency (on top of the ~10ms for inbox read + merge). The total feed read latency budget is typically 50-100ms. Keeping the pipeline within budget requires aggressive pre-computation (features in the feature store, not computed at serving time), model simplicity (inference must be <10ms), and parallelism (feature fetch and model inference can overlap for different candidates).
Decision levers
Push/pull boundary model
Do not tune from follower count alone. Use a compact cost model: push cost ≈ active_followers × author_post_rate, constrained by fan-out write capacity and freshness SLO; pull cost ≈ follower_feed_reads × merge_cost, constrained by the p99 number of pull-path authors each reader follows. Lower the push boundary when the fan-out queue backs up or author post rate rises. Raise it when read-time pull latency grows. Use hysteresis and cooldowns so authors do not oscillate between paths.
Inbox size bound
Bounded inboxes (typically 800 entries) control Redis memory usage. The bound determines how far back a user can scroll in their pre-materialised feed before the system falls back to pull-based deep scroll. 800 entries × ~80 bytes = 64KB per user. For 200M users: ~12.8TB. Increasing the bound improves deep-scroll UX but linearly increases memory cost. Decreasing it saves memory but triggers more fallback reads.
Spread fan-out vs instant fan-out
For mid-tier high-follower authors who remain on the push path, choose between instant fan-out (faster delivery, risk of queue saturation) and spread fan-out (chunked over time, steady queue utilisation). Spread is safer but adds latency. The choice depends on the product's freshness SLA and current queue depth: if users expect near-immediate delivery, instant with rate-limited workers; if a short delay is acceptable, spread.
Ranking model complexity
The ranking model sits on the read path and must score 100-200 candidates in <20ms. Gradient-boosted trees (XGBoost/LightGBM) inference in <5ms for 200 items. Small neural nets add expressiveness but need GPU inference or heavy optimisation. The ranking budget is the difference between the feed read SLA (e.g., 100ms) and the merge + data fetch time (~50ms). Overshoot the budget and feed reads lag; undershoot and engagement drops.
Multi-region vs single-region
Single-region is simpler but adds cross-region latency for global users (100-200ms for US-to-Asia). Multi-region with regional inboxes eliminates this but requires cross-region fan-out routing, regional Redis clusters, and celebrity timeline replication. The break-even point is typically ~30% of DAU outside the primary region. Below that, CDN-cached feed responses may suffice; above it, regional inboxes are worth the operational cost.
Failure modes
Without a push/pull boundary, one high-follower post can generate tens of millions of fan-out tasks. The queue saturates, workers fall behind, and unrelated users' feeds go stale. Fix: route that author through pull or spread fan-out based on the cost model. The post goes to the author's timeline; followers pull it at read time or receive it through controlled chunks.
In pure pull, a user following 5,000 accounts triggers 5,000 timeline reads per feed refresh. Even with 99% cache hit rate, that's 50 cache misses at 20ms each = 1 second of tail latency. Fix: hybrid with push for normal-follower-count authors. The inbox handles 95% of the content; pull handles only the small number of celebrities.
Without bounded inboxes, memory grows linearly with post volume × follower count. At Twitter scale, unbounded inboxes would consume petabytes. Fix: ZREMRANGEBYRANK to trim inboxes to 800 entries. Deep scroll falls back to author timeline reads. Add TTL (7 days) to expire inactive users' inboxes.
An author hovering around the push/pull boundary gets reclassified on every measurement window. Some posts fan out, others do not, creating inconsistent delivery and operational noise. Fix: hysteresis — require sustained cost above the boundary before promotion to pull, sustained lower cost before demotion to push, and a cooldown before any reversal.
When a user unfollows an account, their inbox still contains that account's posts. Without cleanup, the feed shows content from unfollowed users. Fix: lazy filter at read time (check each candidate's author against current follow set) ensures immediate correctness. Nightly sweep removes stale entries for memory reclamation.
If the fan-out queue (SQS, Kafka, etc.) goes down, no new posts reach inboxes. Feeds go stale within minutes. Fix: replicated queue with at-least-once delivery. Fan-out workers are idempotent (ZADD with same member+score is a no-op). Dead-letter queue for poison messages. Circuit breaker on the fan-out path with fallback to pull for all users during outage.
Case studies
Twitter / X
The canonical push-pull hybrid
Raffi Krikorian's QCon talk "Timelines at Scale" is the classic public reference for Twitter-style timelines. The talk summary describes handling thousands of events per second: tweets, social-graph mutations, and direct messages.
The architectural lesson is the push/pull split. For most authors, Twitter could materialise home timelines ahead of reads by copying tweet ids into follower timeline storage. For high-follower accounts, write-time fan-out becomes the expensive path; their tweets are better stored once and merged into a reader's home timeline when that reader asks for it.
The talk is useful because it names the real constraint: power-law reach means the average author and the outlier author are different workloads. It does not give a universal threshold to copy. Derive your own boundary from active followers, post rate, feed-read rate, fan-out capacity, and freshness SLO.
Key technical details to carry forward: materialised inboxes must be bounded, writes must be idempotent, high-follower content needs a pull path, and feed reads need a merge/filter/rank step that can enforce current blocks, mutes, privacy, deletes, and moderation state.
Ranking changes what fan-out is for
Instagram is useful here as a public example of the feed moving from pure recency toward ranking. Its official ranking explainer says Feed and Stories start from recent posts by accounts you follow, then use signals about the post, the person who posted, your activity and your interaction history to predict which posts matter most.
That explainer is about ranking; Instagram has not published its fan-out thresholds or storage internals in it. It still teaches the key interaction: fan-out produces candidates, ranking decides order. A pre-materialised inbox is valuable because it gives the ranking system a fast candidate set; it is not the final feed order once personalisation is involved.
The interview move is to keep these concerns separate. First choose how candidates arrive cheaply (push, pull, or hybrid). Then describe how ranking filters and orders them using product-specific features.
Professional activity feed (illustrative)
Professional activity feed with bounded fan-out
Illustrative example: a professional network mixes structurally different content types — posts, job changes, endorsements, work anniversaries, article shares — each with different fan-out characteristics.
If the primary social graph is mutual and product-bounded, push can dominate for connection-sourced content because every write has a known ceiling. If the same product also supports unbounded follows for companies, creators, or topics, those sources may need the pull path because their reach is no longer bounded by the connection graph.
The ranking pipeline then handles heterogeneous candidates. A job-change notification, a long-form article, and a short post have different engagement signals and different time decay. Fan-out decides how candidates arrive; ranking decides whether they are relevant enough to show.
This is a useful interview pattern because it prevents over-generalising from open social graphs. Sometimes the product model itself gives you a safe push boundary. Sometimes a second relationship type reintroduces the high-follower problem.
Decision table
Fan-out approach comparison
| Approach | Best for | Write cost | Read cost | Complexity |
|---|---|---|---|---|
| Push (fan-out on write) | Bounded active followers | O(followers) per post | O(1) — read inbox | Low — single path |
| Pull (fan-out on read) | High-follower or low-read sources | O(1) — append to timeline | O(following) — merge K streams | Low — single path |
| Hybrid (push + pull) | Skewed large-scale feeds | O(followers) on push path; O(1) on pull path | O(1) inbox + O(K) pulled authors | High — two paths + merge + boundary tuning |
| Activity stream (bounded push) | System events, team feeds | O(group size) — bounded | O(1) — read inbox | Low — push only, no threshold needed |
- K = number of pull-path authors considered for this reader
- Boundary = active followers × post rate vs follower feed reads × merge cost, constrained by capacity and freshness
Drills
Why can't you use pure fan-out on write for a Twitter-scale system?Reveal
Because the follower distribution is power-law. A high-follower author can generate tens of millions of inbox writes from one post, saturating the fan-out queue and delaying unrelated delivery. The write amplification (eligible active followers × posts) is unsustainable for the top tail, even though push works well for the median author.
What data structure should you use for the inbox and why?Reveal
Redis sorted set (ZSET). Members are post ids, scores are timestamps. ZADD is O(log N) for insertion. ZREVRANGE returns the top-K by recency in O(K + log N). ZREMRANGEBYRANK trims to bounded size atomically. The ZSET provides insertion, retrieval, and trimming in a single data structure with sub-millisecond latency.
How does the merge-K step work at read time in the hybrid model?Reveal
Initialise a max-heap with the most recent entry from each source (inbox + K celebrity timelines). Extract the max, add it to the result, and push the next entry from that source. Repeat for page_size items. Time complexity: O(N log K) where N = page size and K = number of sources. For typical values (N=50, K=15), this is ~200 comparisons — trivial.
What happens when a normal user goes viral and crosses the push/pull boundary?Reveal
The routing model detects sustained cost above the safe push boundary. Future posts route through pull or controlled spread fan-out. In-flight fan-out for the current post can continue idempotently, but new chunks obey backpressure. Hysteresis prevents oscillation: promote only after sustained high cost, demote only after sustained lower cost, with a cooldown period.
Why use lazy filter + nightly sweep for unfollows instead of eager deletion?Reveal
Eager deletion requires scanning the entire inbox to find posts by the unfollowed author — O(N) per unfollow, expensive at scale. Lazy filter checks the follow graph at read time (cached, O(1) per candidate) and hides unfollowed content immediately. The nightly sweep batch-processes stale entries for memory reclamation. This trades temporary space overhead for correctness + low write-path cost.
How would you handle fan-out in a multi-region deployment?Reveal
Partition followers by region in the follow graph. When a post arrives, the region router dispatches fan-out tasks to regional workers, each writing to regional inbox Redis clusters. Feed reads are served from the local region. Pull-path author timelines may cross regions — cache them locally with a short TTL or replicate hot timelines. This keeps the common feed read path local while making freshness tradeoffs explicit.
When to reach for this