Scenario-Based AI Engineering Questions

Course Content

Scenario-Based AI Engineering Questions

26 sections · 146 lessons

80% of queries hit 5% of documents. Your vector DB has uniform sharding — and the hot shards keep crashing. How do you design retrieval infrastructure for skewed query distributions?


Where each match-night query stopsExact and semantic cache — about 55 percentHot tier: last 14 days, 6 replicas in RAMCold tier: the other 95 percent of articlesAdmission control and one capped retry
Adding equal shards never helps when the popular documents live together; tiers and caches move the load away from them.

What you need to know

Why hot shards appear

Sharding splits an index across machines. How queries reach shards decides where skew hurts:

  • Scatter-gather sends every query to every shard. Skew in which documents are popular does not make one shard hotter, but every shard pays for every query.
  • Routed sharding (by tenant, category, date or region) sends a query only to the shards that own the data. If the popular documents sit together — this week's news, one big tenant, one product category — those shards take most of the traffic.

With 80% of queries on 5% of documents, a routed shard holding that 5% can receive many times its fair share. Adding more equal shards does not help, because the popular data stays together.

The fixes, in order

  1. Cache the head — an exact-match cache on normalised queries, then a semantic cache (return a stored result when a new query's embedding is very close to a past one). Invalidate entries when their documents change.
  2. Replicate hot shards — the load is reads, so add replicas and balance across them. Splitting a hot shard only helps if the skew is spread inside it.
  3. Tier by popularity — keep the hot 5% in a small index fully in memory with many replicas; the cold 95% sits on a cheaper tier. Query the hot tier first, and fall through when results are weak.
  4. Shed load safely — per-tenant rate limits, a bounded queue and timeouts, so overload means slower answers, not crashes. Cap client retries.
FixRemoves load fromWatch out for
CacheEverything behind itStale results; set short TTLs and invalidate on update
ReplicasThe hot shardCost; replicas must stay in sync
Hot tierCold tier and big shardsKeeping the hot set current as popularity moves
Admission controlCascading failureSome requests get slower or rejected

A semantic cache check:

Python
def cached_search(query: str, threshold: float = 0.95):    q = embed(query)    hit = cache_index.search(q, limit=1)              # small in-memory index of past queries    if hit and hit[0].score >= threshold and not stale(hit[0].payload):        return hit[0].payload["results"]    results = tiered_search(q)                        # hot tier first, then cold    cache_index.upsert(vector=q, payload={"results": results, "doc_ids": [r.id for r in results]})    return results

Set the threshold carefully: too low and "refund policy for gold members" gets the answer cached for "refund policy for silver members". Test it on real query pairs.

Make hotness a measured thing

Popularity moves. Track queries per second per document and per shard, and rebuild the hot set daily (or hourly during events). The metrics to manage are p99 latency per shard and the load-imbalance ratio (busiest shard's load divided by the average).

A real-life example

Scenario, numbers made up. A cricket news app runs semantic search over 20M articles on 8 shards, routed by publish month. On match nights, 80% of queries are about the current series, all stored in one shard. That shard hits 100% CPU, crashes, and client retries knock over its replica.

The team adds an exact plus semantic cache (hit rate about 55% on match nights), moves the last 14 days of articles into a hot tier with six in-memory replicas, and caps retries at one with a 300 ms timeout. During the next final, p99 search latency stays around 180 ms instead of timing out, and no shard crashes. The hot tier is rebuilt every hour during tournaments.

Follow-up questions to expect

  • "Why not just re-shard by hash?" — Hash sharding spreads hot documents across shards, which fixes placement skew, but you then lose routing, so every query fans out to every shard. It is a valid choice when most queries are not filtered.
  • "How do you keep the cache correct when documents change?" — Store document IDs with each cached result and invalidate every entry that contains an updated document, plus a short TTL as a safety net.
  • "What if the skew changes daily?" — Compute the hot set from recent traffic on a schedule, and let the cache absorb sudden spikes before the tier catches up.