Distributed retrieval and sharding

Split the index across machines and query them all: scatter-gather. Each shard returns its own top k and a coordinator merges them. Two things decide whether it works: how you shard, and the fact that your latency is now the slowest shard's, not the average.

Overview

Shards and replicas are different things

A shard holds a slice of the corpus. Sharding handles data that will not fit — memory or index build time — and every query must visit every shard.

A replica holds a complete copy of a shard. Replication handles query volume and failure, and a query goes to one replica of each shard.

They solve different problems and are often confused. If the index fits on one machine and you are simply serving too many queries, you need replicas and no sharding at all — which is much simpler, because there is nothing to merge.

Parameters

Visualisation

—

Readout

What to watch

  • Every shard is queried; the coordinator merges local top-k lists.
  • Random sharding spreads relevant documents evenly — that is the point.
  • Latency is the slowest shard, not the average.

Distributed retrieval and sharding: A Practical Guide

How do you scale retrieval past one machine? What are the trade-offs between sharding strategies?

Shard randomly

The instinct is to shard by topic or tenant so a query only touches one shard. For a multi-tenant system where every query is scoped to one tenant, that is right — it is really many small indexes.

For a general corpus it is a mistake. If all the finance documents live on shard 2, a finance query's true top 50 are all there, and taking only the local top 5 discards forty-five of them. Random sharding spreads relevant documents evenly so a modest per-shard k captures nearly all of them.

The fix if you must shard semantically is to over-fetch from every shard, which costs the bandwidth and latency the routing was meant to save.

Tail latency, and the merge

Scatter-gather waits for the slowest shard, so p99 response time is roughly the p99 of any shard. Add shards and the chance one is slow rises — the classic result is that latency gets worse as you scale out, not better.

Mitigations: hedged requests (ask two replicas, take the first answer), a deadline after which a late shard is dropped and the degraded result served, and keeping shard counts modest.

The merge itself needs care. Scores must be comparable across shards — fine for cosine similarity, not fine for BM25, whose IDF depends on corpus statistics that differ per shard unless they are shared globally. That is a real bug in hand-rolled distributed lexical search.

The merge, and the scores that do not survive it

Merging local top-k lists assumes the scores are comparable across shards, and that assumption quietly fails for lexical retrieval.

BM25's idf term depends on corpus statistics — document frequency and average document length — which differ per shard. A term that is rare on shard 1 and common on shard 3 gets different weights, so a merged ranking is comparing numbers computed on different scales. The fix is global statistics: compute df and avgdl across the whole corpus and distribute them, which is an extra coordination step people discover only after the rankings look wrong.

Dense retrieval escapes this. Cosine similarity between a query and a document vector involves no corpus statistics at all, so shard-local scores are directly comparable. It is one of the underrated operational advantages of vector search.

Replication, failure and the degraded answer

Sharding splits data; replication copies it. They solve different problems, and conflating them is the most common design error here: if the index fits on one machine and you are simply serving too many queries, you want replicas and no sharding at all, and there is then nothing to merge.

Once shards exist, so does partial failure. Scatter-gather waits for everyone, so a single slow or dead shard degrades every query. The practical answers are a deadline — serve what returned in time and mark the result degraded — and hedged requests, asking two replicas and taking whichever answers first, which trades a few percent extra load for a much better tail.

Decide explicitly whether a degraded answer is acceptable. For search it usually is. For a RAG answer that will be presented as authoritative, silently dropping a shard means silently dropping evidence, and the user has no way to know.

When one machine is not enough

A vector index has to fit in memory to be fast. That gives a hard ceiling per machine:

memory ≈ vectors × dimensions × bytes × index overhead

For 100 million 768-dimensional vectors at 4 bytes each, that is 307GB before the HNSW graph's 1.5–2× overhead — so 450–600GB. Beyond a single machine.

Three separate reasons drive distribution, and they want different architectures:

Capacity. The index does not fit. Split it — sharding.

Throughput. Too many queries per second. Copy it — replication.

Availability. A machine failure must not take the system down. Replication again.

Most production systems need both: shards for capacity, replicas of each shard for throughput and resilience.

Sharding, and how to split

StrategySplit byQuery cost
RandomHash of the idEvery shard, every query
SemanticCluster of the vectorOne or a few shards
MetadataTenant, language, dateOnly matching shards
TimeIngestion periodRecent shards first

Random sharding distributes evenly and requires querying every shard for every request, then merging. Simple, predictable, and the cost grows with shard count.

Metadata sharding is usually the best choice when the data has a natural partition. Shard by tenant, and each query touches one shard. Shard by language, and a query in French touches the French shard. The saving is proportional to the selectivity.

Semantic sharding clusters similar vectors together so a query only needs the nearest clusters. Attractive in principle, and it risks recall loss when the answer sits in an unqueried shard, and it needs rebalancing as the distribution shifts.

Time-based sharding suits corpora where recency matters: query the recent shards first, and older ones only if needed. It also makes retention easy — drop an old shard.

Scatter-gather, and the detail that breaks it

With random or semantic sharding, a query fans out and the results are merged:

  1. Send the query to all (or the selected) shards.
  2. Each shard returns its local top k.
  3. Merge, sort by score, take the global top k.

The correctness condition is easy to get wrong: each shard must return its own top k, not top k/n. If the global top 10 all happen to live in one shard, asking each of five shards for 2 results loses 8 of them.

