Case Study: Inside a Distributed Vector Database
How a distributed vector database separates storage from compute, how inserts and searches flow through it, and the production failure patterns that follow from that design, using Milvus as the worked example.
The earlier lessons treated the vector store as a box that holds an index. At hundreds of millions of vectors, that box is a distributed system in its own right, and interviewers increasingly ask how it works inside. This lesson takes one apart. It uses Milvus as the concrete example because its architecture is openly documented, but the same ideas appear in most systems built to scale past one node.
The core idea: separate storage from compute
A single-node vector store keeps everything on one machine: the write-ahead log, the data files, the index and the query engine. That machine is both the capacity limit and the failure domain.
Distributed vector databases split those responsibilities apart:
- Access layer. Stateless proxies are the single endpoint clients talk to. They validate requests, route writes and searches, and merge partial results. You can add or replace them freely.
- Coordinators. A small control plane that owns the catalog (schemas), hands out IDs and logical timestamps, decides which worker holds which data, and schedules background jobs. Coordinators assign work but do not do it.
- Workers, split by job. Query nodes serve searches from memory. Data nodes turn the incoming write stream into durable files. Index nodes build ANN indexes. Each pool scales independently, so an index-build backlog doesn't need more query nodes.
- Durable storage. A strongly consistent metadata store (etcd in Milvus), a log broker (Kafka or Pulsar) holding recent mutations, and object storage (S3 or MinIO) holding segment files and index artifacts.
Because every piece of durable state lives in the bottom layer, any worker can die and a replacement rebuilds itself: it loads its assigned segments from object storage and replays recent changes from the log.
The write path: log first
An insert does not go straight into an index:
- The proxy validates the rows, assigns them to a shard, and gets a logical timestamp so every worker agrees on the order of operations.
- The rows are published to the shard's channel in the log broker. At this point the insert is acknowledged.
- Two consumers read the same channel independently. Data nodes buffer rows and flush them to object storage as segment files. Query nodes add the rows to an in-memory growing segment, so recent data is searchable before it's indexed.
- When a segment is sealed and flushed, an index node builds its ANN index, and the query nodes swap the growing data for the indexed segment.
So "acknowledged", "searchable" and "indexed" are three different milestones, each owned by a different component. That is the root of most freshness questions: a read-your-writes guarantee means waiting for the query node to consume up to the write's timestamp.
The read path: scatter, wait, merge
A search request carries a vector, a filter, top-k and a consistency level:
- The proxy sends the search to the shard leaders in one loaded replica, not to every replica.
- Each query node checks its consumed position in the log against the timestamp the consistency level requires. Strong consistency may wait for the stream to catch up; bounded or eventual consistency accepts slightly stale data and doesn't wait.
- Each node searches its sealed (indexed) and growing segments in parallel and returns its own top-k.
- The proxy merges the partial results into the global top-k.
This is the scatter-gather from the sizing lesson, with one addition: search latency can be time spent waiting for freshness, not time spent computing distances.
Segments are the unit of everything
Data is managed in segments: batches of rows that move from growing to sealed, flushed, indexed and eventually compacted. Segments are immutable once sealed, which is why deletes are tombstones followed by background compaction. The freshness lesson steps through that lifecycle.
Segment count has costs of its own. Many small segments, often caused by forcing a flush after every small batch, mean more index builds, more scheduling and more metadata keys.
What goes wrong in production
The architecture explains the failure patterns operators actually hit. Four come up again and again:
| Symptom | Mechanism | Lesson |
|---|---|---|
| One proxy is overloaded while the fleet average looks fine | gRPC sends many requests over one long-lived connection, and a Layer 4 load balancer only places connections | Balance requests, not connections, and compare the busiest instance with the average |
| Search p99 rises during a bulk insert, with flat search QPS | Query nodes also consume the write stream to keep recent data searchable | Budget query-node memory and CPU for ingestion, not just for the loaded index |
| Everything lags when the log broker is throttled | Both consumers depend on the broker, and slow consumers delay freshness | Find the limiting broker resource before tuning the database |
| The metadata store fills up as tenants are added, while workers are idle | Every collection, partition and segment adds metadata, and every metadata replica holds the full copy | Choose the tenancy boundary deliberately: shared collections with a partition key, or separate clusters |
Multi-tenancy is a metadata problem
A collection per tenant looks clean, but each collection brings shards, partitions, segments and metadata. Milvus, for example, caps the total of shards × partitions across all collections at 65,536 by default. With 4,000 tenants at 2 shards and 4 partitions each, that is already 32,000 units, while query nodes may be almost idle.
The options, from lightest to heaviest isolation:
- Shared collection with a partition key. Compatible tenants share physical partitions. The application must enforce a tenant filter on every request, because the partition key is not a security boundary.
- Separate collections. Use these when tenants need different schemas or index settings. They are still one cluster with one control plane.
- Separate clusters with explicit tenant-to-cluster routing. This is the only option that separates metadata capacity and failure domains, and the application now owns routing and migration.
This mirrors the tenant isolation models in the permissions and multi-tenancy lesson, seen from the database's side.
Explore it hands-on
The Milvus Architecture Lab is an interactive version of this lesson for one real system. It has a clickable cluster map, animated insert, search, index-build and recovery flows, a replica memory model, and five production incident patterns traced from symptom to lesson.
Key takeaways
- Distributed vector databases separate stateless access, coordination, specialised workers and durable storage, so each can scale and fail independently.
- Writes go to a log first. Persistence and query serving consume that log independently, so "acknowledged", "searchable" and "indexed" are different milestones.
- Searches scatter to shard leaders in one replica, wait for the required freshness, then merge top-k results.
- Durable state lives in a metadata store, a log and object storage, which is why a replacement node can rebuild itself.
- The usual production limits are hot connections, inserts that load the query layer, and metadata growth, not raw vector math.