ARTICLE · INTELLIGENCE

战地情报 · 详情页

来自尧图项目组的一线实战观察与深度解析

Milvus Flush 全链路设计解析:从 SDK 请求到数据落盘的完整流程

Milvus Flush 全链路设计解析:从 SDK 请求到数据落盘的完整流程 Milvus Flush 全链路设计解析从 SDK 请求到数据落盘的完整流程【免费下载链接】milvusMilvus is a high-performance, cloud-native vector database built for scalable vector ANN search项目地址: https://gitcode.com/GitHub_Trending/mi/milvusFlush持久化刷写是 Milvus 中保证已插入数据可见且落入持久化存储的关键操作也是理解向量数据库写入链路、Segment 状态机与异步协调机制的理想入口。本文以仓库设计文档 docs/design-docs/design_docs/20211109-milvus_flush_collections.md 为骨架系统拆解 Milvus 2.0 中一次 Flush 请求从 SDK、Proxy、DataCoord 到 DataNode 的完整执行链路并结合当前仓库源码对照其实现演化读完后你将掌握 Flush 的 RPC 协议字段、Proxy 任务队列执行模型、Segment 状态机Growing/Sealed/Flushed以及时间戳水位推进这一核心协调机制。一、为什么需要 Flush理解 Milvus 的写入模型Milvus 面向海量向量数据的高吞吐写入天然采用了先写入内存、异步批量落盘的架构用户通过 SDK 插入的数据首先被组织进内存中的Growing Segment增长中的段Segment 内部的数据经过多级缓存与消息流MsgStream中转最终由 DataNode 以刷写Sync/Flush的方式写入对象存储S3/MinIO 等中的持久化文件落在对象存储上并完成元数据登记后的段才进入Flushed状态可以被查询节点加载、参与索引构建与压缩Compaction。因此在默认异步刷写策略下插入请求返回成功并不等同于数据已经持久化。Flush操作的使命就是为用户提供一个确定性边界调用它可以确保在此之前插入的、属于指定 Collection 的数据被密封并最终写入持久化存储。需要说明的是本文解析的 原设计文档 诞生于 Milvus 2.0 时代描述的组件拓扑Proxy / DataCoord / DataNode / MsgStream与异步协调思路至今仍是理解 Milvus 数据面架构的最佳教材文末会对照当前仓库源码说明该机制在新版本中的演化形态。二、Flush 全链路总览一次请求的四次跳跃一次 Flush 请求在 Milvus 中要穿越四个角色执行流程如下SDK → Proxy客户端通过 gRPC 调用Flush(FlushRequest)声明要刷写的collection_namesProxy 内部把 RPC 请求包装为FlushTask推入 DDL 任务队列由后台服务以PreExecute → Execute → PostExecute三阶段执行Proxy → DataCoordExecute阶段由 Proxy 转发内部FlushRequest请求 DataCoord密封Seal指定 Collection 下的所有 Growing SegmentDataCoord → DataNodeDataCoord 依据 DataNode 上报的消费时间戳判断安全时机通过FlushSegments通知 DataNode 把已密封的 Segment 落盘之后这些 Segment 的状态推进到Flushed。值得注意的是Flush 是异步语义。SDK 收到 Flush 响应时只代表段已被 DataCoord 密封真正的持久化尚未完成需要 SDK 配合轮询段状态才能确认最终落盘。三、第一跳SDK → Proxy 的 FlushRequest文档中定义的 SDK 侧接口如下对应 MilvusService 的公共 gRPC 服务service MilvusService { ... rpc Flush(FlushRequest) returns (FlushResponse) {} ... } message FlushRequest { common.MsgBase base 1; string db_name 2; repeated string collection_names 3; } message FlushResponse{ common.Status status 1; string db_name 2; mapstring, schema.LongArray coll_segIDs 3; }几个字段的语义要点base通用消息头携带消息 ID、时间戳TSO等可追踪信息collection_names一次可以批量刷写多个 Collection这是Flush支持批量操作的基础coll_segIDsmapstring, LongArray键是 Collection 名值是本次被密封的 Segment ID 数组。SDK 拿到这批 Segment ID 后正是靠它们去轮询状态详见第五节。从当前仓库源码可以印证这条链路的演进在 internal/proxy/task_flush.go 中flushTask直接内嵌了*milvuspb.FlushRequest保留了批量刷写多个 Collection 的能力其响应结构中携带的CollSegIDs、CollSealTimes、CollFlushTs等字段正是coll_segIDs这一设计在新协议中的延续与丰富。四、Proxy 侧FlushTask 与任务队列的三阶段执行Proxy 收到 Flush RPC 后会将其包装成 DDL 任务并交给统一的任务调度模型处理。设计文档给出了当时 Proxy 任务系统的核心抽象type task interface { TraceCtx() context.Context ID() UniqueID // return ReqID SetID(uid UniqueID) // set ReqID Name() string Type() commonpb.MsgType BeginTs() Timestamp EndTs() Timestamp SetTs(ts Timestamp) OnEnqueue() error PreExecute(ctx context.Context) error Execute(ctx context.Context) error PostExecute(ctx context.Context) error WaitToFinish() error Notify(err error) } type FlushTask struct { Condition *milvuspb.FlushRequest ctx context.Context dataCoord types.DataCoord result *milvuspb.FlushResponse }这段接口定义揭示了 Proxy 处理 DDL 的通用框架所有任务创建/删除 Collection、建索引、Flush 等都实现统一task接口具备TraceCtx链路追踪上下文、ID请求 ID、BeginTs/EndTsTSO 时间戳区间、WaitToFinish阻塞等待任务完成等能力FlushTask组合了FlushRequest与对 DataCoord 的 RPC client 引用任务被OnEnqueue入队后Proxy 有一个后台服务不断从 DDL 队列取出任务按PreExecute → Execute → PostExecute三阶段执行。设计文档明确指出FlushTask在三阶段中的行为阶段行为PreExecute不做任何事直接返回Execute通过 gRPC 向 DataCoord 发送内部FlushRequest并同步等待响应PostExecute不做任何事直接返回对照当前源码internal/proxy/task_flush.goflushTask的PreExecute与PostExecute同样为空实现但Execute已发生显著演化——新版 Proxy 在启用 Streaming 服务时会通过task_flush_streaming.go的sendManualFlushToWAL向 WAL 追加一条带BarrierTimeTick屏障时间戳的手动 Flush 消息再调用mixCoord.Flush取回已刷写段信息详见 internal/proxy/task_flush_streaming.go。这条WAL 屏障 时间戳对齐的新路径本质上是用流式消息系统强化了旧设计中密封即屏障的语义保证屏障之前的插入都已被处理。五、第二跳Proxy → DataCoord密封所有 Growing SegmentProxy 在Execute阶段向 DataCoord 发起内部 Flush RPC。该接口定义在 DatapbData 面协议中注意它与 SDK 侧协议不同内部协议直接使用dbID、collectionID数据库与集合的数字 ID而不携带 Collection 名字service DataCoord { ... rpc Flush(FlushRequest) returns (FlushResponse) {} ... } message FlushRequest { common.MsgBase base 1; int64 dbID 2; int64 collectionID 4; } message FlushResponse { common.Status status 1; int64 dbID 2; int64 collectionID 3; repeated int64 segmentIDs 4; }当前仓库中该协议的实体定义位于 pkg/proto/data_coord.proto其中DataCoordService仍保留rpc Flush(FlushRequest) returns (FlushResponse)第 32 行、rpc GetSegmentInfo(...)第 43 行、rpc FlushSegments(FlushSegmentsRequest) returns(common.Status)第 142 行等关键接口与设计文档一一对应。DataCoord 收到请求后执行最关键的一步SealAllSegments——将该 Collection 下所有仍处于Growing状态的 Segment 全部密封为Sealed且此后不再为这些 Segment 分配新的 ID主键/插入条目。随后 DataCoord 把密封好的 Segment ID 列表返回给 Proxy。当前实现可见于 internal/datacoord/services.goDataCoord.Flush会先分配一个时间戳作为timeOfSeal基准再调用flushCollection逐 vchannel 执行密封。其中有一处值得注意的实现细节在未启用 Streaming 服务的路径下internal/datacoord/services.go才会遍历VChannelNames调用segmentManager.SealAllSegments进行显式密封并且密封前会先采集各 vchannel 的 checkpoint保证 checkpoint 早于 Segment 的结束时间戳这是后续安全判定落盘的重要前提。六、异步语义衍生出的两个核心问题设计文档在这里特别点明了 Flush 异步模型的复杂性。SDK 收到 Flush 响应仅仅意味着DataCoord 已密封这批 Segment但此时存在两个悬而未决的问题段仍驻留内存被密封的 Segment 可能还在 DataNode 的内存中尚未真正写入对象存储已分配 ID 尚未消费完DataCoord 不再为密封段分配新 ID但此前已分配的 ID对应正在消息流中传播的插入数据可能仍未被 DataNode 消费完毕——如果贸然落盘会丢失这批仍在途的数据。设计文档给出的解法分别对应两套机制Segment 状态轮询与时间戳水位推进下面两节分别展开。七、问题一的解法GetSegmentInfo 轮询与 SegmentState 状态机针对段仍在内存问题SDK 需要周期性地向 DataCoord 发送GetSegmentInfo请求直到所有密封段都进入Flushed状态。相关协议如下service DataCoord { ... rpc GetSegmentInfo(GetSegmentInfoRequest) returns (GetSegmentInfoResponse) {} ... } message GetSegmentInfoRequest { common.MsgBase base 1; repeated int64 segmentIDs 2; } message GetSegmentInfoResponse { common.Status status 1; repeated SegmentInfo infos 2; } message SegmentInfo { int64 ID 1; int64 collectionID 2; int64 partitionID 3; string insert_channel 4; int64 num_of_rows 5; common.SegmentState state 6; msgpb.MsgPosition dml_position 7; int64 max_row_num 8; uint64 last_expire_time 9; msgpb.MsgPosition start_position 10; } enum SegmentState { SegmentStateNone 0; NotExist 1; Growing 2; Sealed 3; Flushed 4; Flushing 5; }SegmentInfo 中除了状态外还携带了分区的归属partitionID、所属写入通道insert_channel、行数num_of_rows、DML 位置dml_position等信息为客户端提供了判断是否可安全查询/落盘完成的依据。由此可提炼出 Milvus Segment 生命周期状态机Growing ──(Flush 密封 Seal)──▶ Sealed ──(数据落盘完成)──▶ Flushed │ ▼ Flushing落盘进行中用于内部流转Growing段正在内存中接收写入可追加数据Sealed段已被密封不再接受新数据但数据可能尚未落盘Flushing段正处于刷写过程中的中间状态Flushed段已持久化到对象存储这是客户端轮询期望到达的终态。从当前仓库实现看段状态的管理集中在 internal/datacoord/segment_manager.go 的SealAllSegments及其配套的元数据记录中而对 Flushed 状态的达成在新架构下已由internal/flushcommon/writebuffer写缓冲管理器与同步管理器协同推进段真正落盘后才会在元数据中被标记为 Flushed从而被下游的GetSegmentInfo观察到。八、问题二的解法DataNodeTtMsg 时间戳上报与水位判定第二个问题的本质是如何证明段内所有已分配 ID 的数据都已被 DataNode 消费。设计文档的解法依赖时间戳对账DataNode 每从 MsgStream 消费一个数据包就会向 DataCoord 上报一次该通道的消费时间戳协议消息为DataNodeTtMsg。message DataNodeTtMsg { common.MsgBase base 1; string channel_name 2; uint64 timestamp 3; }DataCoord 中有一个后台服务startDataNodeTsLoop专门处理DataNodeTtMsg其判定逻辑分两步从DataNodeTtMsg中提取channel_name过滤出挂载在该通道上的所有已密封 Segment用DataNodeTtMsg.timestamp与该 Segment 进入Sealed状态的时间戳做比较若消费时间戳已超过密封时间戳说明该 Segment 对应的全部 ID 都已被 DataNode 消费完毕消息流中没有遗漏数据此时才允许通知 DataNode 将段写入持久化存储。这个机制非常精巧——它把数据是否全部抵达 DataNode这样一个分布式一致性问题规约为消息通道消费水位是否越过密封屏障时间戳的比较运算。当前版本源码中DataCoord 的flushCollection在密封前先记录各 vchannel checkpoint、并以分配的时间戳作为timeOfSeal正是同一思想的延续checkpoint / 时间戳必须早于段的结束位点才能保证水位判定成立。九、最后一跳FlushSegments 通知 DataNode 落盘当水位判定通过后DataCoord 通过以下 RPC 通知 DataNode 执行真正的持久化service DataNode { ... rpc FlushSegments(FlushSegmentsRequest) returns(common.Status) {} ... } message FlushSegmentsRequest { common.MsgBase base 1; int64 dbID 2; int64 collectionID 3; repeated int64 segmentIDs 4; }DataNode 收到FlushSegments请求后会把对应的Sealed段数据含数据文件、统计信息、删除位图等写入对象存储并登记元数据随后段状态推进到Flushed——至此SDK 在轮询GetSegmentInfo时才能观察到终态一次 Flush 才真正闭环。该接口在当前协议中同样可见于 pkg/proto/data_coord.protoDataNode 服务侧。十、从 2.0 设计到当前代码的演化对照设计文档描述的是 Milvus 2.0 以DataCoord MsgStream DataNode 手动 Flush为主的架构。对照当前仓库这条链路的职责划分与核心协议被完整保留但实现载体发生了明显升级维度2.0 设计文档所述当前仓库形态Proxy 任务执行FlushTask入 DDL 队列三阶段执行flushTask保留三阶段框架见 internal/proxy/task_flush.go密封信号直接向 DataCoord 发 Flush RPC启用 Streaming 时先向 WAL 追加带屏障时间戳的手动 Flush 消息再同步 DataCoord见 internal/proxy/task_flush_streaming.go段落盘驱动DataCoord 依据DataNodeTtMsg水位调FlushSegmentsDataNode 侧落盘能力演化为 internal/flushcommon/writebuffer写入缓冲与 syncmgr 同步管理器协同实现协调组件DataCoord 单一协调演进为 MixCoord 聚合协调DataCoord.Flush服务端实现见 internal/datacoord/services.go从源码结构看无论架构如何演进密封Seal→ 确认水位时间戳对账→ 持久化Flush to storage→ 状态置为 Flushed这条流水线始终是 Flush 的核心骨架而Sealed/Flushing/Flushed的 Segment 状态机也依然是全链路对外呈现的统一事实模型。这种设计理念长期稳定、底层载体持续重构的特征使得该设计文档至今仍是理解 Milvus 数据写入与持久化机制不可替代的入门材料。延伸阅读若希望沿同一体系继续深入可参考仓库中与之配套的设计文档与源码DataNode 恢复设计理解 DataNode 重启后如何基于通道 checkpoint 恢复消费位点与 Flush 的水位判定互为表里DataNode 流图恢复设计DataNode 内部 flowgraph 管道的恢复机制创建索引设计Flush 到Flushed后段才具备参与索引构建的前提两者构成完整的数据生命周期数据面协议定义pkg/proto/data_coord.protoDataCoord Flush 服务端实现internal/datacoord/services.go、段管理internal/datacoord/segment_manager.goProxy 侧任务实现internal/proxy/task_flush.go 与 internal/proxy/task_flush_streaming.go。【免费下载链接】milvusMilvus is a high-performance, cloud-native vector database built for scalable vector ANN search项目地址: https://gitcode.com/GitHub_Trending/mi/milvus创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
RELATED READING

延伸阅读

更多一线实战笔记与深度复盘,助您持续精进