Hiring BarSupport

Design a Search Engine — Staff-Level Case Study

Case study57 min read7 diagrams

Technologies referenced in this case study: Elasticsearch · Apache Kafka · PostgreSQL · Redis · Flink & Stream Processing

Related: Database Indexing · Sharding & Partitioning · Scaling Reads · Data Pipeline Patterns · Web Crawler · Database Sharding

How to Use This Case Study#

Organized for interview use first, reference second. Read front-to-back once. Return to individual sections for targeted review.

ModeTimeWhat to Read
Quick Review15 minExecutive Summary → Interview Walkthrough → Fault Lines table → Active Drills 1–3
Targeted Study1–2 hrsExecutive Summary → Walkthrough → Section 3 (Fault Lines) → Section 4 (Failures) → weak-spot Deep Dives
Deep Dive3+ hrsEverything, including the Principal Lens and appendices
What is a Search Engine? — Why interviewers pick this topic

A search engine turns documents into an inverted index (term → list of documents containing it), then answers queries by intersecting posting lists, scoring candidates and returning the top K. Around that core sit an ingestion pipeline (how fast a new or updated document becomes searchable), a serving tier (sharded, replicated, fan-out and merge) and a ranking stack (lexical scoring, learned re-ranking, business rules).

Before vs After — Marketplace search scenario:

With SQL LIKE on the product table:
t=0:       Catalog at 40M products; query "red running shoes"
t=+0s:     WHERE title ILIKE '%red%running%shoes%' → full table scan, 9 s
t=+1wk:    Search traffic 800 QPS; primary DB CPU 100%; checkout slows
t=+2wk:    "running shoe red" returns nothing — word order and stemming unsupported
t=+1mo:    Conversion from search down 18%; a read replica dedicated to LIKE queries

With an inverted index (20 shards × 3 replicas, NRT refresh 1 s):
t=0:       Same query → 3 posting lists intersected, BM25 top-1000, re-ranked to top-50
t=+0s:     p50 25 ms, p99 120 ms at 3,000 QPS; primary DB untouched
t=+1s:     Price change on a product visible in search within ~2 s via CDC
t=+1mo:    Learned ranking lifts search conversion 6%; zero impact on checkout DB

Why interviewers reach for this question: Search combines a read-heavy fan-out serving problem, a write-heavy freshness pipeline and a subjective quality problem (relevance) that can't be unit-tested. It shows whether you can reason about sharding for scatter-gather, tail latency, freshness vs indexing cost, and — the Staff part — how you change the index or the ranking without breaking a system every product surface depends on.

Mechanics Refresher: Index and Retrieval Options
TechniqueHow It WorksProsCons
Inverted index + BM25Term → postings (doc IDs, freqs, positions); score by term frequency, inverse document frequency, length normalizationFast, explainable, exact term matchingMisses synonyms/semantics
Segments (Lucene-style LSM)Immutable segments written on refresh, merged in background; deletes are tombstonesHigh write throughput; lock-free readsRefresh/merge cost; deletes linger until merge
Document-partitioned shardingEach shard indexes a subset of docs; query fans out to all shardsEven writes; shard-local scoringEvery query hits every shard
Term-partitioned shardingEach shard owns a subset of termsQuery touches only shards for its termsHot terms, multi-term intersection across shards, painful writes
Vector (ANN) retrievalEmbeddings + HNSW/IVF indexSemantic matchesMemory-heavy; recall/latency tuning; harder to explain
Hybrid retrievalLexical + vector candidates merged, then re-rankedBest recallTwo indexes to operate

For most production systems: document-partitioned Lucene-family indexes (Elasticsearch/OpenSearch/Solr or Vespa), BM25 for candidate generation, a learned re-ranker on the top few hundred, and vector retrieval added as a second candidate source when semantic recall matters. The engine is almost never the interview question — freshness, fan-out tail latency, reindexing and ranking governance are.


Executive Summary

If you only read one section, read this. Everything in the case study flows from the contrast below.

What This Interview Actually Tests#

Search is not an inverted-index question. Everyone can describe posting lists.

It is a fan-out, freshness and relevance-ownership question that tests:

  • Whether you shard for the query path (scatter-gather tail latency) and not just for data size
  • Whether you state a freshness SLO and pay for it deliberately (refresh interval, indexing pipeline)
  • Whether you can reindex a live index with a schema or analyzer change without downtime
  • Whether you treat ranking as a product surface with offline metrics, online experiments and owners
  • Whether the index is a derived store you can rebuild — and who owns the rebuild

The key insight: The index is a cache of the source of truth, optimized for a query shape. Staff engineers design it to be rebuilt, swapped and ranked independently — and they own the gap between what's in the database and what's findable.

The L5 vs L6 Contrast — Start Here#

BehaviorSenior (L5)Staff (L6)Principal (L7)
First move"Elasticsearch with an inverted index""What corpus, what freshness, what query mix, and how is relevance judged?""Is search one shared platform with per-product ranking, or does each product run its own cluster? Who owns relevance org-wide?"
Sharding"Shard by doc ID, add replicas"Sizes shards (10–50 GB), counts fan-out cost, uses replicas for QPS, routing for tenant-scoped queriesSets cluster topology standards (per-tenant vs shared, cell size) and capacity cost per QPS
Freshness"Near real-time, 1 s refresh"Freshness SLO per field (price 5 s, reviews 1 h); CDC pipeline; bulk vs NRT paths; knows refresh costNegotiates freshness tiers with product owners and prices each tier
Reindexing"Reindex the data"Versioned indexes + alias swap; dual-write/backfill; verify doc counts and query parityMakes index versioning a platform capability; schema-change SLA for 40 teams
Ranking"BM25 then sort by popularity"Multi-stage: retrieval → L1 → learned L2; offline NDCG + online A/B; guardrail metricsRelevance governance: who can change ranking, experiment policy, business-rule budgets
OwnershipSearch team owns everythingSource teams own document contracts; search platform owns cluster + pipeline; product owns ranking goalsRedraws boundaries: relevance as a product org; platform with SLOs; per-surface ranking teams
Why "first move" separates levels

L5: Leads with the engine. Reasonable — Elasticsearch is a fine default — but it skips the questions that change the design: a 5B-document web corpus with daily freshness is a different system from a 40M-product catalog with 5-second price freshness, and both differ from per-tenant enterprise search with access control.

L6: "Before the engine: corpus size and growth, update rate, the freshness SLO per field, query QPS and latency target, and how we measure relevance. I'll assume e-commerce: 50M products, 2K updates/s, price freshness ≤ 5 s, 5K QPS peak at p99 ≤ 200 ms, and relevance judged by search conversion."

L7: Asks whether search is a platform the company provides to 10 product surfaces, and who owns relevance quality when three surfaces want conflicting ranking.

Why "reindexing" separates levels

L5: "Reindex the data" — as if it's a command. At 50M docs and 2K updates/s, a full reindex takes hours, and updates arriving during it must not be lost; an analyzer change alters every document's tokens; a mapping change is impossible in place.

L6: "Every index is versioned: products_v7. Readers query an alias. To change the mapping, I create v8, backfill it from the source of truth via bulk indexing, dual-write live updates to v7 and v8 from the CDC consumer, verify doc counts and run a query-parity test on the top 10K queries, then atomically swap the alias. Rollback is swapping back — v7 stays current for 48 hours."

Why "ranking" separates levels

L5: Ranking is a formula: BM25 plus a popularity boost. Changes ship when they "look better" on a few queries.

L6: Ranking is a pipeline with measurement: cheap retrieval of ~1,000 candidates per shard, a lightweight L1 score, a learned L2 re-ranker on the top ~200, business rules last. Changes are evaluated offline (NDCG@10 on judged queries) and then online via interleaving/A/B with guardrails (zero-result rate, latency, revenue). "No ranking change ships on anecdotes."

The Staff Positions#

PositionRationale
The index is a derived store — always rebuildableSource of truth lives elsewhere; rebuild is a routine operation, not a disaster
Readers query aliases, never physical indexesZero-downtime reindex and instant rollback
Document-partitioned shards of 10–50 GBRecovery and rebalancing stay fast; fan-out stays bounded
Replicas scale QPS; shards scale dataDon't add shards to fix query throughput — it increases fan-out
Freshness SLO per field, not per indexPrice in 5 s, reviews in 1 h; don't pay NRT cost for everything
Multi-stage ranking with offline + online evaluationExpensive models only on the top few hundred; changes gated by metrics
Degrade by dropping expensive stages, not resultsUnder load, skip the L2 re-ranker before timing out queries

The Three Intents#

IntentConstraintStrategyFailure ModeCorrectness Bar
E-commerce / product searchFreshness of price/stock; conversion-driven rankingDoc-partitioned ES/OpenSearch, CDC updates, learned re-rankingStale price/availability; ranking regressions cost revenuePrice/stock freshness ≤ 5 s; zero-result rate < 3%
Web-scale / content searchBillions of docs; query latency at huge fan-outTiered index (fresh tier + base tier), batch index builds, massive replicationTail latency from fan-out; index build failuresCoverage and relevance; freshness hours–days (news minutes)
Multi-tenant enterprise / SaaS searchTenant isolation, access control per documentRouting by tenant, per-tenant indexes for large tenants, ACL filtering at query timeData leakage across tenants/permissions; noisy tenantsZero permission leaks; per-tenant latency SLO

