Fan-out-on-write vs pull, hybrid celebrity tier, timeline materialization, hydration caches, ranking altitude, and cost-per-render levers.
Lesson 5 · Worked design
Newsfeed at scale
Fan-out is the product
“Design Twitter/Instagram feed” is really: who pays for fan-out — writers or readers? Pure fan-out-on-write (push) makes reads fast and writes expensive. Pure fan-out-on-read (pull) makes writes cheap and reads heavy. Production systems land on a hybrid, usually with a celebrity tier. Highscalability’s archival Twitter numbers (reported historical): ~150M MAU, ~300K QPS timeline generation class, ~4K tweets/s peak, hybrid fanout after pure push broke on celebrities. Facebook TAO (classic paper/blog era) is the graph-read anchor. This lesson makes the tradeoff numerical and operational. Measure fan-out lag as a product metric: p50/p99 time from post ACK to 95% of eligible timelines updated. If lag spikes, user-visible freshness degrades even when HTTP 200s look fine. Page on lag, not only on error rate.
Scope MVP: follow graph, post text/image, home timeline chronological or lightly ranked, like counts eventual. Out of scope for first pass: full ads auction, live video, DMs (lesson 6). Example NFRs: 100–500M DAU class, heavy read:write on home vs posts, p99 home feed < 200–300ms, eventual consistency OK for others’ likes, freshness within a few seconds for normal users.
Capacity that forces fan-out choice
python Feed integrity and safety sit next to fan-out: mute/block filters apply at hydrate or candidate stage; legal takedowns must hide content even if timeline IDs still point at them (SoT filter again); rate limits stop follow/unfollow thrash and post spam from poisoning fan-out workers. If the interviewer is from a trust-heavy company, spend two minutes here after hybrid fan-out is solid. If they are classic infra, keep it to one sentence and return to lag SLOs and cost per home render.
1# Illustrative interview math (heuristics — label as such)2dau = 100_000_0003posts_per_user_day = 0.24post_write_qps = dau * posts_per_user_day / 86_400 # ~230 QPS posts avg5peak_posts = post_write_qps * 10 # event spikes67avg_followers = 2008# Pure push: each post fans out to ~200 timeline inserts9fanout_write_qps = post_write_qps * avg_followers # ~46K QPS timeline writes1011# Home reads: 50% open feed × 10 reads/day12read_qps = dau * 0.5 * 10 / 86_400 # ~6K avg → ~30K peak1314# Celebrity: 30M followers → one post = 30M writes — system death if naive push15# DECISIONS FORCED:16# 1) fan-out service decoupled via Kafka from post write ACK17# 2) celebrity detection on write → pull/merge path18# 3) timeline as IDs not full bodies; hydrate from post cache19# 4) inactive-user skip or lazy fan-out to cut write amp
Fan-out-on-write (push). On post create: lookup followers → enqueue fan-out jobs → insert post_id into each follower’s timeline store (Redis sorted set, Cassandra wide row, etc.). Home read: read precomputed timeline slice + hydrate posts from post service/cache. Pros: fast reads, simple home query. Cons: write amplification, slow fan-out for large graphs, wasted work for inactive users. Cap timeline length (e.g. last ~800 IDs) so storage does not grow without bound per user. Fan-out-on-read (pull). On post create: write to poster’s outbox only. Home read: fetch following list → query recent posts from each (or from a smaller set) → merge-sort by time → rank. Pros: cheap writes. Cons: heavy reads, fan-in merge latency, harder caching of personalized home. Early Instagram-style stories in prep materials often start here, then grow hybrid under load. Hybrid (production default story).Normal users: push to followers’ timelines. Celebrities / high-follower accounts: do not push to all; leave posts in celebrity outbox; at read time, merge celebrity posts into the precomputed timeline. Threshold (e.g. 10k–100k followers) is a product/infra knob. Inactive users: skip push or push lazily on login (optimization). LinkedIn’s feed architecture posts (log-derived feeds / Venice-class systems in eng blogs) are useful anchors for “derived feed data plane,” not only Twitter lore.
python
1def on_post_created(post):2 followers = graph.followers(post.author_id)3 if is_celebrity(post.author_id):4 outbox.add(post.author_id, post.id)5 return6 # chunk fan-out to workers7 for chunk in chunks(followers, 1000):8 fanout_q.push(FanoutJob(post.id, chunk))910def home_timeline(user_id, cursor):11 ids = timeline_store.page(user_id, cursor, n=50) # precomputed12 celeb_ids = merge_celeb_outboxes(user_id, cursor) # pull path13 ids = merge_by_time(ids, celeb_ids)[:50]14 return hydrate(ids) # post cache / multi-get1516FANOUT MODEL TRADEOFF17 Push (write) Pull (read) Hybrid18Read lat very fast heavy merge mixed19Write cost O(followers) O(1) mixed20Storage O(users×recent) O(posts) mixed21Freshness bounded by fanout live mixed22Pick <~10k followers celebs / sparse production social
Data model sketch
posts(post_id, author_id, text, media_refs, ts, …). follows(follower_id, followee_id). timeline(user_id, ts, post_id) — the materialized home. outbox(author_id, ts, post_id) for celebrities/pull. engagement(post_id, user_id, action, ts) eventual. Counters (likes) often separate with cache. On follow: optional async backfill of followee’s last K posts into follower timeline.
Hydration & cache tiers
Timeline stores IDs only (cheap). Hydration multi-gets posts from a post cache (Redis) backed by DB. User cards and media URLs from other caches/CDN. Cache hierarchy: CDN for media → edge/app caches for hot posts → timeline store → origin DBs. Stampede control on viral posts. Deletes: mark deleted on post service (source of truth); hydrate filters deleted; optional async scrub of timeline IDs — never synchronously delete from millions of timelines before ACK.
Ranking (altitude for classic SWE)
Chronological MVP is fine and often preferred in classic SWE rounds. Ranked feeds: candidate generation (timeline IDs + suggestions) → feature fetch → rank model → filters (blocks, mutes). Pinterest/Instagram eng blogs (2024–2025 class) show multi-stage ranking; keep ML shallow here — feature/store latency budgets matter. For classic SWE rounds, say “ranking service with heuristics first; ML later” unless the company is ML-native (then pointer: companion ml-system-design track).
Cost per feed render & failure modes
Cost drivers: fan-out writes, hydrate multi-gets, rank feature fetches, media bandwidth. Optimizations: trim timeline length, sample inactive, compress IDs, batch hydrate, CDN media. Failures: fan-out lag → delayed home visibility; partial fan-out → idempotent retry by (post_id, follower_id); hydration miss storm → soft TTL + singleflight; graph down → cached following set; ranking down → chronological fallback. Storage reality check: 500M users × ~1KB × hundreds of IDs does not all fit in one Redis box — cluster + SSD-backed stores appear in real systems.
If write amplification explodes, you cut push for the heavy tail (celebrities/inactive) — you do not “add more Kafka partitions” and hope.
End-to-end interview narration
1) Scope feed + follow. 2) Estimate post QPS × avg followers → show push cost. 3) Introduce celebrity hybrid. 4) API: post, follow, getHome(cursor). 5) Diagram: post service, graph, fan-out workers, timeline store, hydrate cache. 6) Deep dive fan-out job design + idempotency. 7) Ranking/cache/cost. 8) Failure modes. 9) What changes at 10× DAU.
Interview answers — newsfeed
01Celebrity 50M followers posts? → no full push; outbox + pull merge; rate-limit fanout workers.
02Why not only reverse-chrono forever? → product ranking; still keep chrono fallback.
1POST /posts {text, media_refs[]} → {post_id}2POST /follow {target_user_id}3GET /feed?cursor&limit=50 → {items[], next_cursor}4POST /like {post_id} # eventual counter56# Cursor = (ts, post_id) for stable pagination under concurrent inserts7# Never use OFFSET for hot home feeds at scale
Fan-out job design (the real deep dive)
Job payload: {post_id, follower_chunk[], attempt}. Workers are idempotent on (post_id, follower_id). Chunk size (e.g. 500–2000) balances Kafka message size vs parallelism. Lag SLO: p99 fan-out complete under N seconds for non-celebs. DLQ for poison graph edges. Backpressure: if timeline store saturates, slow consumers and surface a freshness metric on home. Celebrities short-circuit before enqueue.
code
1FANOUT WORKER CHECKLIST2[ ] Idempotent writes to timeline3[ ] Chunked follower lists4[ ] Celebrity short-circuit before enqueue5[ ] Inactive-user skip policy6[ ] Lag metric + alert7[ ] DLQ + replay8[ ] Partial failure retry without duplicating entire fanout
Graph service is its own scale problem: follower lists for celebs are huge — store differently (do not load 30M IDs into one job). Use sharded edge tables, and for celebs never materialize full follower fanout. Home for a user who follows many celebs: merge bounded outboxes (top K by recency) not unbounded.
What changes at 10× DAU
At 10× you revisit: timeline storage medium (memory vs SSD-backed), more aggressive inactive skip, higher celebrity threshold (or dynamic), multi-region home caches, and stricter cost budgets per home render. Ranking may move more work offline (precompute candidates). You do not “just add Kafka partitions” as the only lever — that is the mid-level trap from earlier misconceptions.
Counter and engagement path
Likes/views are classic eventual counters: write to a write-optimized path (Redis INCR + async durable), periodically flush, accept brief disagreement across devices. Do not put counter updates on the critical path of post create or home hydrate. For integrity (no double-like), keep a user×post existence key with TTL or durable set.
Checkpoint
A user with 20M followers posts. Pure fan-out-on-write would…
ABe fine because writes are always cheaper than readsBCreate catastrophic write amplification; hybrid should pull/merge this author’s posts at read timeCOnly require a bigger load balancer
Home timeline p99 spikes when hydrate multi-get misses. Best first fix?
AStore full post bodies duplicated in every follower timeline row foreverBStrengthen post cache (higher hit rate, single-flight, replica-friendly multi-get) and keep timeline as IDsCDisable the feed
Write amplification explodes after a growth spurt. What do you cut first?
APush fan-out to inactive users and raise celebrity threshold handling; keep read p99 via hybridBTurn all reads into full graph scans with no cacheCDrop durability on posts
AOnly delete from author’s outbox; leave stale IDs forever in all timelines without filteringBMark deleted on post service (source of truth); hydrate filters deleted; optional async scrub of timeline IDsCSynchronous delete from millions of timelines before ACK to user
Frame: “Home feed for a Twitter-class follow graph — MVP chronological + mute/block, ranking later.” Estimate: post QPS × avg followers = push cost; call out 30M-follower celebrity as existence proof against pure push. Design: post service, graph, fanout workers, timeline store of IDs, hydrate cache, celebrity outbox merge. Deep dive: fanout job idempotency and lag SLO. Close: delete/filter, inactive skip, cost per home render, ranking altitude pointer to ML track. Historical anchors: Twitter hybrid fanout (archival reported), TAO-style graph reads, modern multi-stage ranking blogs — always labeled.
code
1FEED DECISION LOG2Signal Forces3----------------------------- ----------------------------------4Read:write >> 1 materialize or heavy cache on read5Power-law followers hybrid / celebrity pull6Inactive majority skip or lazy fanout7Viral post hydrate misses single-flight post cache8Edits/deletes SoT filter on hydrate9Ranking uncertainty chrono MVP + filters first1010× DAU storage medium + stricter skip1112Push cost example: 230 post QPS × 200 followers ≈ 46k timeline writes/s avg (heuristic).
Common trap answers to reverse: “Kafka fans out for free” (no), “store full post in every timeline forever” (edit/delete/storage nightmare), “global strong consistency for likes” (unnecessary), “ML ranker before fanout correctness” (wrong altitude). If you only fix fanout correctness and hydrate caching under time pressure, you still pass most classic SWE feed rounds.
Newsfeed math + decision card
code
1HEURISTIC WORKED EXAMPLE2DAU 1e8; posts/user/day 0.2 → post QPS ~230 avg3avg followers 200 → pure push ~46k timeline writes/s4home: 50% open * 10 reads → ~6k RPS avg → ~30k peak5celeb 3e7 followers → one post 3e7 writes → FORBIDDEN pure push6THEREFORE: hybrid; fanout service on Kafka; timeline IDs; hydrate cache78MODEL9posts(SoT) | follows | timeline(user,ts,post_id) | outbox(author,ts,post_id)10API: POST /posts, POST /follow, GET /feed?cursor1112FANOUT WORKER13chunk followers; idempotent (post_id, follower_id); lag SLO; DLQ14skip inactive; short-circuit celebrities before enqueue1516READ PATH17page timeline IDs → merge celeb outboxes → hydrate posts → filter deleted/mute18ranking: chrono MVP; ML later / companion track1920COST LEVERS21trim timeline length; skip inactive; CDN media; batch hydrate22page on fanout lag + home p99 + hydrate miss rate2324TRAPS25Kafka ≠ free fanout; full body denorm everywhere; sync delete N timelines
Can you defend hybrid fan-out with numbers, sketch timeline+outbox models, and name cost/failure levers?
New to itGetting thereConfident
Takeaways
Push vs pull is a cost shift; hybrid is the grown-up answer.
Celebrities break pure push — call it out with math.
Timeline stores IDs; hydrate from post cache; media on CDN.
Fan-out jobs must be idempotent; deletes filter on hydrate.
Ranking altitude: chronological → heuristics → ML (companion track).
Next: realtime chat — WebSockets, presence, ordering, multi-device, persistence anchors from Discord’s engineering blog (as reported).