The problem
A 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. Read that as three constraints that fight each other, and then read the word that does the damage: billions. At one million vectors, a vector database is a library. It runs in-process, fits on a laptop, and the only interesting question is recall. At one billion vectors it becomes a distributed systems problem with a hard latency ceiling, and the questions change completely: how much RAM do I need, how many nodes fan out per query, which shard is hot, what happens during a compaction, and can I afford to keep exact vectors online. Three things push teams off the managed-serverless path and onto open source at that size, and they are the real subject of this case study: unit economics at scale, control over index parameters and quantization, and data residency. The tradeoff is that I now own rebuilds, rebalancing, and on-call for an index type most of my organisation has never operated.Sizing it before choosing anything
The dimensionality and vector count decide the architecture, so I start with bytes.
Two conclusions fall out of that table before I have chosen a product.
First, 1B vectors at fp32 is 3.07 TB of RAM for the vectors alone, and the sub-100 ms promise means most of it must be resident and warm. On 64 GB nodes that is roughly 48 nodes of pure vector storage before I account for the graph, the segments, the operating-system reserve, or replication. That is the invoice that ends the “keep everything exact in memory” design.
Second, quantization changes the shape of the problem, not just its size. The same billion vectors are 768 GB at int8 and 96 GB as 96-byte product-quantization codes. int8 is about a dozen 64 GB nodes. PQ codes are a handful of nodes, or one large one. The catch, from Scaling a Vector Database, is that compression is paid for in recall, and I only find out how much by measuring it on a query set that looks like real traffic.
So the capacity plan for a billion-scale build is a two-stage layout: compressed codes in RAM for candidate generation, exact or int8 vectors on fast SSD for reranking a few hundred candidates per query. That is the single most consequential design decision in this case study, because it decouples my RAM bill from my corpus size.
Reference architecture, all open source
The stack I would assemble, and the one this site’s own tooling vocabulary assumes:1
Ingest and parse
Documents come from object storage or a change stream. A parser that understands structure — Docling for PDFs and office documents — matters more than people expect, because a mangled table becomes a vector nobody can retrieve. Parse, normalise, keep the source URI and version in metadata.
2
Chunk
250 to 400 tokens per chunk with a deliberate overlap. The overlap multiplies vector count by 1.2x, which is a cost line, not a rounding error.
3
Embed
An open embedding model served locally — BGE via sentence-transformers or Hugging Face — so that model version, dimensionality, and normalisation are under my control. Serving embeddings myself also keeps the corpus from leaving the building.
4
Index and store
Milvus when the target is genuinely billion-scale, because it is an open-source, highly distributed system engineered to scale smoothly to billions of vectors with optional GPU acceleration. Qdrant when the workload is filter-heavy, since it is a Rust-powered engine with robust payload filtering and high-performance similarity matching. Weaviate when I want built-in vectorization and multi-modal search so the query and the corpus are guaranteed to be embedded by the same model version.
5
Hybrid retrieve
BM25 for lexical precision on identifiers, prices, and error codes; dense retrieval for paraphrase; Reciprocal Rank Fusion to merge two incomparable score scales; a cross-encoder reranker on the fused top-n.
6
Serve
Asynchronous FastAPI behind a router that owns fan-out and merge, with Redis in front for exact and semantic answer caching — which is the cheapest capacity expansion available, since a cache hit costs nothing in the index.
7
Operate
Kubernetes or ECS/Fargate for the data tier, Terraform for the environment, GitHub Actions for CI, and OpenTelemetry, Prometheus, and Grafana for traces and dashboards. Without per-stage latency histograms I cannot find the shard that is slow, and that is the whole ballgame at this size.
The latency budget
Sub-100 ms is a budget, so I spend it explicitly:
The parallel-fan-out row hides the sharpest trap in the design. Adding shards does not add latency linearly, but it does multiply 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, or 7.7%. A p99 target across a fanned-out system is therefore not achievable by tuning per-shard p99 alone. I need hedged requests with cancellation, a per-shard deadline that degrades to partial results, or fewer shards holding more data each.
Recall is an architecture decision, not a tuning knob
Two-stage search gives me cheap capacity and a quality problem: compressed codes rank approximately, so recall@10 after the candidate stage caps the recall of everything downstream. The open-source answer is to widen the candidate pool and re-rank it, then make retrieval hybrid so that failures in one signal are covered by the other. From the retrieval pipeline in RAGs and the stack named in Sample Resume: hybrid BM25 plus dense retrieval, RRF fusion, cross-encoder reranking, metadata filtering, and confidence scoring — a configuration that improved retrieval precision by about 35% while significantly reducing hallucinations in the build that first made those claims. What I measure before calling it done:- recall@k against a brute-force ground truth on a sampled slice of the corpus. This is the only number that tells me what quantization actually cost.
- Latency histograms per stage, not per request, because the shard and the reranker fail differently.
- Recall under the real filter predicate. A filtered query that rescapes after the ANN search can quietly return three results instead of ten, and no aggregate metric will show it.
- Behaviour at the boundary, where IVF cells split neighbours and HNSW graph degree limits the reachable set.
HNSW, IVF, and ScaNN each trade tiny bits of accuracy for massive speed gains, but they fail in different places. HNSW’s memory is the graph itself and it degrades on deletion-heavy workloads. IVF loses evidence at cell boundaries no matter how much nprobe I buy. ScaNN buys recall at a given compression ratio, which is exactly the currency a billion-scale system spends. Choose against the workload’s failure pattern, not against a benchmark table.
The tradeoffs, stated plainly
Failure modes I expect to meet
- Hot partitions. One tenant or one time window owning most of both the vectors and the queries. The fix order is: replicate the hot shard, split by a second dimension such as ingestion date, cache hot queries at the application layer, then rebalance — because rebalancing means moving terabytes.
- Index rebuild cost. Changing
M,ef, or the quantization scheme is a rebuild of a billion-vector index, which is measured in days. I would rather over-provision the first choice than plan a clean migration. - Compaction and OOM. Segment merges spike memory. Nodes sized to exactly the steady-state footprint are nodes that fall over during their own maintenance.
- Embedding model drift. Query and corpus embedded by different model versions silently destroys recall, and nothing errors. Pin the version, store it in metadata, and refuse to serve across versions without re-embedding.
- Deletion and freshness. A replica that lags the write path serves vectors for documents a user already removed. The delete path has to be as designed as the insert path.