An approximate nearest neighbour index on one machine is a solved problem: pick HNSW or IVF, tune a couple of parameters, and you get millisecond queries over a few million vectors. Serving hundreds of millions of vectors to thousands of queries per second, with filters, fresh writes and an embedding model that will change next quarter, is a different problem. The index algorithm becomes one component in a distributed system, and most of the hard decisions are about memory, sharding, freshness and migration.

This article designs that system from first principles. It separates the write and read paths, explains segments, works through the capacity math for a concrete workload, and covers sharding, filtering, re-embedding and how to know your recall in production. Choosing between index types is covered in vector search infrastructure, and how vectors fit into a broader search stack with lexical retrieval and rank fusion is in search system design.

Advertisement

The workload that makes it hard

Work with a concrete target. A document platform holds 100 million chunks, each embedded as a 768-dimensional float32 vector. Peak traffic is 3,000 queries per second, the latency budget for retrieval is a p99 of 80 milliseconds, recall@10 must stay at or above 0.95 against exact search, new or edited documents must be searchable within a minute, and most queries filter by tenant and language. Every one of those numbers pushes the design somewhere different.

Memory decides the node count, QPS decides the replica count, the latency budget decides how many shards a query may touch, freshness decides how writes reach the index, and filters decide whether the index can be shared at all. Keep all five in view: a design that meets four of them usually fails the fifth in production.

Two paths, one handoff

The core decision is to split writes from reads. Graph indexes such as HNSW are expensive to build and awkward to mutate: inserts take locks on neighbour lists, deletes leave holes the graph routes around, and build work competes with queries for CPU and memory bandwidth. So production systems build indexes in one place and serve them from another.

Write path and read path meet only at sealed segments in object storageSource DBchange eventsEmbedding servicemodel v3, batchedDurable logid, vector, attrsIndex builderseal, build, compactObject storagesealed segmentsGrowing segmentrecent writes, flatClientquery text + filterCoordinatorembed, fan out, mergeShard 1 replicagraph + int8Shard N replicagraph + int8Rerankfp32 or cross-encoderuploadtailloadQuery nodes never build indexes; builders never serve queries. Each scales on its own.
The write path turns change events into sealed, indexed segments. The read path loads segments and answers queries. A small growing segment bridges the freshness gap.

On the write path, change events from the source database go through an embedding service, which batches them for GPU efficiency, into a durable log keyed by document id. Index builders consume the log. Recent records collect in a growing segment that is searched by brute force, which is cheap while it is small. When it reaches a size threshold, say a million vectors, the builder seals it, builds a graph and quantized codes for it, and uploads the result to object storage. Query nodes then download and serve it.

The read path is a coordinator plus query nodes. The coordinator embeds the query with the same model version as the index, fans it out to shards, merges results and reranks. Query nodes hold sealed segments in memory and scan their share of the growing segment. Because the two paths meet only at immutable segments, each scales independently: a large backfill adds builders and does not slow queries.

Advertisement

Segments, deletes and compaction

Immutable segments make updates simple and searches slightly harder. An update is a delete plus an insert: the new vector goes into the growing segment, and the old id is marked in a per-segment deletion bitmap. Searches skip deleted ids while walking the graph. This is the same log-structured pattern Lucene and LSM trees use, and it inherits the same problem: as deletes accumulate, segments carry dead vectors that still cost memory and still sit on graph paths, so recall and latency drift.

Compaction fixes it. A background job picks segments with a high deleted fraction, or many small segments, rebuilds one clean segment from their live vectors, uploads it, and atomically swaps the segment list. Set triggers on the deleted fraction, for example 20 percent, and on segment count per shard, because every extra segment is another index to search per query. Budget the transient memory: during the swap, a query node briefly holds both old and new segments.

Capacity math, worked

Raw float32 vectors cost 768 x 4 = 3,072 bytes each, so 100 million of them take 307.2 GB, or 286.1 GiB. An HNSW graph with M = 16 stores up to 32 neighbour ids of 4 bytes on the base layer, about 128 bytes per vector, or 12.8 GB, with the upper layers adding little. Holding float32 vectors plus graph in RAM needs about 298 GiB before any headroom.

