Milvus 架构设计:从向量索引到分布式检索,AI 原生存储引擎的内部机制
一、十亿向量检索的延迟困境:传统数据库为何无法承载相似度查询
随着大模型应用加速落地,向量检索已经成为 AI 系统中的关键基础设施。在 RAG(检索增强生成)场景中,系统往往需要在毫秒级时间内,从数百万到数十亿条文本向量中,快速找出与查询向量最相似的 Top-K 结果。这类“高维近似最近邻”(ANN)搜索,与传统数据库擅长的精确匹配、条件筛选和范围查询,在底层计算模型上有本质差异。

传统数据库依赖的 B+ 树索引,本质上仍是基于精确比较和范围扫描,对向量相似度检索几乎无能为力。核心原因在于“维度灾难”——在高维空间里,数据点之间的距离分布会逐渐趋于均匀,导致基于空间划分的索引结构(如 R 树)剪枝能力显著下降。举例来说,一张包含 1000 万条 768 维向量的数据表,如果采用暴力扫描方式计算余弦相似度,查询延迟通常会达到秒级,这显然无法满足在线检索服务对低延迟的要求。
Milvus 正是为解决大规模向量相似度搜索而设计的 AI 原生向量数据库。它从底层架构开始,就围绕向量检索的 I/O 模式、索引构建方式和计算特性进行优化,而不是在传统关系型数据库或通用存储引擎之上,简单叠加一个向量索引插件。
二、存算分离与段式存储:Milvus的分布式架构内核
Milvus 2.x 采用云原生的存算分离架构,核心组件可以概括为 Coordinator、Worker Node 和 Storage Layer 三层。要真正理解 Milvus 的分布式架构,最直观的方式就是从数据写入与查询的数据流路径切入。
flowchart TB
Client[客户端请求] --> Proxy[Proxy
请求路由与结果聚合]
Proxy --> QCoord[Query Coordinator
查询节点调度]
Proxy --> DCoord[Data Coordinator
数据节点调度]
Proxy --> ICoord[Index Coordinator
索引节点调度]
QCoord --> QNode[Query Node
向量检索执行]
DCoord --> DNode[Data Node
数据写入与持久化]
ICoord --> INode[Index Node
向量索引构建]
DNode --> WAL[Write-Ahead Log
消息队列 Kafka/Pulsar]
WAL --> QNode
DNode --> ObjStore[对象存储 S3/MinIO
段文件持久化]
INode --> ObjStore
QNode --> ObjStore
ObjStore --> SegData[段文件
binlog 格式]
SegData --> GrowSeg[增长段
可变,在内存中]
SegData --> SealSeg[密封段
不可变,已刷盘]
QNode --> Search[向量检索
IVF_FLAT / HNSW / SCANN]
QNode --> Retrieve[标量过滤
Prune + ANN]
style QNode fill:#e1f5fe
style INode fill:#e8f5e9
style DNode fill:#fff3e0
段式存储(Segment)是 Milvus 数据组织与管理的核心单元。每个 Collection 中的数据都会被切分为多个 Segment,每个 Segment 同时保存一组向量数据以及对应的标量字段。Milvus 中的 Segment 主要分为两种状态:
增长段(Growing Segment):用于持续接收新写入的数据,数据主要驻留在 Query Node 的内存中。其写入链路是:客户端 -> Proxy -> WAL(消息队列)-> Data Node 持久化到对象存储 -> Query Node 订阅 WAL 并在内存中构建增长段。
密封段(Sealed Segment):当增长段达到设定阈值后会被封存,数据变为不可修改。随后由 Index Node 为密封段构建向量索引,索引完成后再替换原始数据文件,最终由 Query Node 加载索引文件并对外提供向量检索服务。
这种架构设计的核心优势在于:写入链路与查询链路实现了解耦。写入过程主要依赖 WAL 追加和对象存储落盘,不会直接干扰在线检索;而查询过程则以加载对象存储中的索引文件到内存为主,也不会被高并发写入流量明显拖慢。这种存算分离设计非常适合高吞吐写入与低延迟搜索并存的 AI 检索场景。
向量索引的选择机制,是决定 Milvus 检索性能、召回率和资源消耗的关键。Milvus 支持多种 ANN 索引类型,不同索引分别在召回率、查询延迟、内存占用和构建成本之间进行权衡。
| 索引类型 | 召回率 | 查询延迟 | 内存占用 | 构建速度 | 适用场景 |
|---|---|---|---|---|---|
| FLAT | 100% | 高(暴力扫描) | 最高 | 无需构建 | 小数据集、精确结果 |
| IVF_FLAT | 90%-95% | 中 | 中 | 快 | 通用场景 |
| IVF_PQ | 85%-92% | 低 | 低 | 中 | 大规模、内存受限 |
| HNSW | 95%-99% | 极低 | 高 | 慢 | 低延迟要求 |
| SCANN | 93%-97% | 低 | 中 | 中 | Google推荐方案 |
HNSW(Hierarchical Navigable Small World)是当前工业界最常见的高性能向量索引方案之一。它的核心思想是构建一个多层导航图:底层包含全部向量节点,上层则进行逐层稀疏采样。查询时会从最稀疏的顶层开始搜索,再逐层向下逼近目标向量,逻辑上类似跳表结构。HNSW 的查询复杂度通常可近似为 O(log N),但代价是较高的内存开销,通常约为原始向量数据的 1.5-2 倍,因为还需要额外存储图结构的邻接关系。
三、生产级向量检索服务:Milvus集群部署与调优实践
下面这段代码展示了一个基于 Milvus Python SDK 封装的生产级向量检索服务,覆盖了连接池管理、Collection 创建、索引构建、混合检索以及性能监控等关键能力,适合作为向量数据库服务化落地的参考实现。
import time
import logging
import numpy as np
from typing import List, Dict, Optional, Tuple
from dataclasses import dataclass
from contextlib import contextmanager
from pymilvus import (
connections, Collection, FieldSchema, CollectionSchema,
DataType, utility, AnnSearchRequest, WeightedRanker,
)
from pymilvus.orm.types import CONSISTENCY_STRONG
logger = logging.getLogger("milvus_service")
@dataclass
class VectorSearchResult:
"""向量检索结果"""
ids: List[int] # 结果 ID 列表
distances: List[float] # 距离/相似度列表
scores: List[float] # 归一化分数 [0, 1]
latency_ms: float # 查询延迟(毫秒)
recall_hint: Optional[float] # 召回率估算(仅 FLAT 对比时可用)
class MilvusVectorService:
"""
Milvus 向量检索服务封装
支持连接池、索引管理、混合检索(向量 + 标量过滤)
"""
# 连接池配置
_connections = {}
_pool_lock = __import__('threading').Lock()
def __init__(
self,
alias: str = "default",
host: str = "localhost",
port: int = 19530,
collection_name: str = "embeddings",
dim: int = 768,
):
self.alias = alias
self.host = host
self.port = port
self.collection_name = collection_name
self.dim = dim
self._connect()
def _connect(self) -> None:
"""建立 Milvus 连接,支持连接复用"""
with self._pool_lock:
if self.alias not in self._connections:
try:
connections.connect(
alias=self.alias,
host=self.host,
port=self.port,
timeout=10,
)
self._connections[self.alias] = True
logger.info("Milvus 连接成功: %s:%d", self.host, self.port)
except Exception as e:
logger.error("Milvus 连接失败: %s", e)
raise
@contextmanager
def _get_collection(self) -> Collection:
"""获取 Collection 对象的上下文管理器"""
if not utility.has_collection(self.collection_name):
self._create_collection()
collection = Collection(self.collection_name)
try:
yield collection
finally:
# 释放 Collection 引用,不释放连接
pass
def _create_collection(self) -> None:
"""创建 Collection,定义向量字段与标量字段"""
fields = [
FieldSchema(name="id", dtype=DataType.INT64, is_primary=True, auto_id=True),
FieldSchema(name="embedding", dtype=DataType.FLOAT_VECTOR, dim=self.dim),
FieldSchema(name="source_id", dtype=DataType.INT64, description="原始数据源 ID,用于标量过滤"),
FieldSchema(name="category", dtype=DataType.VARCHAR, max_length=64, description="数据分类标签"),
FieldSchema(name="created_at", dtype=DataType.INT64, description="创建时间戳"),
]
schema = CollectionSchema(
fields=fields,
description="向量嵌入存储集合",
enable_dynamic_field=False, # 禁用动态字段,减少存储开销
)
collection = Collection(
name=self.collection_name,
schema=schema,
consistency_level=CONSISTENCY_STRONG,
)
logger.info("Collection 创建成功: %s", self.collection_name)
def create_index(
self,
index_type: str = "HNSW",
metric_type: str = "COSINE",
params: Optional[Dict] = None,
) -> None:
"""为向量字段创建索引。HNSW 默认参数:M=16, efConstruction=256"""
default_params = {
"HNSW": {"M": 16, "efConstruction": 256},
"IVF_FLAT": {"nlist": 1024},
"IVF_PQ": {"nlist": 1024, "m": 48, "nbits": 8},
}
if params is None:
params = default_params.get(index_type, {})
with self._get_collection() as collection:
index_params = {
"index_type": index_type,
"metric_type": metric_type,
"params": params,
}
collection.create_index(
field_name="embedding",
index_params=index_params,
)
logger.info("索引创建完成: type=%s, metric=%s, params=%s", index_type, metric_type, params)
def search(
self,
query_vector: np.ndarray,
top_k: int = 10,
filter_expr: Optional[str] = None,
search_params: Optional[Dict] = None,
output_fields: Optional[List[str]] = None,
) -> VectorSearchResult:
"""
执行向量检索
支持标量过滤(混合检索)和自定义搜索参数
"""
if query_vector.ndim == 1:
query_vector = query_vector.reshape(1, -1)
# HNSW 搜索参数:ef 越大召回率越高但延迟越大
if search_params is None:
search_params = {"metric_type": "COSINE", "params": {"ef": 128}}
default_output = ["source_id", "category", "created_at"]
if output_fields is None:
output_fields = default_output
with self._get_collection() as collection:
# 确保 Collection 已加载到内存
if not collection.is_loaded:
collection.load()
start_time = time.monotonic()
try:
results = collection.search(
data=query_vector.tolist(),
anns_field="embedding",
param=search_params,
limit=top_k,
expr=filter_expr,
output_fields=output_fields,
consistency_level=CONSISTENCY_STRONG,
)
except Exception as e:
logger.error("向量检索失败: %s", e)
raise
latency_ms = (time.monotonic() - start_time) * 1000
# 解析检索结果
ids, distances, scores = [], [], []
if results and len(results) > 0:
for hit in results[0]:
ids.append(hit.id)
distances.append(hit.distance)
scores.append(float(hit.distance))
logger.info("向量检索完成: top_k=%d, 返回=%d, 延迟=%.2fms", top_k, len(ids), latency_ms)
return VectorSearchResult(
ids=ids, distances=distances, scores=scores,
latency_ms=latency_ms, recall_hint=None,
)
def hybrid_search(
self,
query_vectors: List[np.ndarray],
weights: List[float],
top_k: int = 10,
filter_expr: Optional[str] = None,
) -> VectorSearchResult:
"""
多向量混合检索(多路召回 + 加权融合)
典型场景:文本向量 + 图像向量联合检索
"""
search_requests = []
for i, vec in enumerate(query_vectors):
if vec.ndim == 1:
vec = vec.reshape(1, -1)
req = AnnSearchRequest(
data=vec.tolist(),
anns_field="embedding",
param={"metric_type": "COSINE", "params": {"ef": 128}},
limit=top_k * 2,
expr=filter_expr,
)
search_requests.append(req)
with self._get_collection() as collection:
if not collection.is_loaded:
collection.load()
start_time = time.monotonic()
try:
results = collection.hybrid_search(
reqs=search_requests,
ranker=WeightedRanker(*weights),
limit=top_k,
output_fields=["source_id", "category"],
)
except Exception as e:
logger.error("混合检索失败: %s", e)
raise
latency_ms = (time.monotonic() - start_time) * 1000
ids, distances, scores = [], [], []
if results and len(results) > 0:
for hit in results[0]:
ids.append(hit.id)
distances.append(hit.distance)
scores.append(float(hit.distance))
return VectorSearchResult(
ids=ids, distances=distances, scores=scores,
latency_ms=latency_ms, recall_hint=None,
)
def insert(
self,
embeddings: np.ndarray,
source_ids: List[int],
categories: List[str],
timestamps: Optional[List[int]] = None,
) -> List[int]:
"""批量插入向量数据,返回插入的 ID 列表"""
if timestamps is None:
timestamps = [int(time.time())] * len(source_ids)
data = [
embeddings.tolist(),
source_ids,
categories,
timestamps,
]
with self._get_collection() as collection:
try:
insert_result = collection.insert(data)
# 触发 Flush 确保数据持久化(生产环境按需调整频率)
collection.flush()
logger.info("向量插入完成: count=%d", insert_result.insert_count)
return insert_result.primary_keys
except Exception as e:
logger.error("向量插入失败: %s", e)
raise
# ==================== 使用示例 ====================
def demo():
"""演示 Milvus 向量检索服务的完整工作流"""
service = MilvusVectorService(
host="localhost", port=19530, collection_name="demo_embeddings", dim=768,
)
# 创建 HNSW 索引
service.create_index(index_type="HNSW", metric_type="COSINE")
# 批量插入向量
num_vectors = 10000
embeddings = np.random.randn(num_vectors, 768).astype(np.float32)
# 归一化向量(余弦相似度要求)
embeddings = embeddings / np.linalg.norm(embeddings, axis=1, keepdims=True)
source_ids = list(range(num_vectors))
categories = ["tech"] * (num_vectors // 2) + ["science"] * (num_vectors - num_vectors // 2)
service.insert(embeddings, source_ids, categories)
# 单向量检索
query = np.random.randn(768).astype(np.float32)
query = query / np.linalg.norm(query)
result = service.search(
query_vector=query, top_k=10, filter_expr='category == "tech"',
)
print(f"检索结果: {len(result.ids)} 条, 延迟: {result.latency_ms:.2f}ms")
if __name__ == "__main__":
logging.basicConfig(level=logging.INFO)
demo()
这段代码中有几个生产环境里非常关键的设计点。首先,连接池通过类级别字典统一管理,避免每次读写或搜索都重新建立 Milvus 连接,从而降低连接抖动和额外开销。其次,`hybrid_search` 方法支持多路召回后的加权融合,这正是 RAG 系统、语义搜索系统和多模态检索中常见的实现方式,例如“文本向量 + 关键词向量”或“文本向量 + 图像向量”的联合检索。再者,搜索参数中的 `ef=128` 实际上是 HNSW 的搜索宽度控制参数——ef 越大,召回率通常越高,但查询时延也会同步上升。128 这个经验值通常是在千万级 768 维向量数据集的压测中,在约 95% 召回率与 10ms 左右延迟之间取得较好平衡的配置。
四、Milvus的性能边界与架构权衡
向量数据库的架构设计,本质上是在召回率、查询延迟、吞吐能力和资源成本之间寻找平衡点,Milvus 同样无法回避这些工程上的取舍。
**召回率与延迟之间很难同时极致优化。** 以 HNSW 为例,ef 参数直接决定搜索宽度。通常 ef=64 时,Top-10 召回率大约在 90% 左右;而 ef=256 时,召回率可提升到约 98%,但查询延迟往往会增加 3-4 倍。因此在生产环境中,必须依据业务对搜索准确性的真实需求来设置 ef。比如推荐系统、内容分发场景通常 90%-95% 的召回率已经足够,而法律检索、风控审查、合规分析等场景则可能需要接近 99% 的高召回。
**内存消耗是向量检索系统最硬的资源约束。** HNSW 索引通常需要整体加载到内存中才能发挥低延迟优势。以 1 亿条 768 维 float32 向量为例,原始数据体积大约为 286GB,再叠加 HNSW 图结构后,整体内存需求可能接近 500GB。对于超大规模数据集,通常不得不采用 IVF_PQ 等压缩型索引,以降低内存占用和节点成本,但压缩也会带来大约 5%-10% 的召回率损失。也就是说,尽管 Milvus 的存算分离架构降低了底层存储成本,但在实际查询时,索引仍然必须加载到 Query Node 内存,因此内存开销依旧是 Milvus 集群中最重要的成本项之一。
**标量过滤的执行顺序会直接影响搜索性能与结果完整性。** Milvus 支持在向量检索过程中附加标量过滤条件,例如 category、时间戳、业务状态等字段。但过滤发生在检索前还是检索后,会显著影响最终效率。Pre-filter 模式会先根据标量条件缩小候选集,再执行 ANN 搜索,适合过滤后数据规模大幅缩减的情况;Post-filter 模式则是先进行向量召回,再对结果做过滤,更适合条件宽松的场景。Milvus 2.x 默认采用 Post-filter,当过滤条件过于严格时,比如排除掉 99% 的候选数据,就可能导致最终返回结果不足 Top-K,因此实际部署时需要结合业务特征评估过滤策略。
**一致性级别的配置也会影响检索时延。** Milvus 支持强一致、有界一致、会话一致和最终一致四种一致性级别。强一致能够保证数据写入后立刻可查,但代价通常是额外增加 10-20ms 左右的延迟,因为系统需要等待 WAL 同步完成。对于大多数 RAG、语义搜索和知识库问答场景来说,会话一致通常已经足够——即同一客户端写入后可以立即检索到自身写入的数据,但不保证其他客户端在同一时刻也能马上可见。
五、总结
Milvus 的架构设计体现了 AI 原生存储引擎的几个核心理念:存算分离、段式管理、索引可插拔以及面向高维向量检索的专用优化。对于低延迟、高并发相似度搜索场景,HNSW 往往比 IVF 系列索引表现更好,但它的高内存占用是无法回避的约束;而 IVF_PQ 更适合超大规模、内存受限的向量数据库部署场景,不过需要接受一定程度的召回率下降。生产环境中的关键决策通常包括:根据数据规模与召回目标选择合适的向量索引类型,根据查询 QPS 和并发负载规划 Query Node 数量,并结合标量过滤条件的严格程度设计合理的检索策略。向量数据库并不是传统数据库的替代者,而是 AI 应用、RAG 系统、语义检索与推荐系统中,专门为高维近似最近邻搜索优化的基础存储引擎。只有真正理解它的 I/O 模型、索引机制和分布式计算特性,才能做出更合适的架构选型与性能调优决策。
