Sharding, Updates, and Stale Embeddings — Production Vector Search
A vector index that works on a laptop fails in production in predictable ways. It outgrows one node, scatter-gather inflates p99, tombstones erode recall, and embeddings go stale when the text or the model changes. Mixing two embedding models fails silently: no error, just garbage neighbors. Hash-shard, return full k per shard, and cut over models with a second index and an alias.
- 1Gist
- 2Maps
- 3Q&A
- 4Sandbox
Voice readout needs Web Speech Synthesis in this browser.
Question ladder
L1
Why must each shard return the full k?
Answer
A shard may hold several of the true global top k. If each shard returns only k divided by the shard count, the merge drops true neighbors and recall falls.
L2
Hash sharding or semantic routing?
Answer
Hash sharding hits every shard and stays balanced. Semantic routing probes the nearest shard centroids, saves fan-out, misses boundary neighbors, and creates hot topics. Most systems accept the fan-out.
L3
What is a tombstone in a vector index?
Answer
The vector stays in the graph or the list and is filtered at query time. Without oversampling, the result count shrinks as deletes pile up. Past roughly 20 to 30 percent, compact or rebuild.
L4
How do you upgrade an embedding model without downtime?
Answer
Build a second index, dual-write new documents with the new model, backfill old documents, shadow-query and compare, flip an alias, keep the old index for rollback, then delete it. Never mix the two vector spaces.
L5
Recall dropped 10 points in a month with no deploy. What do you check?
Answer
Deleted ratio, IVF list skew, Lucene segment count, a new content type outside the training distribution, and the model id distribution. Then compact or rebuild.
L6
How do you know an embedding is stale?
Answer
Store model id, model version, content hash, and embedded-at on every vector. A periodic job compares the hash to the source. Alert on mismatches, on lag, and on the share of vectors from a deprecated model.
L7
Why might two replicas disagree?
Answer
HNSW builds are order-dependent, and segment layouts differ. Approximate results can differ at the margin. Pin a session to a replica if the product needs stability, or accept the jitter and measure recall per replica.
Failure modes
Per-shard k too small
Shipping k divided by the shard count drops global recall. The safe default is k per shard, then a merge.
Tombstones past 30 percent
Oversampling hides the short list. Every query still walks dead nodes, and the links around them rot.
Mixed model spaces
A v2 query on a v1 index is near-random and raises no error. Refuse the model id at the API.
Fusing a partial backfill
Rank-fusing an incomplete new index with the complete old one is worse than serving the old index alone.
Misconceptions
Semantic sharding removes fan-out for free.
It is IVF at cluster scale. Boundary queries need several shards, hot topics make hot shards, and drift forces a re-shard.
Dual-read during a model migration is always safer.
It pays when new documents exist only in the new index. Fusing a partial re-embedding of the old corpus with the old index lowers recall.
A snapshot of the index is a freshness guarantee.
Restoring an index built with a retired model is a staleness bug waiting to happen. The model id has to travel with the bytes.
Interviewer traps
Re-teaching generic Elasticsearch shard colors.
Point at the sharding lesson for health and routing, then say what is different for kNN: per-shard k, per-segment graphs, and model ids.
Designing a consistent-hash ring for the vectors.
The partition-key lesson owns rebalancing. Here the question is fan-out versus recall, and whether the embedding space is still the one you queried.
Design scenario
Same prompt for every reader.
Requirements
Recall stays measured on a frozen query set. Queries from the new model never hit the old index. Rollback is an alias flip.
Traffic / scale
2 billion vectors, online queries, a steady update stream, and a backfill that takes long enough to overlap with new writes.
Latency
p99 is roughly the slowest shard. Fewer, larger shards beat many tiny HNSW graphs.
Consistency
One model id per index. Content hash and embedded-at sit on every vector.
Availability
The live alias stays on the old index until backfill is complete and the shadow eval passes. The old index remains for a rollback window.
Failure assumptions
- A writer embeds new text with v2 and inserts it into the v1 index.
- Deletes reach 40 percent before anyone compacts.
- A restore brings back vectors from a retired model.
Constraints
- Do the memory math before the shard count.
- Do not fuse partial v2 results with v1 as the user-facing path.
Prompt
Run a 2 billion vector index through a model upgrade without downtime, on nodes that cannot hold fp32 HNSW.
API
What does the query service check before it embeds, and what does the alias point at?
Data
Which fields sit next to the vector, and where do tombstones live?
Architecture
How many shards does a query touch, and how do you cut back to v1?
How a query meets the shards
Prefer
Hash sharding, full k per shard
Every query touches every shard. Each shard returns k, not k divided by the shard count. The coordinator merges by score. Load stays even.
- A shard can hold several of the true global neighbors.
- Fewer, larger shards usually beat many tiny HNSW graphs.
- Replicas add QPS. They do not change the fan-out of a single query.
Alternative
Semantic routing to a few shards
Send the query to the nearest shard centroids. You skip most of the fan-out, and you miss neighbors that fell across a boundary unless you probe more shards.
- Hot topics become hot shards.
- Drift forces a re-shard, the same way IVF centroids go stale.
- Use it when fan-out cost dominates and you can live with the skew.
A model upgrade that does not mix spaces
Serve the old index until the new one is complete. Shadow reads measure recall. The alias moves once.
- 1
Keep writing v1
The live alias still serves the old embedding space. New documents are dual-written into the new index with v2 vectors. - 2
Backfill in batches
Old documents are re-embedded with v2. A partial green index is not a user-facing result list. - 3
Shadow, then flip
Compare v2 to the exact oracle and to labeled relevance. Flip the alias at full backfill. Keep v1 for the rollback window, then delete it. - 4
Refuse a mismatch
Store model id and content hash beside every vector. A query from a different model id does not enter the index.
Overview
A vector index that works on a laptop fails in production in predictable ways. It outgrows one node and must be sharded. Scatter-gather fan-out inflates p99. Deletes and updates leave tombstones that erode recall. Segment merges and rebuilds compete with queries. Most dangerously, embeddings go stale when the text or the model changes. Mixing vectors from two models fails silently: no error, just garbage neighbors.
This page is what is different for vectors. Cluster health, replicas, and routing as a search-engine problem live on Sharding, Replicas, Routing & Cluster Health. How you pick a shard key and rebalance a relational store lives on Database Sharding & Partitioning — Keys, Hotspots & Rebalancing. Bulk, rollover, and snapshots as ingest live on Indexing Pipelines, Bulk, ILM & Snapshots.
Sharding
Flow
- 1
1. Query vector
- next2. Coordinator
- 2
2. Coordinator
- scatter3a. Shard 1 local top k
- scatter3b. Other shards
- 3
3a. Shard 1 local top k
- next4. Merge to global top k
- 4
3b. Other shards
- next3c. Shard 2 local top k
- 5
3c. Shard 2 local top k
- next3d. Shard 3 local top k
- 6
3d. Shard 3 local top k
- next4. Merge to global top k
- 7
4. Merge to global top k
- next5. Re-rank and fetch docs
- 8
5. Re-rank and fetch docs
Lesson map
Sharding, Updates, and Stale Embeddings — Production Vector Search
A vector index that works on a laptop fails in production in predictable ways. It outgrows one node, scatter-gather inflates p99, tombstones erode recall, and embeddings go stale when the text or the model changes. Mixing two embedding models fails silently: no error, just garbage neighbors. Hash-shard, return full k per shard, and cut over models with a second index and an alias.
Architecture. Architecture
Select a node to see why it exists, or an edge to see the protocol, direction, effect, and consequence.
Mermaid export
flowchart TB q["1. Query vector"] c["2. Coordinator"] s1["3a. Shard 1 local top"] rest["3b. Other shards"] q -->|1. Query vector to 2. Coordinator| c c -->|scatter| s1 c -->|scatter| rest
| Strategy | Routing | Fan-out | Recall risk | Load balance | Notes |
|---|---|---|---|---|---|
| Hash | Document id hash | All shards, every query | Low if each shard returns full k | Even | Default in Elasticsearch, OpenSearch, and most vector databases |
| Semantic | Nearest shard centroids | Only the top few shards | Neighbors across a boundary are missed unless you probe more | Skewed. Hot topics make hot shards | IVF at cluster scale |
| Tenant or partition key | Tenant id | One shard, or one group | None inside the tenant | Skewed by tenant size | Natural for SaaS. Pairs with the filter page |
| Replicas | Any copy | Unchanged | None | Adds QPS | Each replica holds a full copy in RAM |
Rules:
- With hash sharding, each shard returns the full k, or more, not k divided by the shard count. One shard may hold several of the true top k.
num_candidatesin Elasticsearch is per shard. Total work scales with shard count. Fewer, larger shards usually win for kNN, because each HNSW search has a fixed overhead and each segment is searched on its own.- The p99 of the query is roughly the max over shards, so more shards means more stragglers. Replicas and adaptive replica selection help.
The sandbox's seeded hash layout, 8 shards, shows per-shard k of 2 (16 candidates shipped for a global k of 10) at about 0.88 recall. Per-shard k of 5 restores it here. The safe default is still k per shard. Semantic routing to 2 of 8 shards keeps about 0.97 recall and skips 75 percent of the fan-out, and the hottest shard is about 1.6 times the average. You bought efficiency with skew.
Writes, updates, and deletes
| Engine style | Insert | Update | Delete | Maintenance |
|---|---|---|---|---|
| In-memory HNSW | Online insert, a search plus linking | Delete plus insert, or replace a marked slot | Tombstone | Rebuild when tombstones pile up |
| Faiss IVF | Assign to the nearest centroid and append | Remove ids and add | Remove ids | Retrain centroids on drift |
| Lucene | New segment with its own HNSW graph | Delete the old document and index a new one | A live-docs bitset | Merges rebuild graphs. Force-merge read-only indexes |
| pgvector | The row insert updates the index | MVCC, a new row version | Dead tuples | VACUUM. REINDEX after heavy churn |
| DiskANN-style | An in-memory delta plus a periodic merge | Delta plus merge | A delete list | Background consolidation |
Graph inserts are expensive, so many systems buffer writes into a small in-memory segment, searched by brute force or a small graph, and merge on a timer. The visible effect is indexing lag between the write and searchability. That lag is an SLO, for example p99 under 30 seconds.
Tombstones
Deleted ids stay in the structure and are filtered after the search. The sandbox's rule of thumb: at 0 percent deleted you get k results. At 20 percent, asking for k returns about 8. At 50 percent, about 5. At 80 percent, under 2. Oversampling by about 1.5 / (1 - deleted fraction) fills the list again, and every query still pays to visit the dead nodes. Past roughly 20 to 30 percent tombstones, compact or rebuild. Connectivity around the holes degrades even when the count looks fine.
Stale embeddings
Three causes:
- Content changed, the vector did not. The document was edited and the re-embedding job failed or lagged. Search answers about old text.
- The model changed, the index did not. Queries embedded with model v2 hit an index of v1 vectors. The spaces are incompatible. Results are near random, and nothing errors.
- The index structure drifted. IVF centroids were trained on old data, or the graph was built before a large distribution shift.
Defenses:
- Store
model_id,model_version,content_hash, andembedded_atwith every vector. - Refuse, at the API, a query embedding whose model id does not match the index.
- Re-embed only the chunks whose content hash changed.
- Watch lag from source update to vector update, the share of vectors older than N days, and the count of hash mismatches.
How an embedding is produced, and why a chunk boundary changes the vector, stays on Embeddings & Similarity — Dense Vectors, Metrics & Chunking Basics. Stale retrieval as a grounding failure stays on RAG Failure Modes — Hallucination, Stale Indexes, Evals & Grounding.
Blue-green reindex
Decisions
- 1
1. Keep writing model v1
- next2. Dual-write new docs as v2
- 2
2. Dual-write new docs as v2
- next3. Backfill old docs in batches
- 3
3. Backfill old docs in batches
- next4. Alias still serves v1
- 4
4. Alias still serves v1
- next5. Shadow-query v2, compare
- 5
5. Shadow-query v2, compare
- next6. Backfill done, evals pass?
- ?
6. Backfill done, evals pass?
- no3. Backfill old docs in batches
- yes7. Flip alias to v2
- 7
7. Flip alias to v2
- next8. Keep v1 for rollback
- 8
8. Keep v1 for rollback
- next9. Delete v1 after the window
- 9
9. Delete v1 after the window
- Keep writing v1 embeddings into the live index.
- Dual-write new documents into the green index with v2 embeddings.
- Backfill old documents with v2, in batches.
- Serve users from v1 through the alias.
- Shadow-query v2 and compare recall and relevance offline.
- Flip the alias to v2 when backfill is complete and the evals pass.
- Serve from v2. Keep v1 for a fast rollback.
- Delete v1 after the rollback window.
The TypeScript sandbox makes the silent failure numeric. A v2 query against a v1 index scores recall 0 on this seed and throws nothing. During backfill, fusing a partial v2 index with the complete v1 index is worse than serving v1 alone. Shadow-read v2, then flip at 100 percent. Dual-read pays when documents written after a freeze exist only in v2.
Sandbox: scatter-gather, routing, tombstones
Random sharding fans out to every shard, and each shard must return enough neighbors for the global merge. Centroid routing probes fewer shards and skews. Tombstones shrink the result list unless you oversample, and oversampling does not make the dead nodes free.
Press Run. Snippets must be self-contained — no network, files, or native modules.
Sandbox: mixed spaces and a cutover
Two models of the same meaning are different coordinates. Here v2 is a rotation of v1. A v2 query on the v1 index does not throw. While the v2 index is partial, querying it alone is incomplete, and fusing it with the complete v1 index is worse than serving v1. A content hash next to the vector makes a text change detectable.
Press Run. Snippets must be self-contained — no network, files, or native modules.
Seeded reading: v1 on v1 is recall 1. At 25 percent backfill the partial v2 index is about 0.26 and the fused list is about 0.67, both worse than serving v1 at 1.0. At 100 percent both spaces recall 1. The document is stale because the hash and the model id both moved.
Production checklist
- Capacity: bytes per vector times N times replicas, with 30 to 40 percent headroom. The tuning page has the formulas.
- Shard count: as few as fit in RAM per node. Per-shard k equals global k.
- Filters: pick the strategy from selectivity. Exact fallback for tiny tenants. The filter page owns the strategies.
- Writes: an indexing-lag SLO, a buffer and a merge, an alert on the deleted-document ratio.
- Model upgrades: a model id on every vector, blue-green with a shadow eval, an alias flip, a rollback window.
- Monitoring: recall on a frozen query set, p99 per shard, IVF list-size skew, Lucene segment counts, freshness lag.
- Backups: snapshots include vectors. Restoring an index built with an old model is a staleness bug.
Interview Q&A
How do you shard a 2 billion vector index?
Answer
Do the memory math first. int8 HNSW at about 0.9 KB per 768-d vector is about 1.8 TB per copy. Pick a node size, hash-shard to fit with headroom, and add replicas for QPS. Each shard returns the full k. The coordinator merges. Consider IVF-PQ or DiskANN if a re-rank still hits the recall target and you need fewer nodes. Semantic routing only if fan-out cost dominates and you can live with skew.
Why not shard by semantic cluster to avoid fan-out?
Answer
It works. It is IVF at cluster scale. Boundary queries need several shards, hot topics create hot shards, and drift forces a re-shard. Hash sharding is simpler and evenly loaded. Most systems accept the fan-out.
Recall dropped 10 points over a month with no deploys. What happened?
Answer
Candidates: a rising tombstone ratio, IVF centroid drift, more Lucene segments diluting the candidate budget, a new content type outside the training distribution, or a partial re-embedding onto a new model. Check the deleted ratio, segment count, list skew, and the model id distribution. Then compact or rebuild.
How do you upgrade the embedding model without downtime?
Answer
Build a second index with the new model. Dual-write new content. Backfill old content in batches. Shadow-query and compare on a labeled set and on the exact oracle. Flip an alias. Keep the old index for rollback, then delete it. Never mix vectors from both models in one index.
How do you know embeddings are stale?
Answer
Store a content hash and a model id on each vector. Compare them to the source of truth in a periodic job. Measure lag from the source update to the vector update. Alert on mismatches and on the share of vectors embedded by a deprecated model.
Deletes are 40 percent of the index. What do you do?
Answer
Rebuild or force-merge. Tombstones waste RAM, slow the walk, and weaken connectivity. Oversampling only masks the short result list. Schedule compaction from the deleted ratio. For heavy churn, consider time-partitioned indexes you can drop whole.
Why might replicas return slightly different results?
Answer
HNSW builds are order-dependent, and segment layouts differ per replica, so approximate results can differ at the margin. Use a preference or session routing if you need stability, or accept a small nondeterminism and measure recall per replica.
Why is fusing a half-backfilled index with the old one a bad cutover?
Answer
The new index is missing neighbors that still live only in the old space, and the ranks are not the same coordinates. The sandbox shows the fused list losing to the complete old index until backfill hits 100 percent. Serve the old index, shadow the new one, and flip once.
Pitfalls
One box for the live alias, one for the green index, an arrow for dual-write, and a gate that does not flip until backfill and the shadow eval both pass. Mark the field that makes a mixed-model query fail closed. Then write per-shard k next to a three-shard fan-out.