> ## Documentation Index
> Fetch the complete documentation index at: https://authorsnote.askailab.online/llms.txt
> Use this file to discover all available pages before exploring further.

# Scaling a Vector Database

> How a vector store actually grows to billions of embeddings: ANN index choice, the recall/latency/memory triangle, quantization arithmetic, sharding and replication, and the hot-partition failure mode.

A scalable vector database is **a specialized data store designed to house, index, and query billions of high-dimensional vector embeddings with sub-100 millisecond response times**. That sentence holds three separate engineering problems, and most capacity reviews I have been in argued about the wrong one: what index runs on a single node, how do I split and copy that node, and how many bytes per vector am I willing to spend to make the first two affordable.

These are the notes I keep, with the arithmetic written out so I can check it against a real cluster instead of a vendor slide.

## What a vector database is actually storing

**Embeddings** are numeric arrays that capture the semantic meaning of unstructured data like text, images, and audio ([Google Cloud's explainer](https://cloud.google.com/discover/what-is-a-vector-database), [Databricks' explainer](https://www.databricks.com/blog/what-is-vector-database)). **ANN** — approximate nearest neighbor — methods like **HNSW** (Hierarchical Navigable Small World graphs) and **IVF** (Inverted File Index) trade tiny bits of accuracy for massive speed gains ([Truefoundry's vector DB comparison](https://www.truefoundry.com/blog/best-vector-databases), [Google Cloud](https://cloud.google.com/discover/what-is-a-vector-database)). **Quantization** — Product Quantization, or PQ — compresses vectors to reduce memory consumption ([Oracle on Pinecone-style scaling, part 1](https://www.oracle.com/in/database/vector-database/pinecone/), [part 2](https://www.oracle.com/ae/database/vector-database/pinecone/)). **Distributed scaling** uses horizontal partitioning to split datasets across independent nodes or pods so queries run in parallel ([Oracle](https://www.oracle.com/ae/database/vector-database/pinecone/)).

For the survey literature, [this vector-database overview on arXiv](https://arxiv.org/abs/2505.12524) and [this enterprise selection guide](https://medium.com/@amitkharche/vector-databases-choosing-the-right-one-for-scalable-enterprise-genai-0705565ef333) are my two references when a claim needs citing rather than repeating.

## The single-node decision: HNSW, IVF, or ScaNN

Before I shard anything, I need to know what one node can do, because the index choice sets the shape of every number that follows.

| Index | Memory for 768-dim fp32 | Recall\@10 achievable | p50 at 10M vectors | What breaks first |
| :- | :- | :- | :- | :- |
| Brute force (flat) | 3,072 bytes/vector, no overhead | 1.00 by definition | seconds, not milliseconds | Scan cost grows linearly with N |
| HNSW | raw plus link overhead, roughly 1.08x at 768d and M=32 | 0.95 to 0.99 | 5 to 20 ms | Build time and RAM; deletions fragment the graph |
| IVF (plus PQ codes) | raw plus centroid list, e.g. 65,536 centroids = 0.20 GB | 0.85 to 0.97, nprobe-limited | 10 to 40 ms | Recall collapses when the query falls between cells |
| ScaNN (anisotropic PQ) | PQ-scale: 96 bytes/vector at m=96 | 0.95 to 0.98 at high speed | 5 to 15 ms | Tuning surface; the gain is query-time, not storage |

The reasoning behind each row is what I use in a design meeting.

**HNSW** is a layered navigable small-world graph. Search cost is logarithmic in N, which is why latency holds nearly flat as the corpus grows — 10M to 100M vectors costs far less than 10x, because hop count grows with `log N`. The price is that the whole graph lives in RAM. At 768 dimensions a float32 vector is `768 × 4 = 3,072 bytes`, and with M=32 links as 4-byte ids the graph adds roughly `2 × M × 4 = 256` bytes per vector at layer zero, so about 3,328 bytes per vector — only 1.08x the raw vectors here, but a far bigger relative multiplier at 128 dimensions where raw is just 512 bytes. Low-dimensional corpora are where HNSW overhead surprises people.

**IVF** quantises the space into `nlist` Voronoi cells, assigns each vector to its nearest centroid, and searches only `nprobe` cells at query time. For 100M vectors I pick `nlist` in the tens of thousands — at 65,536 cells and `nprobe=16` I inspect roughly 25,000 vectors out of 100M. Centroid storage is trivial: `65,536 × 768 × 4 = 0.20 GB`. The failure mode is structural: a query near a cell boundary loses the neighbour one cell away, and `nprobe` cannot recover it cheaply. Overlapping assignment (each vector in 2 to 8 cells) fixes recall and multiplies memory.

**ScaNN** is Google's answer: anisotropic quantization, which spends distortion budget preferentially along the direction that matters for inner-product ranking rather than uniformly. It buys recall at a given compression ratio, so I keep PQ-sized memory without accepting the recall cliff. The win is query-time scoring efficiency, not storage.

<Note>
  The three-way tension is real and there is no corner of the triangle that satisfies a product: **recall** (did I find the true nearest neighbours), **latency** (how fast, at what tail), **memory** (how many nodes and how much RAM per node). Improving any one of them costs one of the others. My job is not to find the best index; it is to pick which axis I am deliberately paying on, and to write that choice down. A fourth axis sneaks in sideways: ingestion and rebuild cost.
</Note>

## Quantization arithmetic, in bytes

This is the section that decides the size of the bill, so I do it in explicit bytes. The rule: `bytes per vector × number of vectors`, plus index overhead, times replication.

| Vectors | fp32, 768d (3,072 B) | int8, 768d (768 B) | PQ, m=96 at 8 bits (96 B) | fp32 plus HNSW links |
| :- | :- | :- | :- | :- |
| 1,000,000 | 3.07 GB | 0.77 GB | 0.10 GB | 3.33 GB |
| 100,000,000 | 307.2 GB | 76.8 GB | 9.6 GB | 332.8 GB |
| 1,000,000,000 | 3,072 GB (3.07 TB) | 768 GB | 96 GB | 3.33 TB |

Doubling the dimension to 1,536 doubles every column except PQ with a fixed code size: 1B vectors at 1,536d fp32 is 6.14 TB of raw vectors, while 96-byte codes stay at 96 GB. That asymmetry is the entire reason large deployments run a two-stage search: **scan compressed codes to rank candidates, then rerank a small candidate set against exact or int8 vectors held on SSD**. The compressed stage is cheap and wide; the expensive stage sees only the top few hundred.

Three quantization decisions I have to make explicitly:

* **Scalar quantization (fp32 to int8)** is a 4x memory cut with usually small recall loss, and the first thing I try. Calibration matters: outliers in one dimension shrink effective resolution for everything else.
* **Product quantization** is a 32x cut at `m=96` on 768d, but lossy in a ranking-relevant way, so I always measure recall\@k after quantizing rather than assuming it.
* **Binary quantization** is a 32x cut over int8 and only makes sense for high-dimensional embeddings trained to be Hamming-friendly. Without that training, recall falls off a cliff.

The metric to track is not compression ratio but **recall at the latency the product promises**, measured on a held-out query set that looks like real traffic. A 30x memory saving that drops recall\@10 from 0.95 to 0.82 is a quality regression that shows up as bad answers months later.

## Sub-100 milliseconds, decomposed

"Sub-100 ms" is a budget, so I spend it line by line:

| Stage | Typical allowance | Notes |
| :- | :- | :- |
| Client to router | 1 to 3 ms | Same region, connection reused |
| Router fan-out to shards | 5 to 40 ms | Parallel, so it is the max across shards, not the sum |
| Merge and dedupe top-k | 1 to 5 ms | CPU-bound, grows with `k × shards` |
| Metadata filter or rerank | 5 to 30 ms | Cross-encoder reranking is a separate, much bigger number |
| Response back | 1 to 3 ms | Payload size, usually small |

Because the fan-out is parallel, adding shards does not add latency linearly — 8 shards at 8 ms each is about 8 ms plus merge, not 64 ms. But it does add **tail risk**. If one shard in eight has a 1% chance of exceeding 40 ms, the probability that the request exceeds it is `1 - 0.99^8 = 7.7%`. Fan-out multiplies my p99 exposure, so a system that promises p99 under 100 ms either needs fewer shards on the hot path, hedged requests with cancellation, or a per-shard deadline that degrades to partial results.

## Sharding and replication

**Sharding** partitions the vector space or the key space so each node owns a slice. The two strategies behave differently under churn:

* Hash or range partitioning on document id is simple and rebalances predictably, but every query fans out to every shard, because a neighbour can be anywhere.
* Partitioning by vector space (IVF cell assignment, or a cluster-based split) means a query may touch only a handful of shards, which cuts fan-out and tail risk — at the cost of imbalance, since some regions of the space are far denser than others.

**Replication** is how I get read throughput and availability, not capacity. With `R` replicas of a shard, read QPS for that shard scales roughly `R x` until the network or CPU of a node saturates. Replicas also let me rebuild or patch the index one node at a time, which is the only way to change `M` or `ef` without downtime.

<Warning>
  Two consistency traps I plan for explicitly. First, a query that fans out across shards and merges top-k per shard is an approximation of global top-k: with `k=10` per shard across 8 shards I rerank 80 candidates for 10 slots, and if the per-shard `k` is smaller than the global `k` I can lose true neighbours permanently. Second, replication lag makes an answer reference a document the user just deleted. Retrieval freshness has to be part of the deletion path, not just the insert path.
</Warning>

## The hot-partition problem

The failure that hurts most at scale is not average load but skew. One tenant, product line, or time window accumulates most queries and most vectors, and whichever shard owns it becomes the bottleneck while the other nine idle.

Sources of hotness I have actually seen:

* **Tenant-keyed collections.** One enterprise customer with 40M documents and 80% of query volume: sharding by tenant id puts all of it on one node.
* **Time-windowed indexes.** "Latest reviews only" workloads concentrate on the newest segment, which is usually the least optimized and still being compacted.
* **Namespace or filter explosion.** A high-cardinality metadata field used as a pre-filter turns one index into thousands of tiny sub-indexes, and the popular ones dominate. This is where payload filtering stops being a feature and becomes a capacity constraint.
* **Duplicate-heavy corpora.** Near-duplicate embeddings cluster tightly, so the neighbourhood around a hot item gets traversed by nearly every query.

Mitigations, in the order I try them: replicate the hot shard rather than resharding (cheapest, fixes read pressure immediately); split the hot key across sub-shards by a second dimension such as document date; cache hot queries at the application layer — what [Caching LLM Chats to Quickly Answer User Queries Without RAG](/caching-llm-chats-to-quick-answer-user-query-without-rag) is for, and often a better answer than any index change; move to vector-space partitioning so dense regions get more shards; and finally rebalance on a schedule, which means moving terabytes.

<Tip>
  Before any of this, size the corpus. The number of chunks, the dimensionality, and the growth rate come out of the same arithmetic in [Data size estimation](/data-size-estimation). Do it first, because 1M vectors and 1B vectors are different products: the first is a library on one node, the second is a distributed system where the questions above decide the architecture. For the retrieval layer those numbers feed, see [RAGs](/rags) and [Case Study: building high-capacity vector databases with open-source techniques](/case-stude).
</Tip>

## The contenders, as I recorded them

* **Pinecone** — fully managed, cloud-native serverless, built for real-time applications without infrastructure overhead. I trade control and unit economics for not running the thing.
* **Milvus** — open-source and highly distributed, engineered to scale smoothly to billions of vectors, with optional GPU acceleration. Its storage/compute separation is what makes billion-scale tractable self-hosted.
* **Qdrant** — a Rust-powered engine with robust payload filtering and high-performance similarity matching. The filtering matters: index-respecting metadata constraints beat post-filter rescaping for both latency and recall.
* **Weaviate** — an open-source platform offering built-in vectorization and multi-modal search support. Vectorizing inside the database removes a class of drift bugs where the query and stored embeddings came from different model versions.

My version of this comparison: managed serverless wins when the team is small and traffic is spiky; Milvus at billion scale where I want the shards and the storage tier on my terms; Qdrant when the workload is filter-heavy; Weaviate when I want model and index co-versioned.


This documentation is built and hosted on [Mintlify](https://mintlify.com), a developer documentation platform.