Quantization changes the picture. Scalar int8 quantization stores one byte per dimension, 76.8 GB in total, so int8 plus graph is 89.6 GB, or 83.4 GiB. Product quantization with 96 one-byte sub-codes per vector needs only 9.6 GB, at a larger recall cost. The standard pattern is to search on compressed vectors and rerank a small candidate set with the full float32 vectors read from local SSD.

Layout in RAMPer shard budgetShardsNodes at 5 replicas
fp32 + graph, 298 GiB21 GiB1575
int8 + graph, 83.4 GiB21 GiB420

The 21 GiB budget assumes 64 GiB nodes that keep the index to about a third of memory, leaving room for the segment swap during compaction, the growing segment, the page cache for SSD reranking and the process itself. Replicas come from throughput: suppose a load test shows one replica of one shard sustains 600 QPS at the target recall and latency. Every query touches every shard, so each shard needs 3,000 / 600 = 5 replicas. With int8 that is 4 x 5 = 20 nodes, each with about 72 GiB of fp32 vectors on local NVMe for reranking. The same design in float32 needs 75 nodes. Measure your own per-replica QPS; it varies with dimension, the search beam width and hardware.

If memory is still the constraint, disk-resident graphs such as DiskANN keep only compressed codes in RAM and read full vectors and neighbour lists from SSD, trading a few SSD reads per query for a much smaller memory footprint.

Sharding and the tail

There are two ways to shard. Hashing document ids spreads load evenly but forces every query to fan out to every shard. Sharding by tenant lets a filtered query touch one shard, but large tenants create hot shards and small ones waste capacity. Many systems combine them: hash by default, and give very large tenants dedicated shards. The general trade-offs are those in sharding.

Fan-out multiplies tail latency. If each shard answers within its own p99 99 percent of the time, a query that must wait for four shards is slow whenever any one of them is: 1 - 0.99^4, or 3.9 percent of queries. At fifteen shards it is 14 percent. That is another reason to compress and keep shard counts low, and it is why the coordinator should hedge: after the typical per-shard latency, send a duplicate request to another replica of any shard that has not answered, and take the first response.

import asyncio, heapq

async def search(query_text, k=10, tenant=None, oversample=4, timeout_s=0.08):
    qv = await embed(query_text, model=ACTIVE_MODEL)            # same model as the index
    plan = choose_filter_strategy(tenant)                        # see the filtering section
    shards = route(tenant)                                       # all shards, or one for a tenant
    calls = [shard_search(s.pick_replica(), qv, k * oversample, plan) for s in shards]
    done = await asyncio.wait_for(gather_partial(calls), timeout_s)  # keep what arrived
    if len(done) < len(shards):
        metrics.incr("vs.partial_results", len(shards) - len(done))
    candidates = heapq.nsmallest(k * oversample, (h for hits in done for h in hits),
                                 key=lambda h: h.approx_dist)
    exact = await fetch_full_vectors([h.id for h in candidates])     # fp32 from SSD
    rescored = sorted(candidates, key=lambda h: l2(qv, exact[h.id]))
    return rescored[:k]

Two details in that code matter. Each shard returns more than k candidates, here four times as many, because its local top-k is computed on compressed vectors and the global merge needs slack. And the timeout returns partial results rather than failing, with a metric, because for most retrieval uses nine shards of ten are better than an error.

Filtering at scale

Filters are where vector search differs most from its textbook description. Post-filtering, searching first and dropping non-matching results, fails for selective filters: if a tenant owns 0.1 percent of vectors, a top-40 search returns almost nothing from that tenant. Pre-filtering, computing the allowed ids first and searching only among them, breaks graph navigation when the allowed set is small, because the graph's shortcuts pass through disallowed nodes.

Production planners choose per query by estimated selectivity. When the filter matches only a few thousand vectors, skip the index and brute-force them: exact and fast at that size. When it matches a large fraction, search the graph while skipping disallowed nodes but still traversing through them, and raise the search beam width so enough allowed results survive. In between, a partition or separate index per high-volume filter value, such as a language, often wins. Keep filterable attributes in columnar form beside each segment so the planner can count matches cheaply.

Freshness and consistency

The one-minute freshness target comes from the growing segment. A write reaches the durable log in milliseconds, the embedding service adds its batch delay, and query nodes tail the log into their growing segment, so new documents are searchable in seconds without waiting for a seal. The price is brute-force search over the growing segment on every query, which is why its size cap matters.

