Milvus 分片 Shard 机制:数据分片与查询协调节点的交互流程

发布时间:2026/9/15 1:49:03
Milvus 分片 Shard 机制:数据分片与查询协调节点的交互流程 Milvus 分片 Shard 机制数据分片与查询协调节点的交互流程在分布式向量数据库 Milvus 中当单集合Collection的数据规模突破数千万乃至数亿条高维向量时单台物理服务器的内存与算力已经无法容纳全量数据。为了实现水平横向扩展与极致的高吞吐检索Milvus 采用了基于**数据分片Sharding与分布式查询协调Query Coordination**的存储与计算分离架构。在创建集合时有一个至关重要的参数——shards_num分片数量通常设为 2、4 或 8。当一条包含 768 维向量与元数据的记录写入 Milvus 时它到底是如何被路由到特定物理分片的当客户端向 Proxy 协调节点发起一次并发向量搜索时分布式集群内部的QueryNode、Proxy、DataNode 与协调器之间经历了一套怎样精密的“广播扇出Fan-out与多路归并排序Reduce Sort”时序流转Milvus 分布式分片写入与存储拓扑Milvus 的分片基于主键Primary Key的一致性哈希Consistent Hashing[ 客户端写入请求: Insert(iddoc_8892, vector[...], text...) ] | v ------------------------ Proxy 接入协调节点 ------------------------ | 1. 计算主键哈希: Hash(doc_8892) % shards_num | | 2. 判定该数据归属于 【Shard_1】 | -------------------------------------------------------------------- | v 向对应的物理消息通道 (Pulsar/Kafka Topic) 发送写消息 ------------------------ 物理分片通道 (Physical Shard Channels) ------------------------ | [ Shard 0 Channel ] [ Shard 1 Channel ] [ Shard 2 Channel ] | ------------------------------------------------------------------------------------ | v (由 DataNode 消费并批量刷入对象存储 MinIO/S3) ------------------------ 段物理组织 (Segments in Shard 1) -------------------------- | [ Growing Segment 101 ] (正在内存中接收增量写入) | | [ Sealed Segment 102 ] (已写满封口构建 HNSW 索引并持久化至 S3) | | [ Sealed Segment 103 ] (已写满封口构建 HNSW 索引并持久化至 S3) | -------------------------------------------------------------------------------------在线分布式查询的四阶段流转时序Scatter-Gather 模型当用户发起一次search(vector, top_k10)时Proxy 协调节点与分布在不同物理机上的 QueryNode 之间执行经典的Scatter-Gather分散-聚合流程[ 客户端发起查询: Search(query_vec, top_k10) ] | v ------------------------- 阶段一: Proxy 协调节点广播扇出 (Scatter) ------------------------- | 1. 解析查询参数 (efSearch64, top_k10, 标量 Expr) | | 2. 将查询请求同时并发广播Fan-out给管理各个 Shard 的 QueryNode 节点 | ------------------------------------------------------------------------------------------ | ---------------------------------------------------- | | | v v v ---------------- QueryNode 1 ---------------- ---------------- QueryNode 2 ---------------- | 管理 Shard 0 中的多个 Segments | | 管理 Shard 1 中的多个 Segments | | 1. 并发在各个 Segment 内部执行 HNSW 图搜索 | | 1. 并发在各个 Segment 内部执行 HNSW 图搜索 | | 2. 在本地执行小顶堆局部归并排序 | | 2. 在本地执行小顶堆局部归并排序 | | 3. 产出 Shard 0 局部 Top-10 结果集 | | 3. 产出 Shard 1 局部 Top-10 结果集 | -------------------------------------------- -------------------------------------------- | | ------------------------------------------------ | v ------------------------- 阶段三: Proxy 全局多路归并排序 (Gather Reduce) ------------------------- | 1. Proxy 收集所有 QueryNode 返回的局部 Top-10 候选集 (共 shards_num * 10 20 条候选) | | 2. 在内存中利用优先级队列Priority Queue / K-way Merge执行全局归并排序 | | 3. 截断保留最终全局最优的 Top-10 结果 | | 4. 阶段四: 根据文档 ID 回表拉取元数据并返回客户端 | ---------------------------------------------------------------------------------------------------生产环境核心参数配置实战在创建集合时shards_num决定了系统未来的横向扩展潜力from pymilvus import Collection, CollectionSchema, FieldSchema, DataType def create_distributed_sharded_collection(collection_name: str cluster_kb_100m): fields [ FieldSchema(nameid, dtypeDataType.VARCHAR, is_primaryTrue, max_length64), FieldSchema(nametenant_id, dtypeDataType.VARCHAR, max_length32), FieldSchema(namevector, dtypeDataType.FLOAT_VECTOR, dim768) ] schema CollectionSchema(fields, description分布式分片海量知识库) # 核心参数shards_num (生产建议设为 2 的幂次方如 2、4 或 8) # 计算原则: 期望单 Shard 管理的向量数据量在 500 万到 2000 万之间最佳 collection Collection( namecollection_name, schemaschema, # 1 亿数据规模划分为 8 个分片每个分片管理约 1250 万向量 shards_num8 ) print(f✅ 成功创建分布式集合 [{collection_name}]分片数设定为: 8) return collection生产选型与调优三大黄金避坑军规shards_num必须在建表时确定创建后绝对无法在线修改如果建表时设了shards_num1未来数据量膨胀到 1 亿时单分片只能被分配在单台 QueryNode 上无法利用集群多机算力并发加速因此哪怕初期数据量小生产建表也必须至少设为shards_num2或shards_num4警惕分片数量过多引发的“小文件与网络广播灾难”如果数据只有 100 万却盲目配置了shards_num64每次极小的查询都会在内网向 64 个节点发起 RPC 广播与连接握手Proxy 节点的 CPU 全部浪费在 64 路网络包的解包与归并排序上查询延迟不仅没有下降反而从 3ms 恶化到 25ms结合物理副本Replicas实现读写分离一个 Shard 可以被加载到多个 QueryNode 形成多副本Replicas。在读多写少场景下增加副本数能够线性拉升 QPS 吞吐而增加 Shard 数主要用于解决海量数据的容量切分。总结分片机制是分布式向量数据库征服亿级海量数据的基石。“数据按主键一致性哈希落入 Shard查询由 Proxy 并发广播扇出并在内存执行高效 K-way 多路归并”理解了这套物理时序你就能在集群容量规划与扩容调优时做到游刃有余、胸有成竹。