The latency condition is different and equally important: the query takes as long as the slowest shard. With ten shards, the chance that at least one is having a slow moment is ten times higher than with one — so tail latency degrades as shard count grows. This is the standard fan-out problem, and the standard mitigations are hedged requests (send to a replica after a short delay and take whichever answers first) and a timeout with partial results.

Tail latency, and the scores that do not survive the merge

Scatter-gather looks simple until you measure it. Two things bite: a query is only as fast as the slowest shard it waited for, and the per-shard scores you merge are often not comparable. Both are arithmetic, and both are here.

example_01.pyNumPy
Output

Things to try

  1. Switch to semantic sharding. Recall collapses: every relevant document is on one shard, whose local top-k discards the rest.
  2. With semantic sharding, raise per-shard k until recall recovers. You have paid back the bandwidth the topic routing was meant to save.
  3. Raise shards to 40 with random sharding. Recall is fine and p99 latency climbs, because scatter-gather waits for whichever shard was slow.

What to remember

Scatter-gather sends the query to every shard, takes a local top-k from each and merges. Shard randomly rather than by topic, or the relevant documents concentrate in one shard and its local k throws most of them away. Latency becomes the slowest shard's, not the average, so p99 gets worse with every shard added. Shards solve data that will not fit; replicas solve query volume.

Replication

Copying each shard serves more queries and survives failures.

AspectEffect
Read throughputScales roughly linearly with replicas
AvailabilitySurvives replica loss
Write costEvery write goes to every replica
ConsistencyReplicas can lag behind

That last row matters for retrieval specifically. A document indexed a moment ago may be present on one replica and not another, so two identical queries can return different results. For most retrieval that is acceptable; where it is not — a user searching for something they just uploaded — the usual fix is to route that user's reads to the replica that accepted the write, or to accept the lag and say so in the interface.

The routing choice is between load balancing (any replica, best throughput) and consistent routing (the same replica for a session, better cache behaviour and freshness). Consistent routing also improves cache hit rates, which is a meaningful secondary benefit.

The metadata and filtering problem

Distribution interacts badly with filtering unless it is planned for.

A selective filter across random shards means every shard is queried and most return nothing. The work is done and thrown away. Sharding by the filtered attribute instead makes the filter a routing decision.

Permissions must be enforced per shard, inside each shard's search, not merged afterwards. Filtering after the gather step means unauthorised content has already crossed the boundary and has already occupied result slots.

Global statistics break. BM25's IDF term depends on how many documents in the corpus contain a term, and each shard only knows its own counts. With random sharding the local approximation is usually close enough; with metadata sharding it can be badly skewed, and correcting it requires distributing global term statistics.

That last point is a real and under-discussed source of ranking inconsistency in sharded hybrid search.

Operational realities

Rebalancing. Adding a shard means moving data. Consistent hashing limits how much moves; semantic sharding may require re-clustering the whole index.

Rebuilds. Changing the embedding model invalidates everything. Build into a new set of shards and switch atomically rather than migrating in place.

Monitoring per shard. Aggregate latency hides a single slow shard. Track per-shard latency, size and error rate.

Deletions. Tombstones accumulate per shard, so schedule compactions.

Cost. Memory-resident indexes are the expensive part. Reducing dimensions — 1,536 to 384 — is a fourfold saving and often a small accuracy cost, and it is usually a better first move than adding machines.

That last point is worth trying before distributing at all: quantisation and dimension reduction frequently keep a corpus on one machine, which removes every problem above.

Questions people ask

When do I need to shard? When the index no longer fits in one machine's memory after quantisation and dimension reduction. Try those first.

How many shards? As few as fit. Each additional shard adds fan-out cost and tail latency.

Should each shard return k or k/n? k. Returning k/n loses results when relevance is concentrated.

Does sharding hurt recall? Random sharding with per-shard top-k does not. Semantic sharding can, if the answer is in an unqueried shard.

How do I handle multi-tenancy? Shard by tenant where volumes allow — it gives natural isolation and cheap filtering. Very many small tenants are better served by one index with a tenant field.

What about BM25 across shards? IDF becomes shard-local. Acceptable with random sharding; distribute global term statistics if the skew matters.

Recap in one screen

  • Shard for capacity, replicate for throughput and availability — most systems need both.
  • Shard by metadata where a natural partition exists, so a query touches one shard instead of all of them.
  • In scatter-gather, each shard must return its own top k, and latency is set by the slowest shard.
  • Enforce permission filters inside each shard's search, never after the merge.
  • Try quantisation and fewer dimensions first — staying on one machine avoids every problem here.

Recall check

0 of 3

Say the answer out loud before you reveal it — recalling it is what makes it stick, and rereading it is not.

  1. Without scrolling back — what is the one-line takeaway from this module?

  2. What does this module say about “Shards and replicas are different things”?

  3. What does this module say about “Shard randomly”?

Cheat sheet

Distributed retrieval and sharding

Split the index across machines and query them all: scatter-gather. Each shard returns its own top k and a coordinator merges them. Two things decide whether it works: how you shard, and the fact that your latency is now the slowest shard's, not the average.

GEN AI · vizlearn.in/gen_ai/distributed_retrieval_and_sharding.html

About the author

Ashish Jangra builds and maintains VizLearn. Every module here is written and the visualisation behind it hand-built, so the numbers in a readout come from the same code that draws the picture. Corrections are genuinely welcome and get priority over everything else — if a page states something wrong, or an animation misrepresents what the algorithm does, get in touch.