The whole design, in writing
Learn AI system design by building a vector database step by step. An interactive guide to approximate nearest-neighbour (ANN) search — distance metrics, IVF cells, HNSW graphs, product quantization, metadata filtering, and sharding — so you can find the closest of a billion embeddings in milliseconds.
Every step of the build above, written out: the problem each piece solves, the option that was taken and the ones that were not, the numbers, and how it fails in production.
The big idea
Why a database just for vectors?
Embeddings turn text, images and users into points in a few-hundred-dimensional space, where nearest = most similar. The whole game is one query: "given this vector, find the k closest out of a billion." A normal database indexes values you can sort and equal-match — but "closest in 768-D space" isn’t a B-tree lookup. So how do you build an index for proximity?
A vector database is built around approximate nearest-neighbour (ANN) search: specialized in-memory indexes (IVF, HNSW) that return the closest vectors in milliseconds by being approximately right instead of exhaustively exact.
What the new pieces do
- Applicationclient
- A service holding a query vector (an embedded question, image, or user) that needs the k most similar vectors out of millions or billions.
Step 1 · The skeleton
A query vector in, k neighbours out
A service has a query vector and wants the k most similar stored vectors. What sits between the request and the answer?
What’s the minimal shape of a vector-search request?
Equality finds an identical vector, not a similar one — and two embeddings of the same idea are almost never bit-identical. You need ranking by distance, not exact match.
A coordinator accepts the vector, k, and any metadata filters, searches the index, and returns the k nearest ids with their similarity scores.
That’s exact kNN — correct but O(N) per query. Fine for 10k vectors, hopeless at a billion. It’s the baseline we’ll approximate, not ship.
A Query API takes a query vector, k, and optional filters, plans the search, and returns the k nearest ids + scores. Everything else — indexes, segments, shards — hangs off this one contract.
Why this piece earns its place
The contract is the part you cannot change later, so three things in it are worth being deliberate about. First, return ids and scores, not stored payloads. Callers will ask for payloads to save a round trip, and the moment you agree, the vector database is also a document store, with its size limits, its consistency questions and its own compaction problems. Second, the score is index-dependent. A client that hardcodes "drop anything below 0.8" has bound itself to the metric, to normalization and, once compression arrives, to the codebook, so retuning any of those stops being a change on your side and becomes a coordinated release across every caller that ever wrote a threshold down. Publish ranks, or a score carrying the metric that produced it. Third, there is no cheap pagination. This index answers "the k nearest" and nothing else, so "results 100 to 200" means re-running the whole search with k at 200 and throwing most of it away. Cost grows with the offset, not the page. Settle that infinite scroll is not a feature here before someone designs a screen assuming it is.
What the new pieces do
- Query APIbackend
- Accepts a query vector + k + filters, plans the search across shards, merges results, and returns the nearest neighbours with their ids and scores.
Step 2 · Store what you search
Ingest, segments, and the index
Vectors arrive from an embedding pipeline. You need to both keep them (to return and re-rank) and index them (to search). How does the write path work?
A new batch of vectors arrives. Where do they go?
In-place edits to a graph/cell index under concurrent search cause lock contention and fragmentation. Most vector DBs write to immutable segments and merge in the background.
Rebuilding per query throws away the whole point of an index. The structure must persist between queries and grow incrementally.
The ingest path does both: the raw vectors land in an immutable segment (for exact re-rank and rebuilds), and each vector is inserted into the ANN index (a cell assignment or graph links).
An Ingest / Upsert path writes each vector twice: the full vector into an immutable Vector Segment (for exact re-ranking and index rebuilds), and a pointer into the ANN Index (graph links or a cell assignment). Segments merge in the background, so live search never blocks on writes.
Why this piece earns its place
Everything above describes adding. Removing is where this design grows teeth. A segment is immutable, so a delete is a tombstone, and the index cannot honour it either: pulling a node out of an HNSW graph severs the edges its neighbours used to reach that whole region, so implementations leave the node in place and filter it out of the answer. Deleted vectors therefore keep occupying RAM and keep being traversed, and an upsert of an existing id is a delete plus an insert — which makes re-embedding a corpus, the most ordinary operation there is, the worst case for both. Ask when the data is really gone and the honest answer is the merge schedule, so that schedule is a commitment made in the design rather than a knob ops picks later. The other thing worth being explicit about is that the two writes land at different times. Durable and searchable are separate guarantees here, and batching index inserts widens the gap on purpose, so the contract has to say which of the two a successful upsert has bought.
- immutablesegments, merged in bg
- 2 writessegment + index
- rebuildableindex derives from segments
What the new pieces do
- Vector Segmentsstore
- The full-precision vectors and ids on disk/RAM, organised into immutable segments. Used to re-rank candidates exactly and to rebuild the index.
- Ingest / Upsertbus
- The write path: validate incoming vectors, append them to a segment, and insert them into the index (graph links or cell assignment).
- Embeddingsstore
- Vectors produced upstream by an embedding model (text, images, users). The vector DB stores and searches them — it doesn’t create them.
Back of the envelope
- append-only segments
- no in-place mutation under live search
- background merge
- compact small segments, drop deletes
- index follows store
- segment is truth; index is an accelerator
- batch upserts
- amortize index-insert cost over many vectors
Step 3 · The naive baseline
Exact kNN and the distance metric
To rank by similarity you need a number for "how close." And the obvious algorithm — compare the query to every vector — is correct. Why can’t you just ship that?
What kills brute-force exact nearest-neighbour at scale?
Exact search compares the query to every vector. At a billion vectors × hundreds of dims, that’s billions of multiply-adds per query — far too slow and CPU-hungry to serve online.
Exact kNN is, by definition, correct — it’s the ground truth other methods are measured against. Its problem is cost, not correctness.
Distances compute fine in high-D; the issue is doing N of them per query. (High-D does make the geometry harder — which is why approximate indexes work so well.)
Similarity is a distance metric — usually cosine (angle) or dot product for embeddings, L2 for some. Exact kNN computes it against every vector and keeps the top k: correct, and O(N·d) per query. That’s the baseline — and why everything past here is about avoiding it.
Why this piece earns its place
Keeping exact search runnable is not sentiment — it is the measuring instrument for everything after it, and instruments need owners. The part teams get wrong is the query sample. Recall measured on random vectors flatters the index: a random point in 768-D space usually has one unambiguous nearest neighbour, while real queries land where the data is dense and the top k are near-ties, which is exactly where an approximate index drops one. So the sample is real traffic, frozen so that two runs are comparable, and refreshed deliberately when the query mix moves rather than whenever someone reruns a notebook. It is also the O(N·d) cost you just refused to serve, paid offline every time the corpus changes, which makes the sample size a budget line rather than a detail. One more property this step fixes permanently: the metric is a build-time choice, not a request parameter. An index built for dot product does not answer L2 questions, so picking wrong is not a config change, it is a rebuild of everything standing on it.
- cosine / dottypical for embeddings
- O(N·d)exact, per query
- ground truthwhat recall is measured against
What the new pieces do
- ANN Indexindex
- The heart of the system: an in-memory structure (IVF cells or an HNSW graph) that answers approximate nearest-neighbour search in sub-linear time.
Step 4 · Partition the space
IVF — search a few cells, not all
Exact search looks at all N vectors. But the answer is almost always near the query. How do you skip the 99% of vectors that are obviously far away?
How do you avoid scanning vectors that can’t be the answer?
There’s no single ordering of points in 768-D space — "sorted" has no meaning across many dimensions. You need spatial partitioning, not a 1-D sort.
IVF runs k-means to make ~√N centroids ("cells"). A query finds its closest centroids and scans only the vectors in those nprobe cells — a fraction of N.
Random sampling misses the actual neighbours most of the time. The cells aren’t random — they’re chosen so nearby vectors share a cell.
IVF (inverted file) clusters the vectors with k-means into many cells, each with a centroid. A query measures distance to the centroids, picks the closest nprobe cells, and scans only those. Visit more cells → higher recall, slower; fewer → faster, lower recall.
Why this piece earns its place
Two properties of cells matter more than the arithmetic. They are contiguous: the vectors in a cell sit together, so probing one is a sequential read, which makes IVF the easier index to serve from disk or a memory map — not the only one, since graph indexes have been re-engineered for SSD residency, but the one that gets there without special layout work. They are also independent, so probed cells can be scanned on separate threads and the search parallelizes within a single query, which a greedy sequential walk cannot. The part that surprises people is the first stage. Finding the nearest nprobe centroids means comparing the query against every centroid, and at a billion vectors the square root of N is tens of thousands of them — a brute-force scan in miniature, the thing cells were built to avoid. So the cell count is a trade between the two stages, not a dial to turn up: more cells shorten every posting list and lengthen the centroid scan, and past a certain size you index the centroids themselves with a small graph, which is why real systems compose IVF and HNSW instead of choosing between them.
- ~√Ncells (k-means)
- nprobecells scanned per query
- recall ⇄ speedset by nprobe
Back of the envelope
- train centroids
- k-means on a sample defines the cells
- N=1M vectors → √N≈1,000 cells; nprobe=32 scans ~32,000 of them
- ≈3.2% of the store — about 31× fewer vectors than an exact scan
- imbalanced cells hurt
- periodically retrain as data drifts
- measure recall@k
- against an exact baseline, not by vibes
Step 5 · Walk a graph instead
HNSW — navigable small worlds
IVF is great, but its recall plateaus and it needs retraining as data shifts. What if, instead of partitioning space, you could walk straight to the neighbourhood from any starting point?
How can you reach a query’s neighbours in a few hops, from anywhere?
A complete graph is O(N²) edges — impossibly large and slow to traverse. You need few, well-chosen links, not all of them.
k-d trees degrade badly in high dimensions — they end up visiting most of the tree. They work in 2-D/3-D, not in 768-D embedding space.
HNSW builds a layered "small-world" graph: sparse long-range links up top for big jumps, dense local links below for precision. Search greedily hops to ever-closer nodes.
HNSW (Hierarchical Navigable Small World) links each vector to a handful of near neighbours across layers: a sparse top layer for long jumps, denser lower layers for fine approach. Search enters at the top, greedily hops to closer nodes, and descends — reaching the neighbourhood in a few hops. efSearch controls how many candidates it keeps in flight: the speed/recall dial.
Why this piece earns its place
The distinction to hold onto is which of these parameters you can still change after launch. efSearch is per request: a query that comes back unconvincing can be retried wider, and two callers can buy different accuracy from one index. M and the construction effort cannot — they are baked into the edges, so changing your mind means building the graph again. And building it is expensive in a particular way: construction is essentially one search per inserted vector, so a rebuild is the same order of work as answering every query in the corpus once. That is what makes HNSW’s headline property load-bearing rather than merely convenient. Taking inserts without retraining matters because the escape hatch of starting over is priced out of reach. The edges also set the floor on memory, and they will not page gracefully — consecutive hops land in unrelated parts of the adjacency lists, so there is nothing to prefetch. They are per node, and no amount of vector compression touches them, so once the vectors themselves are squeezed the graph decides how large a machine has to be. At a billion vectors the question stops being whether the vectors fit and becomes whether their links do.
- ~O(log N)hops to the region
- M linksper node per layer
- efSearchcandidates in flight
Back of the envelope
- M ≈ 16–64
- more links = better recall, more memory
- N=1M vectors → log₂(N) ≈ 20 hops to the neighbourhood
- versus touching all 1,000,000 vectors in an exact scan
- efSearch tunes recall
- higher = more accurate + slower
- great for updates
- insert links incrementally, no retrain
- RAM-hungry
- the graph edges live in memory
Step 6 · Make a billion fit
Product quantization compresses the vectors
A billion 768-D float32 vectors is ~3 TB — it won’t fit in RAM, and the index needs RAM to be fast. How do you shrink the vectors without losing the ability to rank them?
How do you fit billions of vectors in memory and still rank by distance?
Truncating dimensions throws away the information the distance depends on — recall craters. You want to compress, not amputate.
Product quantization chops the vector into m sub-vectors, k-means-clusters each subspace, and stores tiny codebook ids. A 768-D float32 vector (3 KB) becomes ~64–96 bytes — distances estimated from the codes.
Generic compression saves disk but you must decompress to compute distance — no speedup for search. PQ lets you estimate distance directly from the compressed codes.
Product Quantization (PQ) splits each vector into m sub-vectors, learns a small codebook per subspace, and stores only the codebook ids — turning a 3 KB vector into ~64–96 bytes. Search estimates distances from the codes (fast, approximate), then re-ranks the top candidates against the full vectors in the segments for precision.
Why this piece earns its place
The first question is whether you need this at all. Below the memory ceiling it buys nothing but risk, and the cheaper move most systems should make first is scalar quantization — store each dimension as an int8 rather than a float32 for a flat four-fold saving, with no codebook and far less distortion than an aggressive PQ code. It is not free of training, though: fixing the range each dimension maps onto takes a calibration pass over the data, and that range ships with the codes, because it is the only thing that says what an int8 means. Reach for PQ when four-fold is not enough, and understand that you are taking on the same problem in a larger form. The codebook is the only thing that can interpret the codes, so it must be versioned and shipped with the segments it describes; a codebook restored from a different training run decodes the same bytes into different vectors, and nothing errors. Then there is the read side. Pulling full vectors back means a random read per candidate, so re-rank depth is priced in seeks rather than compute — a reason to lay segments out so a shortlist is a few sequential reads instead of a hundred scattered ones.
- ~3 KB → ~96 Bfloat32 → PQ code
- fits in RAMbillions of codes
- re-rankexact, on the top few
What the new pieces do
- Quantizerservice
- Product quantization: compress each vector into a short code so billions fit in RAM, at the cost of a little precision (recovered by an exact re-rank).
Step 7 · Nearest, but filtered
Metadata filters and the recall trap
Real queries aren’t just "nearest" — they’re "nearest where tenant = X and recency < 30 days." Bolt a filter onto ANN search naively and recall quietly tanks. Why?
You want the nearest vectors that also match a metadata filter. What goes wrong?
Post-filtering: if few of the top-k match, you’re left with 1–2 results — you asked for k and the filter ate them. Fine only when the filter is loose.
Pre-filtering: correct, but if the filter still matches millions you’re back to brute force. Great for very selective filters, costly for loose ones.
There’s no single winner: the planner picks a strategy by selectivity, often over-fetching (k′ > k) before filtering, or evaluating the predicate during graph/cell traversal.
The Query API consults a Metadata Store and picks a strategy by selectivity: pre-filter when the predicate is selective (search only matching rows), over-fetch then post-filter when it’s loose (grab k′ ≫ k so enough survive), or push the predicate into the traversal. Naive post-filtering is the classic way to silently return too few, wrong results.
Why this piece earns its place
The first thing to settle is which attributes can be filtered at all, because that is decided when the index is built. Testing a predicate mid-traversal only works if the attribute sits beside the vectors, as a bitmap or a small column inside the index process; anything living only in the metadata store costs a lookup per candidate, which a post-filter over a shortlist absorbs and a walk cannot. So the fast filters are the ones you chose to copy in, and adding one later is an index rebuild, not a schema migration. The second is that filtering inside a graph fights the graph. Its links were chosen over the whole dataset, and restricting the walk to matching nodes deletes most of them; a small-world graph missing most of its edges is no longer navigable, so the search strands in a region with no eligible neighbours. Implementations keep walking through non-matching nodes as stepping stones and collect only the matching ones. All of which makes the filter language an API commitment: a closed set of indexed attributes stays fast for years, while arbitrary predicates promise a planner, and the statistics to feed it, for as long as the product lives.
What the new pieces do
- Metadata Storestore
- Per-vector attributes (tenant, source, recency, tags) used to constrain a search — "nearest, but only within these rows".
Step 8 · Past one machine
Shard, replicate, scale
One billion-vector index outgrows a single box’s RAM, and one replica can’t serve the QPS. How do you grow past one machine without breaking "k nearest"?
How do you scale a vector index beyond one machine’s memory?
Metadata sharding helps multi-tenant isolation but skews load and doesn’t help a single huge tenant. The general answer is to partition the vectors themselves and search all shards.
A Shard Router fans the query out to every shard, each returns its local top-k, and the coordinator merges them into the global top-k. Add read replicas per shard for QPS and availability.
Vertical scaling hits a RAM ceiling and a single-point-of-failure wall. Past a point you must distribute the index across machines.
A Shard Router partitions vectors across shards and runs scatter/gather: every shard searches its slice and returns a local top-k; the coordinator merges them into the global top-k. Replicate each shard for QPS and failover. Because nearest-neighbour is mergeable, the global answer is just the best of the locals.
Why this piece earns its place
The sentence to have ready is that sharding buys memory, not throughput. Because every query is scattered to every shard, each shard sees the full query rate however many shards there are; doubling the shard count halves the vectors per machine and changes the queries per machine not at all, so the fleet is shards multiplied by replicas. The second consequence is the tail. A scatter-gather response is as slow as its slowest shard, so p99 is a maximum over shards rather than an average, and it gets worse as you add them — the familiar reason a well-behaved service degrades as it grows. What rescues it here is that k-nearest is unusually forgiving: a shard that misses its deadline costs recall, not correctness, so you can hedge the request, or merge without that shard and report the degradation, which beats making every user wait on the unlucky machine. The precondition nobody writes down is that merging assumes the scores are comparable. Best-of-the-locals holds only while every shard shares one metric, one normalization and one generation of the index — which makes an index rebuild a fleet-wide event rather than something you roll one machine at a time.
- scatter / gatherquery all shards, merge
- replicasQPS + failover
- mergeableglobal top-k = best of locals
What the new pieces do
- Shard Routerservice
- Fans the query out to every shard, then merges their local top-k into a global top-k. Lets the index grow past one machine’s RAM.
The payoff
You built a vector database
From "find the closest of a billion vectors" to a system that does it in milliseconds: a read-optimized Query API, an ingest path feeding immutable segments and an ANN index, IVF cells or an HNSW graph for sub-linear search, product quantization to fit RAM, a filter-aware planner, and scatter/gather sharding.
Now starve the search — drop nprobe/efSearch to the floor — and watch recall collapse with no error at all: the index returns k vectors, just not the nearest ones. That’s why you measure recall@k against an exact baseline instead of trusting that "it returned results."
Everything you assembled, in order
- Query API — one contract: query vector + k + filters → nearest ids
- Segments vs index — full vectors are truth; the ANN index is a rebuildable accelerator
- Distance metric — cosine/dot/L2 — normalize, and match the embeddings
- Exact kNN — correct but O(N·d) — the baseline you approximate
- IVF — k-means cells; nprobe trades recall for speed
- HNSW — layered small-world graph; ~log N hops, efSearch dial
- Product quantization — compress to fit RAM, re-rank exactly on the shortlist
- Filters + sharding — a selectivity-aware planner; scatter/gather mergeable top-k