🎯 Staff Move: "I'll design e-commerce product search, because it forces freshness, ranking and reindexing into the same conversation. Web-scale search changes the index build to batch tiers; enterprise search makes permissions the hardest problem — I'll point out where those diverge."

The Five Fault Lines#

#Fault LineThe Tension
1Freshness vs Indexing CostRefresh every second (more segments, merges, CPU) or batch updates (stale results)?
2Sharding: Fan-Out vs Shard SizeMany small shards (fast recovery, high fan-out) or few large ones (low fan-out, slow recovery)?
3Relevance vs LatencyExpensive ranking models improve quality but consume the latency budget
4Reindexing: In-Place vs Versioned RebuildMutate the live index (cheap, risky) or build a new one and swap (safe, 2× capacity)?
5Platform vs Per-Team SearchShared cluster and relevance stack, or each product runs its own?

In the Wild: Real Production Systems#

Why this section belongs here: Citing specific production systems demonstrates you've studied operational reality, not textbook designs.

Google — Tiered Indexes and Incremental Indexing (Percolator / Caffeine)#

Google's Percolator paper (2010) describes moving web indexing from large MapReduce batch rebuilds to incremental processing with cross-row transactions over Bigtable, powering the Caffeine indexing system and cutting the average age of documents in search results substantially. The broader pattern — a fast-updating fresh tier alongside a larger, slower-built base tier — is how web-scale engines balance freshness with build cost.

Staff insight: Freshness and index build cost are a tier decision, not a single knob. Name "fresh tier + base tier" and you've shown you know web-scale search isn't one index.

LinkedIn — Galene#

LinkedIn publicly described Galene, its Lucene-based search architecture, which builds base indexes offline in Hadoop and serves them alongside a live-update index for recent changes, with a federated broker fanning out to verticals (people, jobs, companies) and multi-stage ranking.

Staff insight: Offline-built base index + live update layer + federation is the production answer to "how do you reindex 1B documents without downtime?" — you don't reindex live; you build and ship snapshots.

Twitter — Earlybird#

Twitter's Earlybird (2012 paper) is a real-time inverted index optimized for tweets searchable within seconds: in-memory, append-only posting lists partitioned by time into segments, with concurrency designed for a single writer and many readers.

Staff insight: When freshness is the product, the index layout changes — time-partitioned, append-only, recency-ordered posting lists. Designing for the dominant query ("recent tweets matching X") beats a general-purpose index.

What Interviewers Probe#

