Skip to main content

Facebook News Feed

Sharpened prompt. Design a news feed blending posts from friends, groups, pages, and ads for 2B daily users, ranked by predicted engagement, loading in under 300ms, staying consistent for the user who just posted, and paginating correctly while thousands of new stories arrive mid-scroll.

Where Instagram is a single content type with a follow graph, News Feed is heterogeneous aggregation: many sources with different shapes, merged, deduplicated, ranked, and mixed with paid content — all while the underlying social graph is being mutated constantly.

1. Problem framing​

Functional requirements​

  • Aggregate posts from friends, followed pages, joined groups, and events.
  • Rank by predicted engagement; inject ads at policy-defined slots.
  • Deduplicate: the same link shared by six friends should appear once, aggregated.
  • Stable pagination through an infinite scroll.
  • Read-after-write for the user's own actions.

Non-functional requirements​

PropertyTargetConsequence
Feed loadP99 < 300msParallel leaf fan-out with a strict deadline
ConsistencyRYW for self; seconds of staleness for othersTwo different consistency models in one product
Graph reads~10^9 QPS aggregateA cache tier is the primary datastore in practice
Availability99.99%Every stage must have a degraded fallback

Back-of-the-envelope​

Users: 2B DAU, ~15 feed sessions/day = 30B feed requests/day = 350k/sec avg
Per request: ~500 friends + 100 pages + 50 groups = 650 sources to consider
Graph reads: 350k feed req/s × 650 assoc lookups = 227M graph reads/sec
-> a database cannot do this; a cache with a ~99% hit rate can
Candidates: ~2,000 stories/request considered, 20 returned
Ranking: 350k/s × 2,000 = 700M scorings/sec -> two-stage funnel is mandatory

The 227M graph reads/sec is the number that explains TAO's existence, and quoting it makes the "just use a graph database" answer visibly wrong.

2. High-level architecture​

Aggregator-leaf is the pattern to name. The aggregator knows nothing about content; each leaf owns one source type, applies its own retrieval logic, and returns a small pre-ranked candidate list. Adding "Marketplace listings" to the feed means adding a leaf, not modifying the aggregator — and each leaf scales, fails, and is deadlined independently.

3. Component inventory​

ComponentConcrete choiceWhy this one
Graph cacheTAO (Memcached-based, objects + associations) over sharded MySQL10^9 QPS of point lookups with ~99.8% hit rate; MySQL is the durable backstop
LeavesIndependent services with per-leaf deadlinesFailure isolation: a broken groups leaf degrades the feed rather than breaking it
AggregatorScatter-gather with a hard deadline and partial resultsTail latency is dominated by the slowest leaf, so you must be able to drop it
RankingTwo-stage, embeddings precomputed at write timeSame funnel logic as Instagram
Feature storeIn-memory, per-request batchedFeature fetching, not model inference, is usually the latency bottleneck
DedupSimHash / URL canonicalisation at candidate stageSix friends sharing one article is one story

4. The toughest parts​

4.1 Aggregating heterogeneous sources under a hard deadline​

Why it's hard. A feed request touches five leaf services, each hitting caches and stores. Latency is the maximum of the parallel calls, so your P99 is governed by the worst leaf's P99. With five leaves each at P99 = 150ms, the combined P99 is far worse than 150ms — roughly the 99.8th percentile of a single leaf. Waiting for everything guarantees you miss the budget.

Solution — scatter-gather with a deadline, partial results, and per-leaf circuit breakers.

func (a *Aggregator) Candidates(ctx context.Context, uid int64) []Story {
ctx, cancel := context.WithTimeout(ctx, 120*time.Millisecond) // hard budget
defer cancel()

var mu sync.Mutex
var all []Story
var wg sync.WaitGroup

for _, leaf := range a.leaves {
if leaf.Breaker.IsOpen() { continue } // known-bad leaf: skip entirely
wg.Add(1)
go func(l Leaf) {
defer wg.Done()
// Each leaf gets its own slice of the budget, weighted by importance.
lctx, lcancel := context.WithTimeout(ctx, l.Budget)
defer lcancel()
s, err := l.Fetch(lctx, uid)
if err != nil {
l.Breaker.RecordFailure()
metrics.Inc("leaf.miss", l.Name) // degraded, not failed
return
}
mu.Lock(); all = append(all, s...); mu.Unlock()
}(leaf)
}
wg.Wait()
return all // whatever arrived in time. A feed missing group posts is a feed.
}

