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, Databricks’ explainer). 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, Google Cloud). Quantization — Product Quantization, or PQ — compresses vectors to reduce memory consumption (Oracle on Pinecone-style scaling, part 1, part 2). Distributed scaling uses horizontal partitioning to split datasets across independent nodes or pods so queries run in parallel (Oracle). For the survey literature, this vector-database overview on arXiv and this enterprise selection guide 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.
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.
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.
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.
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=96on 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.
Sub-100 milliseconds, decomposed
“Sub-100 ms” is a budget, so I spend it line by line:
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.
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.
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.
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.