Vendor case study · Milvus 2.5
Milvus Architecture, Explained Interactively
Explore the Milvus 2.5 distributed architecture: trace inserts, searches, index builds and node recovery through the Proxy, coordinators, worker nodes, etcd, the log broker and object storage.
Before the map
The concepts behind the architecture
- Records & collections
- An entity is a record. Its vector is a numeric representation used for similarity search. A collection groups records under a schema: their field definitions.
- Partitions & segments
- A partition groups data within a collection. A segment is a batch of rows managed together. Growing segments accept inserts; sealed segments no longer do.
- Shards & channels
- A shard divides a collection’s write workload. Its virtual channel carries an ordered stream of changes. A physical channel is the underlying log topic, which several virtual channels can share.
- Indexes & replicas
- An index is a search structure built from the data. A replica is one loaded copy of a collection spread across QueryNodes. More nodes share one copy; more replicas create extra copies.
- Metadata & coordination
- Metadata describes schemas, data locations and cluster state. The control plane schedules and coordinates work; workers execute it. DDL means operations such as creating or dropping collections.
- Logs & persistence
- A mutation is an insert, upsert or delete. The log broker records mutations for readers; consumer lag is how far behind a reader is. A flush persists buffered data. Compaction rewrites segments to merge them and drop deleted rows.
Example: a collection with two shards has two virtual change streams. Those streams feed segments, and QueryNodes search the segments. Partitions organize the data, while replicas add extra loaded copies. Official terminology ↗
01 / Inside the cluster
Distributed architecture
Four layers: an access layer of stateless Proxies, a control plane of coordinators, independent worker pools, and durable dependencies that outlive every worker. Select a component, or trace a request through it.
Select a component for its role, dependencies and operating notes, or pick a request to trace.
Compute scales independentlySearch, ingestion and indexing each have their own worker pool.
Storage outlives workersDurable logs and object storage let a replacement node rebuild its state.
Coordinators assign the workThey manage placement and lifecycle; workers execute it.
02 / Follow a request
Watch the cluster work
Play a request from start to finish, or step through one handoff at a time. Select any component to pause and inspect it.
Insert request
What happens
Validate
The Proxy checks the schema and routes accepted rows toward collection shards.
The client uses one endpoint; the Proxy handles routing.
Play animates the handoffs; Next moves at your own pace.Learning animation · timing is illustrative
Every request, step by step
Insert data
- 1
Validate (Client → Proxy). The Proxy checks the schema and routes accepted rows toward collection shards. The client uses one endpoint; the Proxy handles routing.
- 2
Order & assign (RootCoord · DataCoord). Timestamp and segment allocations establish ordering and placement. Allocations are batched; this is not one coordinator round trip per row. Shared ordering lets independent workers agree on operation time.
- 3
Persist the log (Proxy → broker). Mutations enter the durable broker. An insert acknowledgement does not mean the segment index has been built or that every reader has caught up. Acknowledged, searchable and indexed are different milestones.
- 4
Consume in parallel (DataNode + QueryNode). DataNodes prepare durable files while QueryNodes add recent data to growing segments. The two consumers advance independently. Two consumers can process the same stream at different speeds.
- 5
Flush & progress (DataNode → storage). The background pipeline persists segment files. Read visibility depends on the chosen consistency level and consumer progress. Indexing and loading continue as background work.
Search vectors
- 1
Submit search (Client → Proxy). The collection must be loaded. The request carries a vector, filter, top-k and consistency level. Collection load state and request consistency both matter.
- 2
Route to shards (Proxy → QueryNodes). The Proxy routes work to the shard leaders in an available replica; every replica does not execute every request. A search uses the required shards, not every loaded replica.
- 3
Meet visibility (QueryNode + stream). If the query's guarantee timestamp is ahead of what the QueryNode has consumed, execution waits for the stream to catch up. Consumer lag can add waiting time before the search even runs.
- 4
Search segments (QueryNode workers). Workers filter and search the relevant growing and sealed segments in parallel. Loading from storage happens ahead of time, not on every search. Growing data and sealed data both contribute to the result.
- 5
Merge top-k (QueryNodes → Proxy). Partial results are reduced and merged into the final response. Large result sets can make the Proxy or the network the bottleneck. Large result sets make reduction and the network expensive.
Build an index
- 1
Schedule work (DataCoord). Once a segment is ready for indexing, DataCoord assigns a build task to an IndexNode. Scheduling and index construction are separate responsibilities.
- 2
Read segment (Storage → IndexNode). The worker reads persisted field data from object storage. Indexing is its own resource-intensive workload. Storage bandwidth can limit how fast builds progress.
- 3
Build index (IndexNode). Build-time index parameters determine construction work, memory demand and index structure. They are distinct from per-search recall controls. Build settings and search settings serve different purposes.
- 4
Persist artifacts (IndexNode → storage). Completed index files are written back to shared object storage, where any QueryNode can load them. Durable artifacts can be loaded by independent QueryNodes.
- 5
Load & hand off (QueryCoord → QueryNode). For a loaded collection, QueryCoord schedules the sealed, indexed segment and coordinates the handoff from growing data. A finished build and a ready query replica are separate checks.
Recover a QueryNode
- 1
Detect failure (Membership / etcd). A QueryNode failure is observed through cluster membership and health. Recovery time depends on detection and on spare capacity. Failure detection comes before reassignment.
- 2
Reassign work (QueryCoord). QueryCoord restores the desired placement. A healthy second replica keeps serving while the missing work is recovered. Spare capacity and healthy replicas decide the recovery options.
- 3
Reload history (Storage → QueryNode). Replacement workers load the required sealed segments and indexes. Storage bandwidth and warm-up affect readiness. Cold loading adds storage traffic and warm-up time.
- 4
Replay updates (Broker → QueryNode). The assigned channel resumes from its checkpoint and rebuilds recent state. That log history must still be retained by the broker. The broker must still hold the history a recovering node needs.
- 5
Verify service (Client → Proxy → QueryNode). Check loaded replica state and run real insert and search probes. A Running pod alone does not prove recovery. A Running pod alone does not establish successful recovery.
03 / Reference
Every component, and what to watch
The same details as the map, in one place.
Client / SDK · Your application
Application
Client / SDK
Expresses inserts, searches and schema operations against one service endpoint.
- Core function
- Sends vectors, filters and request options to a Proxy through a load balancer.
- Watch in production
- End-to-end latency, timeouts, payload size and retry volume.
Operator’s note
Bound retries and use backoff. An ambiguous timeout needs application-level reconciliation; repeated writes are not a universal exactly-once guarantee.
Under the hood
Collection load state, consistency level, top-k and output fields all affect the request. Compare client latency with server latency before adding workers.
Proxy · Route · reduce
01 / Access layer
Proxy
A stable entry point that hides a distributed cluster from the application.
- Core function
- Validates requests, routes mutations into the log broker, dispatches searches and merges partial results.
- Watch in production
- Request latency, CPU, rejection rate and result size.
Operator’s note
Compare the busiest Proxy with the fleet. Long-lived gRPC connections can concentrate requests behind connection-level balancing, and new pods may stay underused. Extra Proxies do not add QueryNode memory.
Under the hood
Proxies cache schema and routing information. Large top-k values, many output fields or large batches increase serialization and merge cost. Make sure the load balancer distributes requests, not just connections.
RootCoord · Schema · time
02 / Control plane
RootCoord
Provides a shared catalog and ordering service across independent workers.
- Core function
- Coordinates schema and access-control operations, allocates IDs and logical timestamps, and maintains metadata.
- Watch in production
- DDL failures, timestamp allocation latency, leader changes and etcd health.
Operator’s note
Coordinator HA is active/standby. Extra standbys improve recovery; they do not multiply write throughput.
Under the hood
RootCoord's timestamp service gives operations a comparable logical order. A standby must acquire leadership and reload metadata before serving.
QueryCoord · Place · balance
02 / Control plane
QueryCoord
Keeps the desired collection and replica placement aligned with available QueryNodes.
- Core function
- Schedules loading and release of segments, balances work, and coordinates the handoff from growing to sealed data.
- Watch in production
- Load progress, balancing activity, unavailable replicas and node membership.
Operator’s note
A healthy new pod is only the first step. Wait for its assigned segments and channels to become ready.
Under the hood
Placement depends on available nodes, replica groups, resource-group constraints and load. QueryCoord is a control-plane service, not the search executor.
DataCoord · Segment lifecycle
02 / Control plane
DataCoord
Coordinates how the write stream becomes durable, compacted, indexed segments.
- Core function
- Assigns channels (ordered streams of collection changes) and segments (batches of rows), and schedules flush, compaction and index tasks across workers.
- Watch in production
- Pending tasks, channel assignments, flush completion and compaction progress.
Operator’s note
In the 2.5 logical architecture, index coordination belongs here. IndexNodes execute the builds.
Under the hood
A stalled background pipeline may reflect worker capacity or object-store errors. Adding coordinators will not make a blocked upload or index build complete.
QueryNode · Growing + sealed
03 / Worker nodes
QueryNode
Provides the compute and working set used to answer searches.
- Core function
- Consumes recent mutations, loads historical segments, applies filters and searches vectors.
- Watch in production
- Search latency, memory, CPU, growing-segment size, segment loading and stream lag, even when only insert traffic increases.
Operator’s note
Adding nodes spreads one replica's work. Adding collection replicas creates additional searchable copies.
Under the hood
A shard leader (delegator) on a QueryNode coordinates work for its shard. Growing data may use an interim index in 2.5, so it is not always a brute-force scan. Durable state is external, but local memory, caches and mmap files still need recovery time.
DataNode · Flush · compact
03 / Worker nodes
DataNode
Turns a replayable mutation stream into durable segment files.
- Core function
- Consumes assigned channels, buffers rows, writes binlogs (durable segment files) and executes compaction.
- Watch in production
- Channel lag, buffer memory, flush duration and storage upload errors.
Operator’s note
More DataNodes help only when channel assignments and downstream storage allow parallel work.
Under the hood
Channel ownership is the unit of ingestion scheduling. Too few active channels or one hot channel can leave added nodes idle. Flush and compaction compete for CPU, memory and object-storage bandwidth.
IndexNode · Build indexes
03 / Worker nodes
IndexNode
Separates expensive index construction from latency-sensitive search.
- Core function
- Reads segment files, builds an index and writes index artifacts to object storage.
- Watch in production
- Pending and failed builds, task duration, worker memory and storage I/O.
Operator’s note
Add IndexNodes for queued builds only after ruling out failed tasks and storage bottlenecks.
Under the hood
The index type controls the work: graph construction and clustering cost differently. Temporary build memory can exceed the final index size. Search-time ef or nprobe is a different control from build-time parameters.
etcd · Metadata · leases
04 / Durable dependencies
etcd
Keeps a strongly consistent view of the cluster's metadata and membership.
- Core function
- Stores metadata, service registrations and coordination state.
- Watch in production
- Quorum, leader changes, allocated and in-use backend bytes, quota alarms, volume headroom and disk commit latency.
Operator’s note
Every etcd member holds the same keyspace, so extra members do not multiply metadata capacity. Compare allocated size with in-use size after history compaction.
Under the hood
Revision compaction retires old versions; defrag returns unused backend pages to disk. Neither removes current Milvus keys. Separate rootPath prefixes still share one etcd quota. Vectors live in object storage, so an etcd snapshot alone is not a collection backup.
Log broker · Pulsar / Kafka
04 / Durable dependencies
Log broker
Decouples producers from consumers while keeping a replayable mutation history.
- Core function
- Persists recent operations and feeds DataNodes and QueryNodes independently.
- Watch in production
- Publish and fetch delay, consumer backlog, storage, retention and per-broker pressure. On managed Kafka, tell network shaping apart from Kafka quota throttling.
Operator’s note
Retain enough history for the slowest recovering consumer. The broker's durability settings are part of the recovery design.
Under the hood
A virtual channel carries one collection shard's change stream. A physical channel is a broker topic that several virtual channels can share. More consumers cannot split an indivisible channel's work indefinitely.
Object storage · S3 / MinIO
04 / Durable dependencies
Object storage
Lets worker lifetimes and durable data lifetimes be independent.
- Core function
- Holds segment data and index artifacts shared by the compute layer.
- Watch in production
- Request latency, errors, throttling, capacity and credential validity.
Operator’s note
Do not manually delete files that look old. Milvus metadata, retention and garbage collection govern object lifetimes.
Under the hood
A warm loaded search may not fetch its whole segment again, but loading, recovery and indexing depend on object storage.