Consistency is usually eventual, and that is fine for retrieval, with one exception: a user who uploads a document and immediately searches for it. Handle that case explicitly by passing the write's log offset with the query and having the coordinator wait until the serving nodes have consumed that offset, bounded by a short timeout.

Re-embedding: the migration you will run

Embedding models improve, and a vector from one model is meaningless in another model's space. Unlike an index parameter change, a model upgrade means re-embedding every document, and you cannot mix old and new vectors in one index. Treat it as a versioned migration, not an in-place update.

# Re-embedding with a new model: vectors from v3 and v4 are not comparable.
create_collection("docs_v4", dim=1024, index=HNSW(M=16, ef_construction=200))
backfill("docs_v4", source="warehouse.docs", model="embed-v4")      # batch, hours to days
dual_write(["docs_v3", "docs_v4"])                                  # new changes go to both
shadow_read(primary="docs_v3", shadow="docs_v4", sample=0.05)       # compare, never serve v4
assert evaluate(judgements, "docs_v4") >= evaluate(judgements, "docs_v3")
switch_alias("docs", "docs_v4")                                     # coordinator also embeds with v4
keep("docs_v3", days=7); drop("docs_v3")

The alias is the key: clients query docs, and the switch is a single metadata change that is easy to roll back. The coordinator must switch its query embedding model at the same moment, so store the model name with the collection and have the coordinator read it, rather than configuring the two separately. Budget for double capacity during the overlap, and for the embedding cost of the backfill, which for 100 million chunks is a substantial GPU job in its own right.

Knowing your recall in production

Recall is the metric that silently degrades. Compaction lag, deletes, parameter changes and data drift all lower it without raising any error or latency alarm. Measure it directly: sample a few hundred real queries each hour, compute exact top-10 by brute force on one shard, and compare with what the live index returns on that shard.

def recall_probe(sample_queries, shard, k=10):
    """Run hourly on a small sample: exact brute force over one shard vs the live index."""
    total = 0.0
    for q in sample_queries:
        truth = set(brute_force_topk(shard.full_vectors(), q, k))   # slow, exact
        got = set(h.id for h in shard.search(q, k, params=LIVE_PARAMS))
        total += len(truth & got) / k
    metrics.gauge("vs.recall_at_10", total / len(sample_queries), tags={"shard": shard.id})

Alert when recall@10 falls below your target, and log it next to the segment count and deleted fraction so the cause is visible. Offline relevance judgements measure something different, whether the model finds useful documents, and belong in the migration gate above.

Failure modes

  • Model mismatch. The coordinator embeds with a different model or version than the index; results look plausible and are nearly random. Store the model with the collection and check it on every query.
  • Selective filters returning nothing. Post-filtering drops everything for small tenants. Plan by selectivity.
  • Compaction debt. Deleted vectors pile up, recall and latency drift, and memory grows with no new data.
  • Growing segment runaway. A stalled builder lets the brute-force segment grow until every query slows down. Alert on its size.
  • Fan-out tail. Too many shards make p99 a function of the slowest node. Compress, cap shards and hedge.
  • Rebuild storms. A cluster restart makes every node download and load every segment at once. Stagger restarts and keep segments on local disk.

What to do next

  1. Write down vector count, dimension, peak QPS, p99 budget, recall target, freshness and filter selectivity for your workload.
  2. Compute memory for fp32, int8 and PQ layouts, choose one, and derive shards from a per-node budget that leaves room for compaction.
  3. Load-test one shard replica to find its QPS at your recall target, then derive replicas.
  4. Separate index building from serving, with a growing segment for freshness and compaction triggers on deleted fraction and segment count.
  5. Implement a selectivity-based filter planner and oversampled, hedged fan-out with fp32 reranking.
  6. Store the embedding model with each collection, query through an alias, and rehearse a re-embedding migration before you need one.
  7. Run an hourly recall probe per shard and alert on it.
Key takeaway: Vector search at scale is a storage and distributed systems problem wrapped around an ANN algorithm. Split the write path from the read path so that builders and query nodes meet only at immutable segments, keep a small brute-force segment for freshness, and compact away deletes. Size the cluster with arithmetic: quantization sets memory, memory sets shards, measured per-replica throughput sets replicas, and shard count sets the tail. Plan filters by selectivity, treat a model change as a versioned migration behind an alias, and measure recall continuously, because it is the one failure that raises no alarm on its own.