Dashboard
0%
1
Curious builder0 XP earned · 300 to level 2
0 daysFinish a lesson to begin
Badge collection0 of 6 unlocked
51 small wins to finish your pathNext question →

Q38HardSystem design

How would you scale RAG to 100 million+ documents (billions of chunks)?

30-second answerSay your answer out loud first, then reveal.
Offline, distributed batch embedding feeds an index build per shard; at query time a router sends the query to shards 1 to N, each returns top-k from a quantized ANN index, and the merged candidates are rescored with full-precision vectors and reranked by a cross-encoder before the LLM.

Design areas

  1. Storage math: 2B chunks × 768 dims × 4 bytes ≈ 6 TB of raw vectors. In RAM, that's infeasible economically, so you need quantization (int8 → 1.5 TB, PQ or binary → much smaller) plus disk-based ANN or tiered memory.
  2. Sharding strategies:
    • By tenant / domain: queries hit one shard (best when filters are natural).
    • By hash: even load, but every query fans out to all shards (scatter-gather), so tail latency matters.
    • By time: recent data hot, old data cold.
  3. Two-tier retrieval: a cheap first stage (BM25 or binary vectors) over everything → full-precision rescoring → cross-encoder on the top 100.
  4. Indexing throughput: distributed embedding (GPU batch jobs), checkpointed pipelines, back-pressure. Re-embedding the whole corpus takes days and costs real money, so plan model migrations carefully (Q45).
  5. Freshness: a small, fast "delta" index for recent documents merged with the large static index; periodic compaction.
  6. Replication for QPS and availability; consistent snapshots.
  7. Document-level vs chunk-level retrieval: first retrieve documents (using summaries), then chunks within the top documents. This reduces the search space.

Common mistakes

  • Designing one giant HNSW index in RAM without doing the memory math.

Slow is fine. Stopping is the only problem.