Facebook Post Search
Sharpened prompt. Design search over 10 trillion posts where a post is findable within seconds of being written, every result respects the viewer's real-time privacy relationship to the author, results are personalised, and the P99 stays under 300ms — knowing that the same query from two different people must return different results.
Web search caches aggressively because everyone querying "weather" gets the same answer. Social search cannot, because privacy and personalisation make every result set viewer-specific. That single difference generates most of the hard parts here.
1. Problem framing
Functional requirements
- Full-text search over posts, with filters (author, date, group, media type).
- Results respect privacy: public, friends, friends-of-friends, custom lists, blocks.
- Personalised ranking: your friends' posts rank above strangers'.
- Freshness: a post written 5 seconds ago is findable.
Non-functional requirements
| Property | Target | Consequence |
|---|---|---|
| Query latency | P99 < 300ms | Scatter-gather with a hard deadline over many shards |
| Index freshness | < 10s | A separate real-time index tier; you cannot rebuild segments that fast |
| Privacy correctness | 100%, no leaks | Enforced at retrieval, never as a post-filter on a truncated list |
| Mutation rate | 100k writes/sec including edits and privacy changes | The index is as write-heavy as the primary store |
Back-of-the-envelope
Corpus: 10T posts × 200 B of indexable text = 2 PB of raw text
Index: postings ≈ 30–50% of text volume -> ~700 TB - 1 PB
Shards: at 1 TB/shard -> ~1,000 shards, replicated 3× = 3,000 index nodes
Mutations: 100k/sec (new posts, edits, deletes, privacy changes)
Queries: 100k QPS, each fanning out to N shards
-> at 1,000 shards that is 10^8 shard-queries/sec, which is impossible.
Query routing MUST prune shards. This is the key structural constraint.
Recency: >90% of clicked results are from the last 30 days -> tier by time.
The shard-fan-out arithmetic is the number to lead with: it rules out "query every shard" and forces time-tiering and routing, which is the backbone of the design.
2. High-level architecture
Two structural ideas: a real-time tier separate from immutable segments (because Lucene-style segments cannot be updated in place at 100k/sec), and privacy enforced inside the leaf, before top-k truncation.
3. Component inventory
| Component | Concrete choice | Why this one |
|---|---|---|
| Inverted index | Lucene-based (Elasticsearch) or a custom engine (Unicorn/Earlybird lineage) | Skip lists, block-max WAND, and BKD trees for numeric/geo are decades of tuning you should not rewrite |
| Real-time tier | In-memory forward + inverted index, ring-buffered | Insert-cheap, no segment merges; discards after 24h |
| Privacy sets | Roaring bitmaps in a visibility service | Set intersection is the natural operation for "who can see this" |
| Sharding | Document-sharded (by post ID), time-tiered | See 4.5 — term sharding fails on multi-term queries |
| Merger | Bounded min-heap k-way merge with a deadline | Same scatter-gather discipline as the news feed |
| Ranking | Two-stage: per-shard cheap, global heavy | Only ~1,000 docs reach the expensive model |
4. The toughest parts
4.1 Indexing at 100,000 mutations per second
Why it's hard. Inverted indexes are optimised for reads: posting lists are sorted, compressed, and stored in immutable segments. Inserting one document means either rewriting a segment (unthinkable at this rate) or creating a tiny new one — and a million tiny segments destroys query performance because every query must visit every segment. Meanwhile a deletion cannot even be applied in place; it is recorded as a tombstone and only physically removed on merge.
Solution — a two-tier index with different data structures for different ages.
Real-time tier (last 24h, RAM):
- Skip-list or hash-based postings, mutable, insert in microseconds
- Forward index too (doc -> terms), because updates need to know what to remove
- Small: 100k/sec × 86,400 = 8.6B docs/day, but sharded across the fleet
- No compression, no merging. It is a buffer, not a database.
Segment tier (older, SSD/disk):
- Immutable Lucene-style segments, heavily compressed, block-max WAND
- Built in batches every few minutes from the same Kafka stream
- Background tiered merge (like an LSM tree) keeps segment count logarithmic
// Query hits both tiers and merges. The real-time tier is small enough that
// scanning it is cheap; the segment tier is fast because it is immutable.
List<Hit> search(Query q, int k) {
var rt = realtimeIndex.search(q, k); // microseconds, last 24h
var seg = segmentIndex.search(q, k); // milliseconds, everything older
// Deduplicate: a doc may be in both if a segment build just completed.
return mergeTopK(rt, seg, k, /*preferNewer=*/true);
}
Handle updates as delete-plus-insert with a version number: the tombstone is a bitset of deleted doc IDs consulted at query time (liveDocs), and physical removal happens at merge. The forward index in the real-time tier exists precisely so an edit can remove the old terms.
The mental model to state: "This is an LSM tree for text. The real-time tier is the memtable, segments are SSTables, tiered merging is compaction, and tombstones work the same way." That framing lands well because it is exactly right.
4.2 Privacy that cannot leak, at retrieval time
Why it's hard. A post is visible to "friends only." Whether you can see it depends on your live friendship with the author — a relationship that changes constantly and is not a property of the post. The naive design retrieves the top 100 by relevance and then filters by visibility, which is wrong in two distinct ways: it can return fewer than 100 results (you filtered away most of them), and it can leak existence through result counts and timing.
Solution — encode visibility as bitsets and intersect during retrieval, inside the leaf.
Per-post audience bitset (assigned at index time, small integer IDs):
post 12345 -> audience = {PUBLIC=0} public post
post 67890 -> audience = {FRIENDS_OF(user:42)} friends-only
post 11111 -> audience = {LIST(user:42, "Close Friends")} custom list
Per-viewer visibility bitset (assembled per query, cached per session):
viewer 99 -> {PUBLIC, FRIENDS_OF(f) for each friend f of 99, LISTS containing 99}
Retrieval: candidates AND viewer_visibility AND NOT blocked_authors
// Inside the leaf, before top-k. Roaring bitmaps make this a few microseconds.
RoaringBitmap execute(Query q, Viewer v) {
RoaringBitmap docs = intersectPostings(q.terms()); // block-max WAND skipping
docs.and(v.visibilityBitmap()); // privacy: a set intersection
docs.andNot(v.blockedAuthorDocs()); // blocks are subtraction
docs.andNot(segment.deletedDocs()); // tombstones
return docs;
}
The performance question is the viewer bitmap: a user with 5,000 friends implies 5,000 FRIENDS_OF terms. Materialising that per query is too slow, so cache the viewer's visibility bitmap per session (it changes only when they add or remove a friend) and invalidate on graph mutation. A Roaring bitmap covering a 5,000-friend audience is a few tens of KB and intersects at memory bandwidth.
The correctness rule to state plainly: privacy is a retrieval filter, never a post-filter. Filtering after top-k is not merely slower — it silently returns short result sets and it leaks. Interviewers probing this question are usually checking whether you know the difference.
4.3 Caching when every user's results are different
Why it's hard. The single most powerful lever in web search — cache the result set for a popular query — is unavailable. "Concert" from two different people returns different posts, different order, different counts. Result-level cache hit rate is effectively zero. Yet you still need caching, because 100k QPS against 1,000 shards without it is impossible.
Solution — cache the layers that are shared, and only compute the viewer-specific parts per request.
| Layer | Shared across viewers? | Cache |
|---|---|---|
| Posting lists for a term | Yes — identical for everyone | Aggressively, in the OS page cache and a block cache. Highest-value cache in the system |
| Query parse and rewrite | Yes | Small in-process LRU |
| Post content, author, counts | Yes | Hydration cache; a popular post is hydrated once |
| Viewer visibility bitmap | Per viewer, but stable | Per-session cache, invalidated on graph change |
| Intersected result set | No | Not cacheable — this is the per-request work |
| Ranking features (item side) | Yes | Feature store, precomputed at index time |
The reframing: "I can't cache answers, so I cache everything that goes into computing an answer. The per-request work then reduces to a handful of bitmap intersections plus ranking — microseconds of CPU on data that is already resident."
There is one genuine result-cache opportunity worth naming: public-only queries. A query restricted to public content (or issued by a logged-out user) has a viewer-independent result set and can be cached normally. Splitting the public and private paths lets a meaningful share of traffic take the cheap route.
4.4 Ranking a personal corpus
Why it's hard. Relevance in social search is dominated by who wrote it and your relationship to them, not by text similarity. A post from your sibling mentioning "dinner" should outrank a perfect keyword match from a stranger. But social-graph features are viewer-specific, so they cannot be precomputed into the index, and computing them for thousands of candidates per query is too slow.
Solution — split features by where they can live, and by when they can be computed.
def score(post, viewer, query):
# Precomputed at index time, stored in the index. Free at query time.
static = (post.quality_score * 0.15 +
post.engagement_prior * 0.10 +
recency_decay(post.age) * 0.20)
# Computed in the leaf from posting-list data. Cheap.
textual = bm25(query, post) * 0.20
# Viewer-specific — the expensive part. Only for the ~1,000 docs that
# survive per-shard top-k, using bitmap lookups rather than graph queries.
social = (
0.20 * (1.0 if post.author in viewer.friends_bitmap else 0.0) +
0.10 * viewer.affinity.get(post.author, 0.0) + # precomputed affinity table
0.05 * (1.0 if post.group_id in viewer.groups_bitmap else 0.0))
return static + textual + social
Structure the pipeline so social features are applied after per-shard truncation: each of 1,000 shards returns its top 20 by static + textual, the merger takes the global top ~1,000, and only then does the heavy personalised model run. That is 1,000 scorings per query instead of millions.
The affinity table — a precomputed map of viewer → top 500 accounts by interaction — is the trick that makes the social term cheap. It is updated nightly (plus nearline for big changes) and fits in a few KB per user.
4.5 Sharding: by term or by document?
Why it's hard. This is the classic distributed-search fork and it deserves an explicit answer rather than a default.
Term sharding (each shard owns a subset of terms, holding the complete posting list for each). A single-term query touches one shard — wonderfully efficient. But a two-term query requires shipping an entire posting list across the network to intersect it: the posting list for a common word can be hundreds of millions of doc IDs. Worse, term frequency is Zipfian, so the shard owning common words is permanently hot.
Document sharding (each shard owns a subset of documents and a complete index over them). Every query goes to every shard — expensive fan-out — but each shard answers independently with no cross-shard data movement, load is uniform, and adding capacity is trivial.
Solution — document sharding, with time tiering to prune the fan-out.
def route(query) -> list[Shard]:
# Almost all clicked results are recent, so query recent tiers first
# and only descend when they don't satisfy the request.
tiers = [RT_TIER, RECENT_TIER] # last 24h, last 30d
if query.explicit_date_range or query.needs_deep_recall:
tiers.append(ARCHIVE_TIER)
shards = []
for t in tiers:
if query.author_filter: # author-scoped: 1 shard, not 1,000
shards += [t.shard_for_author(query.author_filter)]
elif query.group_filter:
shards += [t.shard_for_group(query.group_filter)]
else:
shards += t.all_shards()
return shards
Three pruning mechanisms turn an impossible fan-out into a feasible one. Time tiering: the recent tier is a small fraction of the corpus and satisfies the overwhelming majority of queries. Entity-aware routing: an author- or group-scoped query maps to a single shard if you also co-locate documents by author within a tier. Early termination: run a first pass over recent shards and only descend into the archive if the result set is too small or the user paginates.
Then add block-max WAND inside each shard so that even within a shard, a query skips blocks of the posting list that cannot contain a top-k document. That is where most of the per-shard speed comes from.
4.6 Privacy changes must propagate faster than posts do
Why it's hard. A user changes a post from Public to Friends-only, or unfriends someone, or blocks an account. Every one of those changes the visibility of potentially millions of documents without touching any document. Re-indexing a user's 10,000 posts because they changed one friendship is unaffordable, and doing it for a block affecting millions of posts is impossible.
Solution — indirection. The index stores audience references, and the viewer's bitmap resolves them at query time.
Post-level change (post 123: public -> friends-only):
Update the post's audience field in the real-time tier. ONE document. Immediate.
Segments: record it in a small mutable overlay consulted at query time until
the next merge folds it in.
Graph-level change (user 42 unfriends user 99):
Update user 99's VISIBILITY BITMAP. ZERO documents touched.
Invalidate the cached bitmap; the next query rebuilds it.
Every one of user 42's friends-only posts becomes invisible to 99 instantly.
Because the index never stores "user 99 can see this," and only stores "this is visible to friends-of-42," a graph change is a single bitmap update on the viewer side. This indirection is the whole reason the scheme scales, and articulating it is the strongest single point you can make on this problem.
Deletion needs stronger guarantees than an overlay: a deleted post must vanish immediately, so maintain a deletion bitset consulted on every query, updated synchronously, and replicated ahead of the merge. Correctness for deletion is a legal requirement (GDPR erasure), not a performance question, so it gets the synchronous path even though everything else is eventually consistent.
4.7 Query understanding: what people actually type
Why it's hard. Social search queries are short, ambiguous, and often navigational. "john" might mean a friend named John, a page, a group, or the word appearing in text. Getting this wrong wastes the entire retrieval budget searching the wrong corpus.
Solution — classify intent and rewrite before retrieval.
def understand(raw: str, viewer) -> Query:
q = normalize(raw) # unicode, case, diacritics
# Entity linking against the viewer's own graph first — the strongest prior.
entities = entity_linker.link(q, scope=viewer.friends | viewer.groups | viewer.pages)
if entities.confidence > 0.8 and entities.covers_whole_query:
return NavigationalQuery(entity=entities.best) # skip text search entirely
return TextQuery(
terms=q.split(),
expansions=synonyms(q) + stems(q), # bounded expansion only
filters=extract_filters(q), # "photos from 2019" -> media+date
intent=intent_classifier.predict(q, viewer), # people | posts | groups | mixed
)
Bound the query expansion aggressively. Synonym expansion multiplies posting-list work, and in a corpus of 10 trillion documents an over-expanded query can turn a 20ms retrieval into a 2-second one. Expand only when the initial result set is too small — an adaptive second pass rather than a speculative first one.
5. What breaks first
| Event | First failure | Mitigation |
|---|---|---|
| Breaking-news query spike | Real-time tier CPU on one topic | Public-query result cache; extra replicas of the RT tier |
| A celebrity blocks many accounts | Bitmap invalidation storm | Bitmaps are per-viewer, so invalidation is bounded by affected viewers |
| Segment merge storm | Query latency doubles during merges | Throttled tiered merging; route reads to replicas not currently merging |
| Deep pagination (page 50) | Cost grows linearly with offset | Cap depth; use search-after cursors rather than offsets |
| Over-expanded rare query | One query consumes a shard's budget | Per-query cost accounting with hard termination |
| Shard loss | Silent recall loss — results are merely worse | Replicate 3×; report partial-result rate as an SLI so silent degradation is visible |
That last row is the subtle one and worth volunteering: in search, a failed shard does not produce an error, it produces slightly worse results that nobody notices. Instrument partial results explicitly or you will ship recall regressions for months.
6. Cheat sheet
- Index: LSM-for-text — mutable in-memory real-time tier (24h) plus immutable merged segments, tombstones for deletes.
- Privacy: audience references in the index, viewer visibility bitmap at query time, intersected before top-k. Never post-filter.
- Graph change = one bitmap update, zero documents re-indexed. That indirection is the core insight.
- Caching: cache posting lists, hydration, and features — not results. Public queries are the one cacheable class.
- Sharding: document-sharded, time-tiered, entity-routed, with block-max WAND inside each shard.
- Ranking: static + textual per shard, social features only on the surviving ~1,000 docs, affinity table precomputed.
- The one-liner: "Web search caches answers; social search can't, because the answer depends on who's asking — so I cache every input to the answer and make privacy a bitmap intersection inside the leaf."