The product decision embedded here is worth stating: a feed missing one source is vastly better than no feed. Rank the leaves by contribution and give the friend-posts leaf a larger budget and a retry, while the events leaf is best-effort. Track "degraded feed rate" as a first-class SLI so silent quality loss is visible.

For merging the returned lists, a tournament tree / bounded min-heap k-way merge keeps the top N without sorting everything: O(N log k) with k ≈ 5.

4.2 227 million graph reads per second​

Why it's hard. Every feed request needs the viewer's friend list, each friend's recent posts, the privacy state of each post, the viewer's group memberships, and the like/comment associations for display. These are billions of tiny point lookups per second on a graph with trillions of edges. A relational join across sharded MySQL is out of the question; so, at this scale, is a general-purpose graph database.

Solution — a purpose-built read-through cache with exactly two data types.

Objects: id -> (type, data) e.g. 12345 -> (post, {...})
Associations: (id1, atype, id2) -> (time, data)
e.g. (user:7, FRIEND, user:9) -> (t, {})
(post:12345, LIKED_BY, user:7) -> (t, {})

Four query shapes, and only four:
assoc_get(id1, atype, id2) -> does this edge exist?
assoc_range(id1, atype, pos, n) -> newest N edges of a type (the workhorse)
assoc_count(id1, atype) -> edge count, maintained incrementally
obj_get(id) -> the object

Constraining the API to those four shapes is what makes the cache tractable: every one is a point lookup or a bounded ordered range, both of which cache perfectly. assoc_count being maintained rather than computed is why a post can display "1.2M likes" without counting anything.

Deployment matters as much as the data model: a two-tier cache (a per-rack follower tier in front of a per-region leader tier) absorbs hot objects locally, and only the leader tier talks to MySQL. Writes go through the leader, which invalidates followers.

The hot object problem remains — a viral post's object is read by tens of millions of requests/sec. Solve it the same way as everywhere else in this playbook: near-cache the object in the web tier with a one-second TTL. Consistent hashing distributes keys, not load; only local caching fixes a single hot key.

4.3 Read-after-write across datacenters​

Why it's hard. Writes go to a primary region. A user in Europe posts (write crosses to us-east), then refreshes (read served locally in eu-west from a replica that is 80ms–2s behind). Their post is missing. They post again. Now there are two. Meanwhile you cannot simply route all reads to the primary — that defeats the entire point of regional replicas.

Solution — remote markers plus sticky sessions, scoped to the writing user only.

# On write, the local region records a marker BEFORE forwarding to the primary.
async def write_through(user_id, obj):
await local_cache.set_marker(f"rw:{user_id}", ttl=10) # "this user has a pending write"
await primary_region.write(obj) # cross-region, ~80ms
# Replication will invalidate the local cache when it lands.

# On read, the marker forces the authoritative path for this user only.
async def read(user_id, key):
if await local_cache.has_marker(f"rw:{user_id}"):
return await primary_region.read(key) # slower, correct
return await local_cache.get(key) # fast path, 99.9% of reads

Three properties make this the right answer. It is scoped: only the user who just wrote pays the cross-region cost, for a few seconds. It is self-healing: the marker TTL exceeds the replication P99, so it expires naturally once replication has caught up. And it is honest about the model: everyone else in the world sees the post a second later, which is correct — social feeds are not a ledger.

Also do the cheap thing from Instagram: write the author's own story into their own feed cache synchronously. It removes the most-noticed instance of the problem before any of this machinery is needed.

4.4 Ranking, ads, and the rules that fight each other​