After You Say...They Will Ask...(What They're Evaluating)
"Elasticsearch""How many shards, how big, and why?"Capacity reasoning, fan-out awareness
"Near real-time""What does 1-second refresh cost at 5K updates/s?"Freshness vs indexing cost
"Add replicas""Does that help write throughput? Tail latency?"Replicas vs shards
"Reindex""Mapping change on a live 2 TB index, zero downtime — walk me through it."Versioned index + alias
"BM25 plus ML""How do you know a ranking change is better?"Evaluation discipline
"Shard by tenant""One tenant is 30% of docs."Hot shard/tenant isolation

System Architecture Overview#

Diagram: System Architecture Overview

Reading the diagram: Writes flow from sources of truth through CDC and enrichment into a bulk indexer; the search cluster never receives writes directly from product code. Every document carries a version (source LSN or updated_at) so out-of-order updates can't regress it. Reads go through an alias, fan out to one copy of each of 20 shards, merge top-K, and re-rank a small candidate set. Signals (popularity, CTR) arrive via a slower batch path — they don't need second-level freshness.

Quick-Reference: The 30-Second Cheat Sheet#

TopicThe L5 AnswerThe L6 Answer — Say This
Index"Inverted index in Elasticsearch""Doc-partitioned Lucene index; 20 shards of ~30 GB; 2 replicas for QPS and availability; readers use an alias."
Freshness"Real-time""Price/stock ≤ 5 s via CDC + 1 s refresh; signals nightly. Per-field SLO, measured as index.freshness_lag_s."
Ranking"BM25 + popularity""Retrieve ~1K per shard with BM25, L1 features, learned L2 on top 200; offline NDCG gate then A/B with guardrails."
Reindex"Reindex the data""Build v8 beside v7, backfill, dual-write, verify counts + query parity, swap alias; rollback = swap back."
Scale reads"More shards""More replicas; shards increase fan-out. Cache head queries; shed L2 under load."
Tenants"Filter by tenant_id""Custom routing so tenant queries hit one shard; big tenants get dedicated indexes; ACL filters are mandatory, tested."

Key Numbers Worth Memorizing#

MetricValueWhy It Matters
Target shard size10–50 GBRecovery/relocation in minutes; bounded per-shard latency
Index size vs raw text~0.3–1.5× (depends on stored fields, doc values, positions)Capacity planning
Refresh interval1 s default (NRT); 30 s+ for bulk-heavyEach refresh creates a segment → merge cost
Bulk request size~5–15 MB / 1–5K docsThroughput sweet spot
Indexing throughput per node~5–30K docs/s (doc size-dependent)Pipeline sizing
Query fan-outN shards × 1 copy eachp99 of query ≈ worst shard
P(≥1 of 20 shards exceeds its p99)1 − 0.99²⁰ ≈ 18%Why hedging/replica selection matters
Latency budget (e-commerce)p99 ≤ 200 ms end-to-end; retrieval ~50–80 ms, re-rank ~15–30 msWhere ranking cost fits
Re-ranking depth100–500 candidatesBeyond this, cost rises faster than quality
JVM heap for ES≤ ~31 GB (compressed oops); rest for OS page cacheMemory sizing
Head query shareTop 1% of queries ≈ 20–40% of trafficResult caching pays off
Zero-result rate healthy range< 2–5%Primary relevance health signal

Interview Walkthrough

The most common mistake: Candidates spend 15 minutes explaining tokenization, TF-IDF and posting-list compression, then have no time for freshness, reindexing or ranking evaluation. The inverted index is table stakes — get past it in 3 minutes.


Phase 1: Requirements & Framing (2–3 min)#

State scope in one sentence:

"Users type a query and get the top 50 relevant products with filters and facets; sellers update products and prices continuously."

Non-functional constraints that pick the design:

"The numbers that decide the design: corpus size and growth, update rate and freshness SLO per field, query QPS and latency, and how relevance is measured. I'll assume 50M products (~2 KB each searchable, ~100 GB raw), 2K updates/s with bursts to 20K during repricing, price/stock visible within 5 s, 5K QPS peak, p99 ≤ 200 ms, and relevance measured by search conversion."

Commit:

"I'll design product search. Freshness matters for price and stock only; popularity signals can be a day old."

🎯 Staff Move: Ask how relevance is judged in Phase 1. It tells the interviewer you know search quality is measured, not asserted — and it sets up the ranking deep dive.


Phase 2: Core Entities & API (1–2 min)#

  • SearchDocument (product_id, version, title, description, brand, category_path, price, in_stock, seller_id, popularity, embedding?)
  • Index version (products_v7) behind alias products
  • Query (text, filters, sort, page_token, user_context)
  • Judgment (query, product_id, grade 0–3) for offline evaluation
Search(q, filters, sort, page_token, ctx) → { results[50], facets, total_estimate, next_page_token }
Suggest(prefix) → [completions]                 # separate, prefix-optimized index
IndexDocument(doc, version)                     # internal: pipeline only, idempotent by (id, version)

"Writes are internal and versioned — no product service writes to the search cluster directly."


Phase 3: High-Level Architecture (≤5 min)#

Diagram: Phase 3: High-Level Architecture (≤5 min)

Walk the flow:

  1. Product changes are captured by CDC into Kafka; enrichment joins seller/category/stock into a denormalized search document.
  2. The bulk indexer writes to the alias's current index with external versioning (older versions rejected).
  3. Queries hit the API: parse, spell-correct, expand synonyms, check the result cache.
  4. Coordinator fans out to one copy of each shard; each returns its top-K by BM25 + features; coordinator merges.
  5. L2 re-ranker scores the top 200 with a learned model; business rules (sponsored slots, diversity) applied last.

🎯 Staff Move: "This is the standard design. What decides quality and reliability is how fresh the index is and what that costs, how we change the index without downtime, and how we know ranking changes are improvements. Let me go there."


Phase 4: Transition to Depth (1 min)#

"Three deep dives: the freshness pipeline and its cost, sharding and tail latency under fan-out, and reindexing plus ranking changes as safe rollouts. I'd start with freshness — it's where stale prices turn into customer complaints. Your preference?"


Phase 5: Deep Dives (25–30 min)#

Deep dive 1: Freshness pipeline (7–8 min)

"End-to-end freshness = CDC lag (~200 ms) + enrichment (~100 ms) + bulk batching (up to 1 s) + refresh interval (1 s) ≈ 2–3 s p50, 5 s p99. The costly part is refresh: each refresh makes a new segment per shard; with 20 shards that's 20 segments/s being created and merged. At 2K updates/s it's fine; at 20K updates/s during a repricing burst, merge pressure competes with queries. So price updates use partial updates, and during declared bulk operations we raise refresh to 5 s for affected indexes."

Versioning: "Each doc carries the source LSN as an external version. If the CDC consumer replays or reorders, older versions are rejected — no stale overwrite."

Deep dive 2: Sharding and fan-out (6–8 min) — 100 GB raw → ~150 GB index → 20 shards of ~7.5 GB now, sized to ~30 GB at 4× growth. Replicas for QPS. Adaptive replica selection, timeouts with partial results, hedged requests for the slowest shard.

Deep dive 3: Reindexing with alias swap (6–8 min) — v8 build, dual-write, parity test, swap, 48-h rollback window.

Deep dive 4: Ranking evaluation (5–7 min) — judged query set, NDCG@10 offline gate, interleaving then A/B with guardrails (latency, zero-result rate, revenue per search).


Phase 6: Wrap-Up (2–3 min)#

"Summary: CDC-fed, versioned documents; 20 shards × 3 copies behind an alias; 1 s refresh with per-field freshness SLOs; multi-stage ranking with offline and online evaluation; versioned reindex with alias swap. Next I'd add hybrid retrieval with embeddings for long-tail queries and query understanding (category prediction). I wouldn't build a custom engine or a term-partitioned index at this scale."

Common Timing Mistakes#

MistakeTime LostFix
Explaining TF-IDF/BM25 math5–7 min"BM25 for candidates; the math isn't the design."
Posting-list compression3–5 min"Lucene handles encoding."
Autocomplete deep dive uninvited5+ min"Autocomplete is a separate prefix index; happy to cover it."
Ignoring writes entirely—Freshness is half the problem
"ML ranking" with no evaluation—Say how you measure

1. The Staff Lens#

1.1 Why This Problem Exists in Staff Interviews#

Search sits between many owners: source-of-truth teams who change schemas, a platform team running clusters, product teams who want their items ranked higher, and a relevance team trying to optimize one metric without breaking five others. The technical core is well-known; the Staff skill is making the index rebuildable, freshness measurable, and ranking changes governed — and keeping p99 flat while fan-out grows.

1.2 The L5 vs L6 Contrast — Visual#

Diagram: 1.2 The L5 vs L6 Contrast — Visual

1.3 The Staff Question That Cuts Through Everything#

"If we had to rebuild the entire index from scratch tomorrow, how long would it take, and would anyone notice?"

If the answer is "days, and yes," the index has become a source of truth by accident — documents written only to search, no versioning, no alias. Every Staff design choice (CDC from the source, versioned documents, aliases, parity tests) exists to make the answer "a few hours, and no."


2. Problem Framing & Intent#

2.1 The Three Intents — Explained#

E-commerce / product search. Corpus 10M–1B items; updates dominated by price/stock; ranking optimized for conversion with many business constraints (sponsored, diversity, availability). Freshness of transactional fields matters more than of text.

Web-scale / content search. Billions of documents; crawl-driven ingestion; index build as a batch pipeline producing immutable snapshots shipped to serving; a fresh tier for news/recent content. Serving is massively replicated; the challenge is fan-out at thousands of shards and index build cost.

Multi-tenant enterprise search. Many tenants with power-law sizes; per-document ACLs; strict isolation. Query must filter by tenant and permissions with zero leakage; small tenants share indexes with routing, large tenants get dedicated ones. The correctness bar is security, not relevance.

🎯 Staff Move: "The permission model is the design driver for enterprise search — a single missed ACL filter is a data breach. For product search, it's freshness of price and stock. I'll name which one I'm optimizing."

2.2 When NOT to Use a Search Engine#

SituationBetter AlternativeWhy
Exact lookups by ID/keyPrimary DB / KV storeSearch engines aren't systems of record
< 1M rows, simple keyword matchPostgreSQL full-text (tsvector + GIN)One fewer system to operate
Structured filters only (no text relevance)DB indexes / OLAPInverted index adds nothing
Strong read-after-write requiredQuery the DB for the writer's own itemsNRT refresh means ~1 s invisibility
Prefix autocomplete onlyDedicated prefix index / trie serviceDifferent data structure, tighter latency
Analytics aggregations over eventsOLAP / columnar storeSearch clusters are expensive analytics engines

"If a user edits a listing and immediately views 'my listings,' I serve that page from the database, not search — refresh lag would look like data loss."

2.3 What the Interviewer Leaves Underspecified#

GapWhy It MattersWhat to Say
Corpus size & growthShard count"At 2× yearly, I size shards for year-2 volume."
Update rate & burstsPipeline + refresh cost"Are there bulk repricing events?"
Freshness per fieldWhere to pay for NRT"Price seconds? Reviews hours?"
Relevance metricEvaluation design"Conversion? CTR? Human judgments?"
PersonalizationCache hit rate, ranking cost"Per-user ranking kills result caching."
Access controlQuery filtering design"Any per-document permissions?"
LanguagesAnalyzer choice per field"Multi-language corpus needs per-language analyzers."

2.4 Precise Terminology#

TermPrecise MeaningCommon Confusion
Inverted indexTerm → postings (doc IDs, freqs, positions)Confused with a B-tree index
AnalyzerTokenizer + filters (lowercase, stemming, synonyms)Changing it requires reindexing
SegmentImmutable mini-index; searchable after refreshNot a shard
RefreshMakes buffered docs searchable (new segment)Not durability — that's the translog/commit
MergeBackground compaction of segments; purges deletesSource of I/O spikes
Shard / replicaPartition of the index / copy of a partitionReplicas don't add write capacity
Recall vs precisionFraction of relevant docs found vs fraction of returned docs relevantTraded by retrieval depth and filters
NDCG@KGraded relevance metric discounted by rankNeeds judgments; not CTR
AliasName pointing to one or more physical indexesEnables atomic swaps
Query-time vs index-timeSynonyms, boosts etc. applied at query vs indexIndex-time changes require reindex

3. The Five Fault Lines#

3.1 Fault Line 1: Freshness vs Indexing Cost#

The tension: Each refresh creates segments that must be merged; each update is a delete + reinsert in Lucene-family engines. Faster freshness means more CPU and I/O on the same nodes serving queries.

StrategyWhat WorksWhat BreaksWho Pays
Nightly full rebuildCheap, clean, optimized segments24 h staleness; price mismatches at checkoutCustomers + support
Hourly incremental batchesModerate costPrice errors for up to 1 hBusiness (mispriced results)
CDC + 1 s refresh (NRT)~2–5 s freshnessMerge pressure; query p99 rises during burstsSearch platform (capacity)
Split: NRT for hot fields, batch for signalsPay for freshness only where neededTwo pipelines to ownSearch platform (complexity)
Serve volatile fields from a side store at query timeInstant price/stockExtra lookup per result (~5–10 ms for 50 items)Latency budget
Diagram: 3.1 Fault Line 1: Freshness vs Indexing Cost

Staff default: Per-field freshness SLOs. CDC + 1 s refresh for fields used in filters/sort (price, in_stock). Display-only volatile data fetched at query time for the final 50 results. Signals batch-updated nightly via partial updates. Raise refresh interval during declared bulk operations.

When to deviate: Web-scale: batch-built base index + small fresh tier. Social/real-time: time-partitioned append-only index (Earlybird-style).

🎯 Staff Move: "I don't want 'real-time search' — I want price and stock within 5 seconds and everything else within a day. That's a third of the indexing cost."

3.2 Fault Line 2: Sharding — Fan-Out vs Shard Size#

The tension: Document-partitioned search sends every query to every shard. More shards → smaller shards (faster recovery, more parallelism per query) but more fan-out (more tail risk, more coordinator merge work, more per-shard overhead).

ConfigurationWhat WorksWhat BreaksWho Pays
Few large shards (5 × 100 GB)Low fan-outRecovery takes an hour; per-shard latency grows; can't spread across many nodesOn-call (slow recovery)
Right-sized (20 × 10–50 GB)Balanced——
Many small (500 × 1 GB)Fast recoveryFan-out 500; heap overhead per shard; cluster-state bloatEvery query (tail) + platform
Routing by tenant/categoryQuery hits 1 shardSkewed shards; cross-tenant queries fan outPlatform (hot shards)
Term-partitionedFew shards per queryMulti-term queries cross shards; updates touch many shardsEveryone (rarely worth it)

Tail latency controls:

  • Adaptive replica selection: route each shard request to the copy with the lowest recent latency/queue.
  • Per-shard timeout with partial results: return with 19 of 20 shards if one exceeds 150 ms; flag partial=true, log it.
  • Hedged requests: after p95 latency, send a duplicate to another replica; take the first.
  • Result cache for head queries (60 s TTL).
Diagram: 3.2 Fault Line 2: Sharding — Fan-Out vs Shard Size

Staff default: Shards sized 10–50 GB at the 18-month projected size; replicas to meet QPS with 2× headroom; replica selection + hedging + partial-result timeouts. Shard count set once per index version (changing it = reindex, or split/shrink operations).

3.3 Fault Line 3: Relevance vs Latency#

The tension: Better models (cross-encoders, many features, personalization) improve ranking but cost milliseconds and CPU/GPU per candidate.

StageCandidatesBudgetTechniqueWho Pays If Too Slow
Retrievalmillions → ~1,000 per shard30–60 msBM25 + filters (+ ANN)All queries
L1 scoring~1,000 → 2005–10 msLinear/GBDT on cheap features, on shardShard CPU
L2 re-rank200 → 5015–30 msGBDT / neural on rich featuresRanking service capacity
Business rules501–2 msSponsored slots, dedupe, diversityRelevance (rules degrade quality)

Staff default: Multi-stage. Put the expensive model only on the top 200. Under load, skip L2 (serve L1 order) before timing out: a slightly worse ranking beats an error page.

Degradation ladder: full pipeline → skip personalization features → skip L2 → reduce per-shard depth to 200 → serve cached results for head queries → reject tail queries with 503. Each step is a flag with a named owner.

3.4 Fault Line 4: Reindexing — In-Place vs Versioned Rebuild#

The tension: Mapping and analyzer changes can't be applied in place. Even when possible, updating 50M documents in the live index competes with queries and can't be rolled back.

Diagram: 3.4 Fault Line 4: Reindexing — In-Place vs Versioned Rebuild
ApproachWhat WorksWhat BreaksWho Pays
In-place update by queryNo extra capacityCan't change mappings/analyzers; no rollback; query impactUsers (latency) + on-call
Versioned index + alias swapZero-downtime; instant rollback~2× storage and indexing during migrationPlatform budget
Offline build + snapshot ship (web-scale)Clean, optimized, reproducibleNeeds build infra; freshness via separate tierPlatform (build pipeline)

Staff default: Always versioned indexes behind aliases; reindex from the source of truth (not from the old index if the old one might be wrong); dual-write via the CDC consumer; parity testing with a judged/top-query set; 48-h rollback window.

Parity test: run the top 10K queries against v7 and v8; compare result overlap (Jaccard@10) and NDCG@10 on judged queries; expected change (for an analyzer change) is declared up front — e.g., "overlap ≥ 0.8 except queries with plurals."

OptionWhat WorksWhat BreaksWho Pays
Each team runs its own clusterAutonomy, independent ranking12 clusters, 12 versions, duplicated pipelines, inconsistent qualityEvery team's on-call
Shared cluster, shared relevance teamConsistent ops & qualityRelevance team bottleneck; noisy neighborsRelevance team
Platform (clusters, pipeline, eval tools) + per-surface ranking ownershipConsistency where it matters, autonomy in rankingContract boundaries to maintainPlatform + surface teams
Managed search serviceNo cluster opsCost at scale; limited controlBudget

Staff default: Search platform owns clusters, indexing pipeline, alias/versioning tooling and evaluation infrastructure; product surfaces own their ranking models and business rules within latency budgets; source teams own document contracts.

🧭 Principal Move: "Relevance is a product with a P&L impact. I'd give each surface a ranking owner and a shared experiment platform — and a rule that no ranking change ships without an experiment readout."


4. Failure Modes & Operational Reality#

4.1 Merge Storm During Bulk Repricing#

Scenario: A seller-tools feature lets merchants reprice entire catalogs. A promotion triggers 30M price updates in 40 minutes.

t=0:       Repricing starts; update rate 2K/s → 12K/s
t=+2min:   Refresh at 1 s creates small segments on all 20 shards; merge threads saturate disk
t=+5min:   Query p99 120 ms → 900 ms; coordinator timeouts cause partial results on 8% of queries
t=+8min:   Search thread pool queues fill; shard.rejected_requests > 0
t=+10min:  Conversion from search drops 11%
t=+15min:  On-call throttles the indexer; freshness lag grows to 20 min for all products
t=+50min:  Backlog drains; p99 recovers

Detection: search.p99_ms > 300; index.merge_time_ms spike; shard.rejected_requests > 0; index.freshness_lag_s growth. Blast radius: All search users; all product freshness during throttling. Mitigation: Priority lanes in the indexer: price/stock partial updates first, bulk repricing lane rate-limited; temporarily raise refresh to 5–10 s. Prevention: Bulk operations are declared events with an ingestion budget; separate indexing and query node roles (or dedicated ingest capacity) at larger scale. Owner: Search platform (indexer lanes); seller-tools team (declares bulk jobs).

🎯 Staff Move: "Freshness and query latency share the same disks. I protect queries first and make bulk writers queue — and I make them tell us before they start."

4.2 Silent Staleness — The Stuck CDC Consumer#

t=0:       Enrichment consumer hits a poison message (malformed seller record); retries forever
t=+0s:     Partition 7 of products topic stops advancing; 1/32 of products stop updating
t=+3h:     Customer sees "$49" in search, "$79" at checkout — 12K complaints over the day
t=+9h:     Support escalates; consumer lag on partition 7 = 3.1M messages

Why it's silent: Search is up, latency is fine, most products update. Only a per-partition lag metric or an end-to-end freshness probe reveals it. Detection: consumer_lag{partition} max > 60 s; synthetic freshness probe (update a canary product every minute, measure time-to-searchable) — index.freshness_lag_s p99 > 30 s pages. Mitigation: Poison messages to DLQ after 3 retries; per-partition lag alerts. Prevention: Canary documents per partition; checkout re-validates price from source of truth (search is never authoritative for price). Owner: Search platform (pipeline); catalog team (bad records).

4.3 Reindex Swap Regression#

Scenario: v8 added a new analyzer with aggressive stemming. Alias swapped after doc-count verification only.

t=0:       Alias swap v7 → v8
t=+10min:  "glasses" now matches "glass" — drinking glasses outrank eyeglasses
t=+1h:     Zero-result rate unchanged, but search conversion −4.5% on eyewear
t=+6h:     Category manager notices; rollback via alias swap in 30 s

Lesson: Counts prove completeness, not quality. Parity tests on top queries + judged NDCG + a small online A/B before full swap. Detection: search.conversion_rate by category vs baseline; search.ctr_at_10 drop; parity jaccard@10 distribution. Owner: Relevance team (quality gate); platform (swap tooling).

4.4 Hot Shard from Custom Routing#

A marketplace routes by seller_id to make "search within store" single-shard. One mega-seller holds 18% of listings → one shard at 3× others' size and query load. Detection: shard.size_gb max/median > 2; shard.query_latency_ms skew. Mitigation: Routing partitions (route a key to a set of N shards), or dedicated index for mega-sellers. Owner: Search platform.

4.5 ACL Filter Missing (Enterprise)#

A new query path for "recent documents" is added without the permission filter. Users see titles of documents they can't open. Detection: Automated permission tests per query path (a canary user must never see canary restricted docs); search.acl_filter_missing_total from a query-builder assertion. Mitigation: All queries built by one library that injects tenant + ACL filters; direct engine queries prohibited. Owner: Search platform + security.

4.6 Coordinator Memory Blow-Up — Deep Pagination#

A scraper requests from=50000&size=100. Each shard returns 50,100 candidates; the coordinator sorts 1M entries per request; heap pressure triggers GC pauses cluster-wide. Mitigation: Cap from + size (e.g., 10K); use search-after cursors; rate-limit anonymous deep pagination. Owner: Search API team.

4.7 Operational Reality Matrix#

FailureDetection SignalBlast RadiusMitigationOwner
Merge stormindex.merge_time_ms, search.p99_msAll queriesIndexer lanes, refresh interval upSearch platform
Stuck pipeline partitionconsumer_lag{partition}, freshness probeSubset of docs staleDLQ, per-partition alertsSearch platform
Reindex quality regressionParity + conversion by categoryAffected query classesAlias rollbackRelevance team
Hot shardshard.size_gb / latency skewQueries touching that shardRouting partitions, dedicated indexSearch platform
ACL leakCanary permission testsSecurity incidentMandatory filter injectionPlatform + security
Deep paginationCoordinator heap, GC timeCluster-widesearch_after, capsSearch API
Node losscluster.status yellow/redReduced capacity, recovery I/OReplicas; throttled recoverySearch platform
Ranking service downrerank.error_rateQuality onlyServe L1 orderRanking team
Zero-result spikesearch.zero_result_rate > 5%RevenueRollback synonyms/analyzer; spell fallbackRelevance team

5. Evaluation Rubric#

5.1 Level-Based Signals#

DimensionSenior (L5)Staff (L6)Principal (L7)
Index designInverted index in ESSized shards/replicas with growth math; doc-partitioned with tail controlsTopology standards across products; cost per QPS and per GB
Freshness"Near real-time"Per-field SLOs; CDC with versioning; refresh cost reasoning; freshness probesFreshness tiers negotiated and priced with product owners
Reindexing"Reindex"Versioned index + alias + dual-write + parity + rollbackSelf-service index-versioning platform with schema-change SLA
RankingBM25 + boostsMulti-stage with budgets; offline + online evaluation; degradation ladderRelevance governance: ownership per surface, experiment policy, rule budgets
OperationsCluster healthFreshness lag, zero-result rate, merge pressure, partial results — with ownersError budgets per surface; game days (node loss, pipeline stall, bad swap)
OwnershipSearch team does it allSource teams own doc contracts; platform owns cluster/pipeline; surfaces own rankingOrg boundaries between platform, relevance and product

5.2 Strong Hire Signals#

SignalWhat It Sounds Like
Rebuildable index"Search is derived; I can rebuild it from the DB in ~3 hours via the backfill path."
Per-field freshness"Price and stock in 5 s; reviews hourly; popularity nightly."
Replicas vs shards"QPS problem → replicas. Data-size problem → shards. More shards makes fan-out worse."
Alias discipline"Readers only see aliases; every mapping change is a new version and a swap."
Measured ranking"Offline NDCG gate, then A/B with guardrails on latency and zero-result rate."

5.3 Lean No-Hire Signals#

SignalWhy It Misses the Bar
App writes directly to ESIndex becomes a source of truth; no replay/rebuild
"Add shards" for QPSMisunderstands fan-out
No freshness mechanismStale prices ignored
Ranking changes by eyeballingNo evaluation discipline
Search authoritative for priceStale index becomes a billing bug

5.4 Common False Positives#

  • Deep Lucene internals (FST, skip lists) ≠ search system design.
  • Vector search enthusiasm — "use embeddings" without recall/latency/cost tradeoffs or a hybrid plan.
  • Elaborate query parsing while ignoring the indexing pipeline.
  • Naming Learning-to-Rank without feature freshness, training data bias or evaluation.

6. Interview Flow & Pivots#

6.1 Typical 45-Minute Shape#

PhaseTimeGoal
Framing0–3 minCorpus, updates, freshness per field, QPS/latency, relevance metric
Entities + API3–5 minSearch document, versioning, alias, search API
High-level design5–10 minCDC pipeline, cluster, query path, ranking stages
Transition10–11 minOffer freshness / sharding & tail / reindex & ranking
Deep dives11–40 minFreshness cost → fan-out tail → reindex → ranking evaluation
Wrap-up40–45 minEvolution (hybrid retrieval), what not to build, owners

6.2 How Interviewers Pivot — And What They're Testing#

Interviewer SaysWhat They're TestingWhere to Go
"Prices are wrong in search"Freshness + authorityCDC, versioning, probes, checkout re-validation
"p99 doubled"Tail latency reasoningFan-out, merges, hedging, replica selection
"We need to change the analyzer"Reindex safetyVersioned index + parity + swap
"Make results more relevant"Ranking disciplineMulti-stage, metrics, experiments
"10× the corpus"Capacity mathShard sizing, tiers, cost
"Add per-document permissions"Security correctnessFilter injection, ACL indexing, tests

6.3 What to Deliberately Skip#

TopicWhy L5 Goes HereWhat L6 Says Instead
BM25 formulaFamiliar"BM25 for candidates; tuning k1/b rarely matters vs features."
Posting-list compressionInteresting"Engine handles it."
Stop-word listsEasy"Default analyzers per language."
Cluster master electionFeels deep"3 dedicated master nodes; standard."

6.4 Follow-Up Questions to Expect#

  1. "How do you guarantee an updated price shows within 5 seconds, and how do you know it does?"
  2. "Walk me through changing the mapping on a live index."
  3. "How many shards and replicas, and why?"
  4. "How do you handle a query that times out on one shard?"
  5. "How do you evaluate a new ranking model before launching it?"
  6. "How would you add semantic/vector search?"
  7. "How do you handle deletions and GDPR erasure in the index?"

7. Active Drills#

Drill 1: The Opening#

Prompt: "Design search for our e-commerce site."

Staff Answer

"Before the engine, five numbers: corpus (say 50M products), update rate (2K/s, bursts to 20K), freshness per field (price/stock ≤ 5 s, else daily), QPS and latency (5K QPS, p99 ≤ 200 ms), and the relevance metric (search conversion, plus a judged query set).

Design: the product DB is the source of truth; CDC into Kafka; enrichment builds a denormalized search document with a version = source LSN; a bulk indexer writes to an index behind an alias. 20 shards sized for 2 years, 2 replicas. Query path: parse and expand, fan out, BM25 + L1 on shard, learned L2 on the top 200, business rules. The deep problems are freshness cost, fan-out tail latency, and safely changing the index and ranking — I'll go into those."

Why this is L6:

  • Numbers and a relevance metric before architecture
  • Index as a derived store behind an alias
  • Names the deep problems unprompted

What L7 adds:

  • Asks whether other surfaces (seller search, support tools) should share the platform and document contract
  • States the ranking ownership model: who decides tradeoffs between conversion and sponsored revenue

Drill 2: Inverted Index Mechanics That Matter#

Prompt: "How does an update to a product become searchable?"

Staff Answer

"In Lucene-family engines, documents are immutable within segments. An update is a delete marker on the old doc plus a new doc in the in-memory buffer; it's written to the translog for durability. On refresh (every 1 s), the buffer becomes a new small segment and is searchable. Background merges combine segments and purge deleted docs. So: frequent updates → many tombstones and small segments → merge I/O. A partial update still reindexes the whole document internally, which is why I keep the search document lean (~2 KB) and avoid giant stored fields."

Why this is L6:

  • Connects mechanics to cost (updates = delete + insert; merges)
  • Draws a design consequence (lean documents)

What L7 adds:

  • Uses this to justify moving ultra-volatile display fields out of the index org-wide

Drill 3: Make Freshness Concrete#

Prompt: "Prove to me a price change shows up within 5 seconds."

Staff Answer

"Budget: CDC capture ~200 ms, Kafka ~50 ms, enrichment ~100 ms (cached seller/category lookups), bulk batching flush at 500 docs or 500 ms, index + refresh ≤ 1 s. p50 ~2 s; p99 ~4 s when healthy.

Proof is measurement: a canary product per Kafka partition gets a price change every minute; a prober searches for it and records time-to-visible → index.freshness_lag_s. SLO: p99 ≤ 5 s over 30 days; page if > 30 s for 5 min. Also per-partition consumer lag alerts, since one stuck partition affects 1/32 of products while averages look fine. And checkout re-validates price from the price service regardless — search is never the authority."

Why this is L6:

  • Decomposed budget with numbers
  • Measures end-to-end with canaries per partition
  • Refuses to make search authoritative

What L7 adds:

  • Publishes freshness SLOs per field as part of the search platform contract, with error budgets that gate risky pipeline changes

Drill 4: The Cluster Is Degraded#

Prompt: "Two of 12 data nodes die. What happens?"

Staff Answer

"With 2 replicas (3 copies) and shard allocation awareness across 3 zones, no shard loses all copies — status yellow, not red. Query capacity drops ~17% and recovery I/O (re-replicating ~500 GB) competes with queries, so I throttle recovery bandwidth (e.g., 100 MB/s per node) and delay allocation 5 min in case the nodes return. If QPS headroom was < 20%, enable the degradation ladder: skip L2, raise result-cache TTL to 5 min for head queries. If they were in the same zone and replicas weren't zone-aware, some shards could go red — partial results with partial=true rather than failing all queries."

Why this is L6:

  • Zone-aware replica placement
  • Throttled recovery vs query capacity
  • Degrade quality before availability

What L7 adds:

  • Capacity policy: N+2 nodes or 30% headroom per cluster, enforced in the capacity review

Drill 5: Hot Tenant / Hot Shard#

Prompt: "Marketplace search 'within store' is slow for our biggest seller."

Staff Answer

"If we route by seller_id so store search hits one shard, a seller with 18% of listings creates one oversized, overloaded shard. Options: routing partitions — route each seller to a group of 4 shards (routing_partition_size), spreading the big seller while keeping small sellers' queries to 4 shards not 20. Or give the top ~10 sellers dedicated indexes, with an alias per seller-group. For global search, the skew matters less because every query fans out anyway — the problem is size imbalance slowing that one shard, which caps global p99 too."

Why this is L6:

  • Knows routing trades fan-out for skew
  • Offers graduated fixes
  • Connects hot shard to global tail latency

What L7 adds:

  • Enterprise seller tier includes dedicated capacity, priced

Prompt: "Build search for a SaaS document product with per-document permissions."

Staff Answer

"Isolation first. Small tenants share indexes with routing by tenant_id; tenants above ~5M docs get dedicated indexes. Permissions: index each document's ACL principals (users, groups) as a field; at query time, expand the user's groups (cached, ~60 s TTL) and inject a filter tenant_id = X AND acl IN (user, groups…). The filter is injected by one query library — no raw engine access. Permission changes are high-priority updates in the pipeline (target < 10 s) because a revoked user seeing a title is a leak. Canary tests: a synthetic user per cluster must never see restricted canary docs; any hit pages security."

Why this is L6:

  • Treats ACL correctness as the primary requirement
  • Centralized filter injection + canary testing
  • Prioritizes permission updates in the pipeline

What L7 adds:

  • Security review sign-off as a gate in the search platform standard; per-tenant encryption keys for regulated customers

Drill 7: Build vs Buy#

Prompt: "Self-host OpenSearch, use a managed service, or a SaaS search API?"

Staff Answer

"SaaS search APIs (Algolia-style) are fastest for < 10M docs with standard ranking — but cost scales per record and per request, and custom ranking/learned models are constrained. Managed Elasticsearch/OpenSearch removes node ops but not index design, pipeline or ranking work. Self-hosting needs ~2–3 engineers for cluster ops at medium scale. For 50M products with custom learned ranking, I'd use managed ES/OpenSearch and invest our engineers in the pipeline and relevance — where differentiation lives. Revisit self-hosting when the managed premium exceeds ~2 FTE of cost."

Why this is L6:

  • Separates cluster ops from search engineering
  • Puts engineers where differentiation lives

What L7 adds:

  • Evaluates license and vendor risk (e.g., Elasticsearch/OpenSearch licensing split) and keeps queries portable

Drill 8: Change Ranking Without an Outage#

Prompt: "The team has a new ranking model. Ship it."

Staff Answer

"Gate 1 offline: NDCG@10 on the judged set (~10K queries graded 0–3) must improve with no regression > 2% on any major category; latency on replayed traffic within the 30 ms L2 budget. Gate 2 shadow: score production traffic without serving; compare features for skew. Gate 3 interleaving on 5% — sensitive to small differences with less traffic. Gate 4 A/B at 10% for 2 weeks: primary metric search conversion; guardrails zero-result rate, p99, sponsored revenue, long-tail query CTR. Ramp 10 → 50 → 100% with a flag. Rollback = flip the flag, instant."

Why this is L6:

  • Four gates with metrics and guardrails
  • Rollback is a flag, not a deploy

What L7 adds:

  • Experiment policy org-wide: required readouts, minimum durations, who can override a failed guardrail

Drill 9: Cost#

Prompt: "The search cluster costs $220K/month. Cut it."

Staff Answer

"Where's the money: (1) replicas sized for peak all day — scale replicas down off-peak or use read-only replica nodes that autoscale; (2) index bloat — _source storing full descriptions, unused fields indexed, positions on fields never phrase-queried; trimming mappings often saves 30–50% of disk and heap; (3) too many shards — per-shard overhead; consolidate small indexes; (4) cold data — out-of-stock products older than 6 months in a smaller, slower tier; (5) result cache for head queries cuts QPS 20–40%. I'd target 30–40% savings without p99 regression, verified by replaying traffic."

Why this is L6:

  • Levers tied to mapping and topology decisions
  • Verification by replay

What L7 adds:

  • Cost per 1K queries and per GB as platform KPIs; showback to surfaces by query volume

Drill 10: Multi-Region#

Prompt: "Serve search from 3 regions."

Staff Answer

"Search is derived, so each region runs its own cluster fed from the regional replica of the CDC stream (or cross-region replicated Kafka) — no cross-region query fan-out. Freshness per region depends on replication lag (~1 s extra). Index versions and alias swaps are coordinated: build v8 in each region, verify each, swap region by region (canary region first). Ranking models are deployed per region with the same flag. For residency-bound corpora (EU seller data), the EU cluster indexes only EU-permitted documents."

Why this is L6:

  • Regional independence because the index is derived
  • Region-by-region swaps as canaries

What L7 adds:

  • Decides whether all regions need full corpora or regional subsets; prices 3× cluster cost vs latency gains

8. Deep Dive Scenarios#

Deep Dive 1: Peak-Traffic Latency Collapse#

Context: Cyber Monday. Search QPS is 3.2× normal. p99 went from 140 ms to 1.4 s; 6% of queries return errors. On-call escalates to you.

Questions to Surface First:

  • Is the bottleneck query CPU, coordinator merge, merges/indexing, or the re-ranker?
  • Is traffic organic or a scraper doing deep pagination?
  • What's the result-cache hit rate?
  • Is indexing volume also elevated (promo repricing)?

Typical L5 Approach: Add nodes. New nodes need shard relocation, which adds recovery I/O during the peak — making things worse for 30+ minutes.

Staff Approach: Apply the degradation ladder immediately: skip L2 re-ranking (saves ~25 ms and a service hop), raise head-query cache TTL from 60 s to 5 min, cap from+size at 1,000, rate-limit anonymous clients. Pause non-critical bulk indexing (keep price/stock lane). Then diagnose: per-shard latency skew → a hot or merging shard; uniform → capacity. Add replicas only if nodes have headroom and relocation can be throttled.

Principal Approach: Peak readiness is a process: load tests at 4× last peak 6 weeks ahead, pre-scaled replicas, bulk-indexing freeze windows agreed with seller tools. The degradation ladder is pre-approved by the business with the expected conversion cost of each step (e.g., skipping L2 ≈ −1.5% conversion — far better than errors).

Staff Approach — Full Reasoning
PhaseWhat to Do
Immediate (0–5 min)Flags: disable L2, cache TTL 5 min, pagination cap, anon rate limits.
Triageshard.query_latency_ms by node/shard; index.merge_time_ms; search.cache_hit_rate; top clients by QPS.
Quick fixPause bulk lane; if a node is hot from merges, route away with replica selection; add replicas with throttled recovery if headroom allows.
GuardrailsKeep index.freshness_lag_s for price/stock under 30 s even while bulk paused.
Post-mortemWhy no pre-scaling? Why did repricing coincide with peak? Was the ladder rehearsed?

Metrics to Watch: search.p99_ms, search.error_rate, shard.rejected_requests, search.cache_hit_rate, rerank.skipped_total, index.freshness_lag_s

Organizational Follow-up: Peak calendar with seller tools and marketing; quarterly ladder drills.

Ownership Question: "Who approves turning off the re-ranker during peak?" Staff answer: The search on-call, via a pre-approved runbook — the relevance team and business agreed on the conversion cost in advance.

Key Takeaway: "Degrade ranking quality before availability. A worse order converts; an error page doesn't."

What clears the Staff bar:

  • Degradation ladder before scaling
  • Protects freshness of transactional fields during shedding
  • Avoids relocation during peak

Deep Dive 2: Silent Relevance Decay#

Context: Over 6 weeks, search conversion dropped 7% with no deploys to the search cluster and no alerts. Leadership asks what happened.

Questions to Surface First:

  • Is the drop uniform or concentrated in categories or query types?
  • Did the zero-result rate or query-reformulation rate change?
  • Did upstream data change (catalog fields, seller-provided titles)?
  • Did the ranking model's feature distributions drift?

Typical L5 Approach: Retrain the ranking model on recent data. May bake the problem into the model if the cause is data quality.

Staff Approach: Slice by category and query class: the drop concentrates in apparel. Root cause: a catalog team changed brand from a normalized field to free text for new listings; brand filters and the brand_match feature silently degraded for ~30% of new apparel items. Fix the document contract (normalized brand via enrichment), backfill affected docs, add feature-distribution monitoring (null rate, cardinality per feature) and a schema contract check at the pipeline boundary.

Principal Approach: Search quality depends on upstream data it doesn't own. Establish document contracts with source teams (field semantics, not just types) enforced at pipeline ingestion, and a relevance health dashboard reviewed weekly by surface owners. Data changes that affect search features require search sign-off — like an API change.

Staff Approach — Full Reasoning
PhaseWhat to Do
ImmediateConfirm with holdout: is the old model on today's data equally bad? (Yes → data, not model.)
TriageFeature null-rate and cardinality by week; apparel brand_normalized null rate 2% → 31%.
Quick fixEnrichment maps free-text brand to canonical IDs; partial-update backfill of 4M docs over 6 hours.
GuardrailsPipeline rejects/flags docs violating contract; alert on feature null-rate change > 5 points.
Post-mortemContract gap between catalog and search; no feature monitoring.

Metrics to Watch: search.conversion_rate{category}, search.zero_result_rate, search.reformulation_rate, feature.null_rate{feature}

Organizational Follow-up: Document contract owned by catalog team, reviewed by search; relevance health weekly review.

Ownership Question: "Who owns 6 weeks of lost conversion?" Staff answer: Shared — catalog changed semantics without a contract; search lacked monitoring to catch it. The fix is a contract both teams sign, not blame.

Key Takeaway: "Ranking regressions are usually data regressions. Monitor features, not just outcomes."

What clears the Staff bar:

  • Separates model from data with a holdout
  • Monitors feature distributions
  • Turns implicit dependencies into contracts

Deep Dive 3: Onboarding a Massive New Corpus#

Context: The company acquires a marketplace with 400M listings (8× current corpus) to be searchable within 3 months, in a unified search experience.

Questions to Surface First:

  • Unified index or federated (separate indexes, merged results)?
  • Do the listings share a schema and relevance features with ours?
  • What's the update rate of the new corpus?
  • What latency budget do we keep?

Typical L5 Approach: Reindex everything into the existing index with more shards. Changes shard count (requires reindex anyway), risks existing quality, and mixes very different ranking signals.

Staff Approach: Federate first: a separate index (acquired_v1, ~160 shards sized to 30 GB) with its own pipeline; the search API queries both in parallel and merges with a calibrated blending model (scores aren't comparable across indexes). Unify later if schemas converge. Capacity: 8× docs → roughly 8× storage and more query CPU; stage by category.

Principal Approach: Decide the long-term platform: one document contract and one ranking stack, or a federated architecture as a permanent feature (as LinkedIn does across verticals). Price both: federation costs a blending layer and ongoing calibration; unification costs a large migration. Set the decision point at 6 months with data.

Staff Approach — Full Reasoning
PhaseWhat to Do
Month 1Map schema to a minimal common search document; build pipeline; index into acquired_v1.
Month 2Federated query with parallel fan-out; blending model trained on combined clicks; offline eval.
Month 3A/B at 5% → 50%; guardrail: our existing categories' conversion must not drop > 0.5%.
AfterDecide unify vs federate based on overlap and quality data.

Metrics to Watch: search.p99_ms (federated adds max of two paths), blend.share_acquired_results, search.conversion_rate by source

Organizational Follow-up: Acquired team owns their document contract; relevance owns blending.

Ownership Question: "Who decides how much acquired inventory ranks above ours?" Staff answer: The blending model decides by relevance — with a business-owned policy on any explicit boosts, reviewed through the experiment process, not by fiat.

Key Takeaway: "Federate to ship, unify when the data says so. Scores from different indexes aren't comparable — calibrate."

What clears the Staff bar:

  • Federation vs unification as a staged decision
  • Knows cross-index score calibration
  • Protects existing quality with guardrails

Deep Dive 4: Post-Mortem — The Synonym Change That Emptied Search#

Context: A merchandiser updated the synonym file (index-time synonyms) and triggered an in-place reindex. A malformed rule mapped "apple" to a single-token pattern that broke tokenization for 2M products; for 3 hours, queries like "iphone case" returned zero results.

Questions to Surface First:

  • Why could a merchandiser change index-time analysis without review?
  • Why was the change applied in place instead of to a new index version?
  • Why did zero-result rate not page?
  • How fast was rollback possible?

Typical L5 Approach: Revert the synonym file and reindex again — another 3 hours of broken results.

Staff Approach: Immediate: if the previous index version still existed, swap the alias back (seconds). It didn't — so fall back to query-time synonyms disabled plus a restore from the last snapshot. Structural fix: synonyms move to query-time (search-time synonym filters, updatable without reindex); any analyzer change goes through versioned index + parity test; synonym edits via a tool with validation and a test-query suite; search.zero_result_rate > 5% for 5 min pages.

Principal Approach: Merchandising controls are a product with guardrails: a rules console with preview, validation, staged rollout and audit log, owned by the relevance team. Define which knobs are business-owned (boosts, synonyms within limits) and which are engineering-owned (analyzers, mappings).

Staff Approach — Full Reasoning
PhaseWhat to Do
ImmediateRestore last snapshot to new index; alias swap; disable offending synonym.
TriageIdentify affected docs and query classes; estimate revenue impact.
Quick fixQuery-time synonyms; validation on synonym file.
GuardrailsZero-result alert; versioned-index rule for any index-time change.
Post-mortemRoot cause: index-time mutable config without review; no rollback path.

Metrics to Watch: search.zero_result_rate, search.conversion_rate, synonym.rules_changed_total

Organizational Follow-up: Rules console with preview and audit; merchandiser training.

Ownership Question: "Should merchandisers be able to change synonyms?" Staff answer: Yes — query-time synonyms through a validated tool with preview and staged rollout. No — index-time analysis, which belongs to engineering with the versioned-index process.

Key Takeaway: "Anything that changes tokenization is a reindex; anything that's a reindex is a new version behind an alias."

What clears the Staff bar:

  • Distinguishes query-time vs index-time config
  • Rollback path via aliases and snapshots
  • Business controls with guardrails, not bans

Deep Dive 5: Multi-Region Expansion#

Context: Expanding to EU and APAC with p99 ≤ 200 ms locally. EU seller data must stay in the EU; global products are searchable everywhere.

Questions to Surface First:

  • Which documents are global vs region-restricted?
  • Is ranking regional (local popularity, language)?
  • How do index version swaps coordinate across regions?
  • What's the failover story if a region's cluster fails?

Typical L5 Approach: One central cluster with a CDN in front. Cache misses cross oceans (~150 ms RTT) and blow the latency budget; residency is violated.

Staff Approach: Regional clusters fed by regional pipelines: global products replicated via Kafka to all regions; EU-restricted docs indexed only in EU. Regional ranking features (local popularity) computed regionally. Index version swaps rolled out region by region, starting with the smallest. Failover: route a region's traffic to the nearest region with degraded latency and without restricted docs — explicitly accepted by legal and product.

Principal Approach: A 3-region search platform roughly triples cluster cost (~$600K/month at this scale). Decide per region whether a full corpus is needed or a regional subset suffices; standardize the data-classification field on every document so residency is enforced by the pipeline, not per team.

Staff Approach — Full Reasoning
PhaseWhat to Do
Q1Classification field on documents; regional pipelines; EU cluster.
Q2APAC cluster; regional ranking features; per-region A/B capability.
Q3Coordinated version rollout tooling; failover drill.
Ongoingresidency.restricted_docs_outside_region must be 0; audited monthly.

Metrics to Watch: search.p99_ms{region}, index.freshness_lag_s{region}, replication.lag_s, residency.violations_total

Organizational Follow-up: Legal signs off classification scheme; source teams set classification on documents.

Ownership Question: "Who decides EU traffic fails over to the US cluster?" Staff answer: The incident commander per a pre-approved plan — restricted EU docs are excluded from US results, which legal approved in advance.

Key Takeaway: "Search is derived, so regions can be independent. Residency is enforced in the pipeline, not the query."

What clears the Staff bar:

  • Regional independence with classification-aware pipelines
  • Region-by-region version rollouts
  • Pre-approved failover semantics

9. Level Expectations Summary#

After studying this case study, you should be able to:

  • Size a document-partitioned index: shards, replicas and growth headroom, with fan-out tail controls
  • Design a CDC-fed, versioned indexing pipeline with per-field freshness SLOs and end-to-end probes
  • Reindex a live index via versioned indexes, dual-write, parity tests and alias swaps with rollback
  • Build multi-stage ranking within a latency budget and evaluate changes offline and online
  • Define a degradation ladder that trades ranking quality for availability
  • Assign ownership across source teams, search platform and ranking owners

The Bar for This Question#

Mid-level (L4): Explains inverted indexes and TF-IDF/BM25; proposes Elasticsearch. Little on freshness, sharding math or ranking evaluation.

Senior (L5): Designs a working ES-based system with shards, replicas and an indexing pipeline; knows NRT refresh; adds popularity boosts. Reindexing and ranking evaluation handled superficially when asked.

Staff+ (L6): Treats the index as a rebuildable derived store behind aliases; sets per-field freshness SLOs with probes; reasons about fan-out tail latency and replicas vs shards; designs safe reindexing and gated ranking changes; assigns owners for documents, pipelines and relevance. The interviewer should learn something from the answer.


10. Staff Insiders: Controversial Opinions#

10.1 "Your Search Index Should Never Be a Source of Truth"#

If search is authoritativeWhat happens
Data written only to ESCan't rebuild; mapping changes impossible
Price read from search at checkoutStale index becomes a billing bug
No versioningReplays and reorders corrupt documents

The Staff position: Everything in the index must be reproducible from sources. If it isn't, you don't have a search engine — you have a fragile database.

Why this matters in interviews: "How would you rebuild it?" is the fastest way for an interviewer to find out whether the design is sound.

10.2 "Adding Shards Is Usually the Wrong Fix"#

Shards scale data; replicas scale QPS. Every added shard is one more request per query and one more chance to hit a slow node.

The Staff position: Right-size shards for data growth once per index version; scale reads with replicas, caching and cheaper ranking.

10.3 "Vector Search Doesn't Replace BM25"#

Lexical retrieval handles exact terms, SKUs, brand names and rare words better than embeddings; embeddings handle paraphrase and intent. Production systems that dropped lexical retrieval often reintroduced it.

The Staff position: Hybrid retrieval, with a re-ranker deciding. Add vectors when long-tail recall is the measured problem.

10.4 "Relevance Is an Org Problem Disguised as an ML Problem"#

The hardest ranking fights are between sponsored revenue, conversion, seller fairness and user trust. Models optimize what you tell them; someone has to own the objective.

The Staff position: Name the objective owner per surface and make business rules budgeted, measured and experimentally tested.

10.5 "Real-Time Indexing for Everything Is a Tax You Don't Need"#

Most fields don't need second-level freshness. Paying NRT cost for reviews and popularity burns capacity that should serve queries.


11. The Principal Lens (L7)#

Why L7 Sees This Problem Differently#

A Staff engineer builds a great product search. A Principal engineer sees seven search clusters across the company — product search, seller search, help center, internal docs, logs, support tooling, and a new semantic search prototype — each with its own pipeline, ranking approach and on-call, and three of them indexing the same catalog differently. The L7 question is what search capabilities the company provides as a platform, who owns relevance for each surface, and where standardization pays for itself.

The Org-Level Fault Line#

Central search platform vs surface-owned search stacks. A central platform (clusters, pipelines, index versioning, experiment and evaluation tooling) gives consistent operations and shared learning but can become a bottleneck and flatten surface-specific needs. Surface-owned stacks move fast but duplicate expensive infrastructure and produce inconsistent quality and security.

The L7 default: Centralize infrastructure and safety (clusters, pipelines, alias/versioning, ACL injection, evaluation and experiment tooling); decentralize ranking objectives and models to surface teams, within latency and experiment standards. Log search and analytics stay out — different access patterns and economics.

Cost Model#

Assumptions: cloud VMs with local NVMe at ~$1.2K/month per data node (16 vCPU/64 GB/2 TB), 3 copies of data, index ≈ 1.2× raw searchable text, engineer $25K/month fully loaded.

ScaleCorpus / QPSInfra ($/month)HeadcountOn-call LoadNotes
Small5M docs, 300 QPS~$8–15K (or SaaS ~$5–20K)1–2 FTE (shared)~1 page/monthSaaS or managed service; BM25 + simple boosts
Medium50M docs, 5K QPS~$120–220K incl. pipeline + ranking6–10 FTE (platform 3–4, relevance 3–6)~6 pages/monthLearned ranking, CDC pipeline, experiment tooling
Large1B docs, 50K QPS, 3 regions~$1.5–3M30–60 FTEDedicated rotations per layerTiered indexes, hybrid retrieval, per-surface ranking teams

Relevance headcount usually exceeds infrastructure headcount at medium scale — and is where the revenue return lives. A 1% conversion lift on a $1B GMV search funnel is worth more than the entire cluster bill.

The 3-Year Evolution Path#

Diagram: The 3-Year Evolution Path

One-Way Doors vs Two-Way Doors#

DecisionDoorReversal Cost
Index as derived store (vs written directly by apps)One-way once apps depend on itUntangling direct writers takes quarters
Document contract & IDs exposed to other surfacesOne-way-ishEvery consumer migrates
Engine choice (ES/OpenSearch vs Vespa vs SaaS)One-way-ish6–12 months; query DSL rewrite
Shard count per index versionTwo-wayNew version + swap
Analyzers, mappingsTwo-way (with aliases)Reindex; hours
Ranking modelTwo-wayFlag flip
Exposing raw engine queries to product teamsOne-wayEvery caller must move to the API

The Standard I'd Write#

RFC-SEARCH-002: Search Indexes and Relevance Changes

Scope: Every production search index serving user-facing or internal product surfaces (excluding log/observability search).

MUST:

  • Index only data reproducible from a source of truth; full rebuild time documented and tested twice a year.
  • Serve reads through aliases; mapping/analyzer changes ship as new index versions with parity testing and a ≥ 24-h rollback window.
  • Carry a source version on every document; reject out-of-order writes.
  • Publish per-field freshness SLOs and measure them with end-to-end probes.
  • Build queries through the platform query library, which enforces tenant and ACL filters.
  • Ship ranking changes only with an experiment readout including latency, zero-result-rate and business guardrails.

SHOULD:

  • Size shards 10–50 GB; replicas for QPS with ≥ 30% headroom.
  • Implement a degradation ladder with pre-approved steps.

Exceptions: Prototypes may skip experiments for ≤ 5% traffic for ≤ 30 days. Other exceptions reviewed by the Search Platform Council within 5 business days.

Success metrics: Zero ACL-leak incidents; median reindex (mapping change) < 1 day end to end; ≥ 90% of ranking launches with experiment readouts; duplicated catalog indexes reduced from 3 to 1 within a year.

What I'd Tell the VP#

"Search drives about 40% of our revenue, and today four teams run their own search systems, two of which index the same catalog differently. I recommend we consolidate the infrastructure — clusters, data pipelines, safety controls and experiment tooling — into one platform, while each product keeps control of its own ranking. That removes duplicate spend, closes the class of incidents where stale or leaked data showed up in search, and lets relevance improvements ship through a consistent testing process. It needs a platform team of about four engineers, funded partly by retiring two redundant clusters. The bigger return is relevance: every percentage point of search conversion is worth more than the entire infrastructure bill."

Principal Interview Signals#

SignalWhat It Sounds Like
Surface portfolio view"We have seven search stacks; I'd platform the infrastructure and keep ranking with the surfaces."
Prices relevance"1% conversion on $1B GMV is $10M/year — that's why relevance headcount outnumbers cluster headcount."
Names one-way doors"Letting apps query the engine directly is the one-way door; I'd force the query library now."
Governs changes"Merchandisers get query-time controls with preview and audit; analyzers stay engineering-owned."
Knows what not to centralize"Log search has different economics — it stays on the observability stack."

Staff answers that L7 interviewers find insufficient:

  • A superb single-surface design with no view of the other six search stacks or shared catalog indexing.
  • "We'll add vector search" without cost per query, recall targets, or who owns the embedding models across surfaces.
  • A ranking experiment process without an owner for the objective when conversion and sponsored revenue conflict.

Appendices

Appendix A: Index and Retrieval Mechanics

A.1 Inverted Index#

term → postings: [(doc_id, tf, [positions]) ...]   sorted by doc_id, delta + block encoded
query "red running shoes": intersect postings(red) ∩ postings(running) ∩ postings(shoe)
                           using skip lists; score survivors; keep top-K heap per shard

A.2 BM25 (for reference)#

score(q,d) = Σ_t IDF(t) · tf(t,d)·(k1+1) / (tf(t,d) + k1·(1 − b + b·|d|/avgdl))
k1 ≈ 1.2, b ≈ 0.75

IDF is per-shard by default; with skewed shards, scores differ across shards — use DFS-style global stats only when it matters (it rarely does with random doc routing).

A.3 Segments, Refresh, Merge#

Writes → in-memory buffer + translog → refresh → new segment (searchable) → flush/commit (durable on disk) → merges (tiered policy) purge deletes.

A.4 Vector Retrieval#

HNSW: ~1–2 KB per 384–768-dim vector plus graph; memory-resident for low latency; recall@10 ~0.9–0.98 tuned by ef_search. Quantization cuts memory 4–8× at small recall cost.

A.5 Hybrid Merge#

Reciprocal rank fusion: score = Σ 1/(k + rank_i), k ≈ 60 — simple, robust, no score calibration needed; then re-rank.

Appendix B: Document Model and Keys

B.1 Search Document#

{ product_id, version (source LSN), title (text + keyword), description (text, no positions if unused),
  brand_id (keyword), category_path (keyword, hierarchical), price (scaled_float), in_stock (bool),
  seller_id (keyword), popularity (float, nightly), region_class (keyword), acl? (keyword[]) }

B.2 Mapping Hygiene#

Disable indexing for display-only fields; avoid dynamic mappings (field explosion); keep _source lean or exclude large fields.

B.3 Deletes and Erasure#

Deletes are tombstones until merge; GDPR erasure requires delete + forced merge on affected segments or reliance on retention of segments bounded by merge policy — document the maximum time-to-erase (e.g., ≤ 30 days) with legal.

Appendix C: Indexing Pipeline Mechanisms

C.1 Comparison#

MechanismFreshnessCostFailure Mode
Dual-write from appSecondsLowInconsistency on partial failure; app coupled to search
CDC → Kafka → indexerSecondsMediumStuck partitions; needs enrichment joins
Periodic batch (updated_at polling)MinutesLowMisses deletes; clock issues
Offline build + snapshot shipHoursHigh build, cheap servingFreshness needs a separate tier

C.2 Idempotent Writes#

External versioning: index(doc, version=lsn, version_type=external) — writes with lower versions are rejected, so replays and reorders are safe.

C.3 Priority Lanes#

Lane 1 (price/stock/ACL, partial updates) → Lane 2 (content edits) → Lane 3 (bulk/backfill, rate-limited). Separate consumer groups and indexer quotas.

Appendix D: Query API Contract and Client Behavior

D.1 Pagination#

Cursor-based (search_after with a sort tiebreaker); from + size capped at 1,000–10,000.

D.2 Partial Results#

Response includes partial: true and shards_failed when a shard times out; clients render results and may log for quality.

D.3 Caching#

Result cache keyed by normalized query + filters + ranking version (not user) for non-personalized traffic; TTL 60 s; invalidation by TTL only — freshness SLO accounts for it.

Appendix E: Observability

E.1 Core Metrics#

Query:     search.qps, search.p50_ms, search.p99_ms, search.error_rate, search.partial_rate,
           search.cache_hit_rate, rerank.latency_ms, rerank.skipped_total
Index:     index.freshness_lag_s (probe), consumer_lag{partition}, index.bulk_rejections,
           index.merge_time_ms, index.segment_count, shard.size_gb
Quality:   search.zero_result_rate, search.ctr_at_10, search.conversion_rate, search.reformulation_rate,
           feature.null_rate{feature}
Cluster:   cluster.status, jvm.heap_pct, gc.pause_ms, shard.rejected_requests

E.2 Critical Alerts#

AlertThresholdSeverity
Latencysearch.p99_ms > 300 for 5 minPage
Freshnessprobe p99 > 30 s for 5 minPage
Zero results> 5% for 5 min (baseline ~2%)Page
Cluster redany red indexPage
Heap> 85% for 10 minTicket → page at 92%
Partition lagany partition > 60 sPage

E.3 Debugging "Search Is Slow"#

  1. Uniform or shard-skewed? 2) Merges/indexing spike? 3) Heavy queries (deep pagination, wildcards, huge aggregations)? 4) Re-ranker latency? 5) GC? Check per-shard latency first — it partitions the problem in one query.
Appendix F: Scale Evolution

F.1 What Works at Each Scale#

ScaleApproach
< 1M docsPostgreSQL full-text or SaaS search
1M–100M docsManaged ES/OpenSearch, CDC pipeline, aliases, learned re-ranking
100M–5B docsTiered indexes (fresh + base), offline builds, dedicated ingest nodes, hybrid retrieval
Multi-regionRegional clusters from replicated pipelines, classification-aware indexing

F.2 What You Don't Build on Day One#

  • A custom search engine
  • Term-partitioned indexes
  • Vector search before measuring long-tail recall gaps
  • Per-user personalization before non-personalized ranking is measured well
Appendix G: Multi-Tenancy, Fairness and Cost

G.1 Tenant Placement#

Tenant sizePlacementIsolation
< 100K docsShared index, routing by tenantFilter + routing
100K–5M docsShared index, routing partitionsFilter + quotas
> 5M docs or regulatedDedicated index (or cluster)Physical

G.2 Query Fairness#

Per-tenant QPS quotas and query cost limits (max clauses, timeouts); expensive query classes (wildcards, deep aggregations) on separate thread pools or clusters.

G.3 Showback#

Cost per surface = share of QPS × query cost + share of index storage; published monthly.

  1. Loading the index…