When do you shard a FAISS index across GPUs instead of replicating it?
answer
- two different scaling axes
- one flag on the cloner options
- full copy versus disjoint slice
- every query touches every slice
- losing one is not symmetric
basics
~20 sReplicate when the index fits on one GPU and you need more queries per second — each device holds a full copy and queries are split between them. Shard when the index does not fit on one GPU: each device holds a slice, every query touches all of them, and results are merged.
solid answer
~50 s`faiss.index_cpu_to_all_gpus(cpu_index, co=co, ngpu=n)` gives you both shapes, chosen by `co.shard` on a `GpuMultipleClonerOptions`. The default, `shard = False`, **replicates**: every GPU holds the whole index, incoming queries are divided across devices, and throughput scales roughly linearly while capacity does not — you still cannot index more vectors than one device holds. Setting `co.shard = True` **shards**: each GPU stores a disjoint slice, every query is broadcast to all devices, each returns its local top-k, and FAISS merges them into a global top-k. That multiplies capacity but not throughput, and each query's latency is bounded by the slowest shard. The decision is therefore mechanical: if the index does not fit on one device, you must shard; if it fits and you are QPS-bound, replicate. Large deployments combine them — shard to fit, then replicate the shard set to scale.
code
python · 8 linesimport faiss
co = faiss.GpuMultipleClonerOptions()
co.shard = True # False (default) replicates a full copy per GPU
co.useFloat16 = True
index = faiss.index_cpu_to_all_gpus(cpu_index, co=co, ngpu=4)
D, I = index.search(queries, 10) # merged global top-k either waygo deeper
Know the distinction in one line: replicas are full copies that raise throughput, shards are slices that raise capacity, and it is a single option on the multi-GPU clone call.
Explain the mechanics — how a query is routed in each shape, that sharded queries fan out and merge per-shard top-k results, and why memory usage differs by a factor of the device count.
Show the operating judgment: pick the shape from what you are short of, compose both at scale, re-measure recall after sharding, and account for a lost shard silently degrading results rather than just capacity.
Own the capacity plan — device count derived from index size and QPS target together, redundancy for shards, and the cost comparison against compressing the index to fit fewer devices.
## Two different scaling axes Multi-GPU FAISS solves two problems that people routinely conflate. **Replication** scales *throughput*: N copies of the same index, queries distributed among them, roughly N times the queries per second. **Sharding** scales *capacity*: one index split into N disjoint pieces, so you can hold N times as many vectors as one device fits. Choosing the wrong one produces a system that is fast and too small, or large and no faster. ## The API ```python co = faiss.GpuMultipleClonerOptions() co.shard = True # False (default) = replicate co.useFloat16 = True index = faiss.index_cpu_to_all_gpus(cpu_index, co=co, ngpu=4) ``` The result behaves like an ordinary FAISS index — `search` returns a merged global top-k either way — so the difference is invisible in the calling code and entirely visible in the resource profile. Underneath, the replicated form is an index that fans queries out to identical copies, and the sharded form is one that broadcasts to slices and merges their results. ## What replication costs and buys Each GPU holds the complete index, so device memory usage is N times the index size and the maximum dataset is still whatever one device holds. In exchange, a batch of queries is split N ways and each device works on its slice independently — near-linear throughput scaling, and no cross-device coordination. One nuance matters operationally: a replicated setup only scales throughput if queries arrive in batches large enough to divide meaningfully. Sending one query at a time to a four-way replicated index leaves three GPUs idle, so a batching layer in front is part of the design, not an optimisation. ## What sharding costs and buys Each GPU holds roughly `n/N` vectors, so the aggregate dataset scales with device count. Every query, however, must visit every shard — each computes its own top-k and the results are merged. Consequences: - **Throughput does not scale.** All devices are busy on every query rather than working on different queries. - **Latency is set by the slowest shard.** Uneven shard sizes or a device shared with another workload drags the whole query. - **A merge step is added** to combine per-shard results into the final top-k. Cheap relative to the search, but real. ## Combining them At scale the two compose: shard until the index fits across a group of devices, then replicate that group as many times as your QPS demands. A 400 GB index at 5,000 QPS is not a choice between shapes — it is four-way sharding for capacity, replicated three times for throughput, twelve devices in total. Framing the answer this way is what distinguishes an engineer who has operated multi-GPU search from one who has read the flag list. ## Correctness and recall notes Sharding an IVF index means each shard has its own inverted lists over its own slice of the data, so the effective number of lists probed per query is per shard. Recall is not automatically identical to the unsharded index at the same probe setting, and a recall measurement should be re-run against the sharded configuration rather than inherited from a single-device benchmark. Ids also need care: the merged results must return globally meaningful ids, so the ids you add must be unique across the whole dataset, not per shard. ## Operational shape Replicas fail independently — losing one costs you throughput, and the remaining copies still answer every query correctly. Shards do not: losing one shard silently removes a slice of the corpus from every result. That asymmetry is the strongest operational argument for preferring replication whenever the index genuinely fits, and for treating a sharded deployment as needing redundancy of its own. ## The interview answer in one line Shard for memory, replicate for speed; check which one you are actually short of before reaching for the flag, and remember that a lost shard degrades correctness while a lost replica only degrades capacity.
- With a four-way replicated FAISS index, why might you still see no throughput gain?Replication splits a batch of queries across devices, so it only helps if the batch is large enough to divide. A stream of single queries keeps one GPU busy and three idle, and per-call overheads dominate. A batching layer that accumulates queries for a few milliseconds before dispatching is what converts replication into real throughput.
- How does sharding change how you measure recall?Each shard runs its own inverted-list search over its own slice, so probing a given number of lists per shard is not equivalent to probing that many over the whole dataset. Recall must be re-measured against the sharded configuration at the probe setting you will actually serve with; carrying over a single-device benchmark number is how sharded deployments quietly ship worse results.
- What is the failure-mode difference between losing a replica and losing a shard?Losing a replica costs throughput only — the surviving copies still contain the entire dataset and every query is still answered correctly. Losing a shard silently removes part of the corpus from every result: queries succeed, latency looks fine, and recall degrades invisibly. Sharded deployments therefore need their own redundancy and a health check that fails loudly rather than degrading.
saying these in an interview costs you the question
- Thinking sharding increases queries per second
- Believing replication lets you index more vectors than one GPU holds
- Assuming a sharded query only hits one GPU
- Expecting replication to help without batching queries
- Treating a lost shard as merely a capacity loss rather than a recall loss