Why it's hard. The ranker maximises predicted engagement. The ads system maximises revenue. Integrity wants to demote borderline content. Product wants source diversity. These objectives conflict, and applying them as sequential filters produces incoherent results — the ranker's top story gets removed by diversity, so slot 1 goes to a mediocre story while a great one sits at slot 4.

Solution — separate scoring from slotting, and treat slotting as a constrained selection problem.

def build_feed(candidates, ads, viewer, n=20):
# 1. Score every organic candidate independently.
for c in candidates:
c.score = model.predict(viewer, c) # p(like)·w1 + p(comment)·w2 + ...

# 2. Apply multiplicative adjustments, never hard filters where a demotion works.
for c in candidates:
c.score *= integrity.multiplier(c) # 0.0 (removed) .. 1.0
c.score *= recency_decay(c.age)

# 3. Slot with constraints — a greedy pass with penalties, not a filter chain.
feed, author_counts, last_type = [], Counter(), None
for slot in range(n):
if slot in AD_SLOTS and ads:
feed.append(ads.pop(0)); continue
best = max(candidates, key=lambda c:
c.score
* (0.5 ** author_counts[c.author]) # diminishing returns per author
* (0.7 if c.type == last_type else 1.0)) # type diversity
feed.append(best); candidates.remove(best)
author_counts[best.author] += 1; last_type = best.type
return feed

The insight to articulate: diversity and integrity are penalties on the objective, not filters after it. A filter chain throws away information; a penalised objective lets a genuinely excellent third post from the same author still win a slot when nothing else comes close.

Ads add a second wrinkle: ad slots are auctioned, and the auction's value must be comparable to organic engagement value. Production systems convert both to a common currency (expected value in a shared unit) so an ad only displaces organic content when it is worth more than what it replaces. That framing — one objective, two sources — is the senior answer.

4.5 Deduplicating the same story from six friends​

Why it's hard. A news article is shared by six friends within an hour. Six near-identical stories dominate the feed. But they are different post objects, by different authors, with different commentary — so ID-based dedup does nothing, and naive URL matching misses because every share carries different tracking parameters.

Solution — canonicalise, fingerprint, cluster, then aggregate into a single story unit.

def dedup_key(post) -> str:
if post.link:
u = urlparse(post.link)
# Strip tracking junk and normalise, or every share looks unique.
q = {k: v for k, v in parse_qsl(u.query)
if not k.startswith(('utm_', 'fbclid', 'gclid', 'ref'))}
return f"link:{u.netloc.removeprefix('www.')}{u.path.rstrip('/')}?{urlencode(sorted(q.items()))}"
if post.attached_media_hash:
return f"media:{post.attached_media_hash}"
# Text: SimHash with a Hamming-distance-3 bucket catches near-duplicate reposts.
return f"text:{simhash_bucket(post.text)}"

def collapse(candidates):
groups = defaultdict(list)
for c in candidates: groups[dedup_key(c)].append(c)
out = []
for _, g in groups.items():
best = max(g, key=lambda c: c.score) # keep the highest-ranked instance
best.aggregated_with = [x.author for x in g if x is not best]
best.score *= 1 + 0.1 * len(g) # social proof: many shares = signal
out.append(best) # renders as "Alice and 5 others shared"
return out

Note the score boost: several friends sharing the same thing is evidence it matters, so collapsing six stories into one should raise that one's rank rather than merely removing five. Turning a deduplication problem into a ranking signal is the move worth showing.

4.6 Pagination that does not repeat or skip​

Why it's hard. Offset pagination (LIMIT 20 OFFSET 40) is broken on a mutating, re-ranked list. Between page 1 and page 2, new stories arrive and scores change, so the user sees duplicates, misses stories, or — with re-ranking — gets a page 2 assembled from an entirely different ordering. This is one of the most common real bugs in feed products and interviewers do ask about it.

Solution — a stateful session cursor that freezes the ranked candidate set.

# Page 1: rank once, persist the ordered ID list, return an opaque cursor.
async def feed_page_one(uid):
ranked = await rank_all(uid) # ~500 story ids, ordered
session = uuid4().hex
await redis.setex(f"feedsession:{uid}:{session}",
1800, msgpack.packb([s.id for s in ranked]))
return Page(items=await hydrate(ranked[:20]),
cursor=encode({"s": session, "off": 20}))

# Page N: slice the frozen list. No re-ranking, no drift, no duplicates.
async def feed_page_n(uid, cursor):
c = decode(cursor)
ids = msgpack.unpackb(await redis.get(f"feedsession:{uid}:{c['s']}") or b"")
if not ids: # session expired
return await feed_page_one(uid) # honest restart beats silent skipping
window = ids[c["off"]:c["off"] + 20]
return Page(items=await hydrate(window), # hydration IS fresh: live counts
cursor=encode({"s": c["s"], "off": c["off"] + 20}))

The crucial split: the ordering is frozen, the content is live. Story IDs come from the frozen list, but likes, comments, and edits are hydrated fresh on every page, so the user never sees stale counts and never sees a duplicate story.

New stories arriving mid-session are surfaced explicitly — a "New posts" pill at the top — rather than being spliced into the middle of an active scroll. That is a product decision that also happens to make the engineering tractable, which is exactly the kind of trade worth pointing out.

4.7 A brand-new user with no graph​

Why it's hard. A user with three friends and no history has almost no candidates and no personalisation signal. An empty feed is the worst possible first experience, and it is self-reinforcing: no engagement means no signal means no better ranking.

Solution — a cold-start ladder that degrades gracefully by available signal.

Signal availableSource of candidates
Nothing but localeGlobally and locally popular, quality-filtered, safe-for-all content
Sign-up interestsPopular content in those topic clusters
A few friendsFriends' posts, plus posts their friends engaged with (2-hop)
Contact/graph importSuggested connections woven into the feed as actionable units
~20 interactionsSwitch to the personalised model; two-tower user embedding is now meaningful

Bootstrap the user embedding from demographic and interest priors rather than zeros, so stage-2 retrieval returns something reasonable immediately. And instrument the transition: measure feed quality separately for users with fewer than 50 interactions, because the aggregate metric will hide a terrible new-user experience behind a good average.

5. What breaks first​

EventFirst failureMitigation
One leaf service degradesAggregator P99 blows the budgetPer-leaf deadlines, circuit breakers, partial results
Viral postHot object in the graph cacheWeb-tier near-cache with a 1s TTL
Cross-region replication lag spikeRYW violations, users see missing postsRemote markers with a TTL above replication P99
Ranking model rollout regressionEngagement drops silentlyShadow scoring + holdout population + automatic rollback on guardrail metrics
Feed session store evictionPagination restarts mid-scroll30-minute TTL, graceful restart, client keeps the scroll position
Graph cache tier lossMySQL sees 227M QPS it cannot serveLoad shedding to a degraded chronological feed; warm the cache before reopening

That last row is the interesting one: the honest answer is that MySQL cannot absorb the cache tier's load, so the fallback is a degraded product (chronological, fewer sources), not a slower version of the same product. Recognising that some failures must be answered with product degradation rather than capacity is a senior instinct.

6. Cheat sheet​

  • Shape: aggregator-leaf. One leaf per source type, independently deadlined and circuit-broken; partial results always beat no results.
  • Graph: TAO-style objects + associations with exactly four query shapes; two-tier cache; assoc_count maintained, never computed.
  • Consistency: remote markers scoped to the writing user, plus a synchronous self-feed write.
  • Ranking: score, then slot. Diversity and integrity are score multipliers, not post-hoc filters. Ads and organic share one currency.
  • Dedup: canonicalise links, SimHash text, collapse into one story and boost it for social proof.
  • Pagination: freeze the ordering in a session cursor, hydrate content fresh.
  • Cold start: a signal ladder from global popularity to personalised, with separately measured new-user quality.
  • The one-liner: "A scatter-gather over independently-owned candidate sources with a hard deadline, ranked once per session and sliced from a frozen list — the interesting parts are the deadline, the graph cache, and the consistency exception for the writer."