news 2026/9/9 13:16:30

Milvus DataNode Flowgraph 恢复机制设计深度解析

作者头像

张小明

前端开发工程师

1.2k 24
文章封面图
Milvus DataNode Flowgraph 恢复机制设计深度解析

Milvus DataNode Flowgraph 恢复机制设计深度解析

【免费下载链接】milvusMilvus is a high-performance, cloud-native vector database built for scalable vector ANN search项目地址: https://gitcode.com/GitHub_Trending/mi/milvus

本设计文档梳理了 Milvus 分布式数据写入链路中 DataNode 侧的核心话题:当一个 vchannel 由 DataNode 以 flowgraph(流图)方式消费时,如何基于位置(position)检查点保存各 segment 的消费进度,并在节点重启、重新拉流后实现精确恢复且不重复处理已刷新数据。读完本文,你将掌握 flowgraph 恢复的整体流程、WatchDmChannels协议的消息结构与"按位置过滤 msgPack"的伪代码思路,并能在当前仓库源码与流图实现中找到这套设计的历史形态与后续演进痕迹。

本文主体依据仓库中的设计文档 20210604-datanode_flowgraph_recovery_design.md(2021 年 6 月成稿)整理,并结合当前仓库源码进行印证与扩充。

1. 设计背景:为什么需要 Flowgraph 恢复

在 Milvus 的数据写入架构中,DataNode 是真正把"消息队列中的数据"落盘为"可检索的 segment 与 binlog"的执行组件。其内部以 flowgraph 为基本执行骨架:一个 vchannel 对应一个 flowgraph,消息流在其中依次经过输入、DML 处理、写入(写 segment / binlog)等节点。该骨架在演进后的代码中位于 internal/flushcommon/pipeline,每个 vchannel 对应一个 DataSyncService,并由 FlowgraphManager 统一管理(AddFlowgraph/RemoveFlowgraph/ClearFlowgraphs/GetFlowgraphService等)。

DataNode 是可能崩溃、重启、被重新调度到其他节点的无状态化组件,因此它必须回答一个核心问题:

当一个 DataNode 在某个 vchannel 上重启消费时,如何知道该 vchannel 上哪些 segment 已经写过哪些数据、应该从哪条消息继续消费、哪些消息需要过滤掉?

这就是 flowgraph recovery 设计要解决的事。整个恢复依赖的并非 Kafka 等消息队列自带的 offset 语义,而是 Milvus 自己的position 检查点体系。

1.1 几条设计常识(Common Sense)

原设计文档开篇先固定了四条贯穿全局的常识性约束,它们决定了恢复方案的形态:

  • A. 一个 message stream 对应一个 vchannel,因此一个 msgPack(一批消息)里只有一个起始位置和一个结束位置。这是后面"用单个 position 描述一段消息"的前提。
  • B. 只有在 DataNode flush 时,DataNode 才会更新每个 segment 的 position。flush 是 checkpoint 的天然锚点。
    • 伴随的一个优化点是:flush 时需要更新的位置包括两类——① 当前正在 flush 的 segment;② 从未被 flush 过的 segment 的 StartPosition。
  • C. DataNode 的自动 flush(auto-flush)同样是合法、有效的一次 flush,它同样要遵守位置更新规则,而不是被当作"临时写盘"。
  • D. DDL 消息(create/drop collection 等)此时也承载在 DML vchannel 中,因此 flowgraph 在处理 DML 数据流时必须同步处理混入的 DDL 事件。这条常识在后续实现中体现为一个独立的 DDL 处理节点 flow_graph_dd_node.go。

提示:文档中称 "message stream" 与 "vchannel",反映了当时消息系统(如 Pulsar/Kafka)与内部虚拟通道的映射关系;在后来版本中数据流进一步抽象为 DmChannel,但"一条逻辑通道被一个 flowgraph 消费、以 position 描述消费进度"的基本模型被保留下来。

2. Flowgraph 中的 Segment 状态模型

原文档的"Segments in Flowgraph"一节用一张状态图说明 flowgraph 内存中的 segment 并非铁板一块,而是处在三种状态之一,每种状态拥有不同的 checkpoint 语义:

  • New Segment:DataNode 第一次遇到、尚未写入任何数据的段。它携带的 checkpoint 是起始位置(start position),表示"从这个位置之后的数据都属于我"。
  • Normal Segment:正在持续接收数据、尚未 flush 的段。它同样维护起始位置与当前数据行数信息。
  • Flushed Segment:已完成一次合法 flush(含 auto-flush)的段,其 binlog 已持久化、checkpoint 已更新,后续 msgPack 到达时不再向它追加。

三者的核心区别可以浓缩为一组标记:isNew(是否新建)、isFlushed(是否已 flush)。段的生命周期即通过 flush 操作在这几种状态间迁移:

  • New Segment 积累数据后成为 Normal Segment;
  • Normal Segment 在触发 flush(含 auto-flush)后成为 Flushed Segment;
  • New Segment 也可以直接因 flush 而结束,未必经过长久的 Normal 阶段。

为什么要区分状态?因为恢复时要分别为不同状态的 segment 保存/使用不同含义的 position(见下一节),New 段的 checkpoint 代表它的起始消息位置,而 Flushed 段则意味着"它占用的消息区间已全部持久化,无需重放"。

3. Flowgraph 恢复(Recovery)设计

恢复机制由"保存检查点"与"从检查点恢复"两段组成。

3.1 A. 保存检查点(Save checkpoints)

每当一个 flowgraph flush 一个 segment 时,需要把以下信息保存下来:

  1. 当前 segment 的 binlog 路径——flush 产物落在对象存储/文件系统上的位置,供后续索引构建与查询读取;
  2. 当前 segment 的 position——即这个 segment 在 vchannel 消息流里"消费到了哪里";
  3. 从 replica(副本/元数据视角)中获得的其他所有 segment 的当前 position——其中有一条重要的补充规则:

    如果一个 segment还没有被 flush 过,则保存 DataNode第一次遇到它时的 position(即它的起始位置,而不是"当前读到的位置")。

保存成功与否的处理策略(原文档规定的兜底语义):

  • 保存成功:flowgraph 将所有这些 segment 的 position 统一更新到 replica(即向协调层上报最新进度)。
  • 保存失败
    • 若是gRPC 失败(此类失败在内部经过多次重试后才浮出到此处),则直接崩溃(crash itself)——因为说明与协调层已经失联,继续运行会造成进度与元数据不一致;
    • 若是普通失败,则重试保存 10 次,若仍失败,同样崩溃自身

这套"崩溃而非静默继续"的策略体现了 DataNode 的定位:它是一个可以从 checkpoint 重新恢复、具备确定性语义的无状态消费者,"宁可重启重建,也不在状态不确定时继续写数据"。

3.2 B. 从一组检查点恢复(Recovery from a set of checkpoints)

恢复的第一步是拿到该 vchannel 内所有 segment 的 position,记作p1, p2, ..., pn。这些 position 通过协调层下发的 Watch 请求带给 DataNode。原设计文档给出了当时设想的WatchDmChannelReq协议:

message VchannelInfo { int64 collectionID = 1; string channelName = 2; msgpb.MsgPosition seek_position = 3; repeated SegmentInfo unflushedSegments = 4; repeated int64 flushedSegments = 5; } message WatchDmChannelsRequest { common.MsgBase base = 1; repeated VchannelInfo vchannels = 2; }

消息语义很清楚:

  • VchannelInfo描述一个待 watch 的 vchannel:它属于哪个 collection、通道名是什么;
  • seek_position是 DataNode 开始消费的起点位置;
  • unflushedSegments给出尚未 flush、需要继续接收数据的 segment 及其元数据(其中包含各自的起始 position,即文档常识 B 与 3.1 中强调的"第一次遇到时的 position");
  • flushedSegments给出已经完成 flush、不需要再接收该消息区间内数据的 segment 集合。

第二件事是基于这些 position 过滤 msgPack——这正是恢复的精髓。原文档配图如下:

过滤算法

假设某 vchannel 内有三个 segments1, s2, s3,对应的位置分别是p1, p2, p3,且这些位置代表了各 segment 在消息时间轴上的覆盖区间。恢复流程:

  1. 将位置按逆序排序p3, p2, p1
  2. 求出各 segment 的重复区间范围
    • s3对应区间(p3 > mp_px > p1)
    • s2对应区间(p2 > mp_px > p1)
    • s1对应区间为zero(从最早开始、无下界);
  3. 从最早的 position 开始 seek,此例即从p1开始拉取消息;
  4. 对于 seek 之后到达的每一个 msgPack,用如下伪代码做过滤:
const filter_threshold = recovery_time // mp means msgPack for mp := seeking(p1) { if mp.position.endtime < filter_threshold { if mp.position < p3 { filter s3 } if mp.position < p2 { filter s2 } } }

理解这段逻辑的关键:

  • filter_threshold = recovery_time:只有结束时间早于恢复时刻的消息才需要考虑去重——恢复时刻之后的新消息天然没有历史包袱,直接正常消费即可;
  • seek 起点选最早p1,是为了不漏掉任何一个还需要继续累积数据的 segment(例如尚未 flush 的s1)的起始数据;
  • 对每条"历史"msgPack,逐级与各已覆盖位置的 segment 比较:凡是落在某个已 flush segment 的时间区间内(mp.position < p3归入s3的覆盖域,mp.position < p2归入s2的覆盖域),就说明这段数据之前已经写进该 segment 并被 flush 持久化过了,本次恢复时应当过滤丢弃,避免重复落盘;
  • 反过来,那些不属于任何已 flush segment 覆盖区间的 msgPack,则正是未 flush 段在重启后缺失的增量数据,需要继续喂给 flowgraph。

一句话概括:"从最早的起点重放,把已经 flush 过的消息区间剪掉,只补尚未持久化的增量。"

4. 与当前仓库代码的对照:设计的落地与演进

这份设计形成于 2021 年 6 月(Milvus 2.0 时代)。以当前仓库为准做对照,可以看到核心模型被保留、而具体实现有清晰演化,以下对照均可在仓库源码中直接验证:

4.1 协议层的演进

上述VchannelInfo/WatchDmChannelsRequest完整保留在 pkg/proto/data_coord.proto 中,但字段已大幅扩充:

message VchannelInfo { int64 collectionID = 1; string channelName = 2; msg.MsgPosition seek_position = 3; repeated SegmentInfo unflushedSegments = 4 [deprecated = true]; repeated SegmentInfo flushedSegments = 5 [deprecated = true]; repeated SegmentInfo dropped_segments = 6 [deprecated = true]; repeated int64 unflushedSegmentIds = 7; repeated int64 flushedSegmentIds = 8; repeated int64 dropped_segmentIds = 9; repeated int64 indexed_segmentIds = 10 [deprecated = true]; repeated SegmentInfo indexed_segments = 11 [deprecated = true]; repeated int64 level_zero_segment_ids = 12; map<int64, int64> partition_stats_versions = 13; msg.MsgPosition delete_checkpoint = 14; }

可见演进方向与设计文档一脉相承但更进一步:

  • seek_position被保留,仍然是 DataNode 恢复消费的基准位置;
  • 携带完整SegmentInfounflushedSegments/flushedSegments/dropped_segments均标记为deprecated = true,改为只下推 ID 列表(unflushedSegmentIds/flushedSegmentIds/dropped_segmentIds),元数据细节改为按需查询——这与"传输轻量化"的思路一致;
  • 后续新增了 L0(level zero)segment、分区统计版本号、delete_checkpoint等新语义字段,说明原设计框架("协调层下发 vchannel 全貌 + DataNode 按位置恢复")足够稳定,可以继续容纳新特性;
  • WatchDmChannelsRequest消息体(base + 一组VchannelInfo)与设计文档一致,在 pkg/proto/data_coord.proto 与 pkg/proto/query_coord.proto 中都有出现,说明 query 侧也复用了同一套 watch 协议。

4.2 流图实现层的演进

设计中的"一个 vchannel 一个 flowgraph"在代码中落地为DataSyncService与一套节点化流水线,集中在 internal/flushcommon/pipeline:

  • 输入节点从 DmChannel 拉取消息包,对应 flow_graph_dmstream_input_node.go;
  • DDL 与时间同步(time tick)作为并行节点存在,对应 flow_graph_dd_node.go 与 flow_graph_time_tick_node.go——与设计常识 D"DDL 消息在 DML vchannel 中"相印证;
  • 写入节点负责将数据累积进 segment 并触发 flush,对应 flow_graph_write_node.go;
  • flowgraph 的注册、查询与清理由 FlowgraphManager 承担,每个 channel 以typeutil.ConcurrentMap[string, *DataSyncService]维护,AddFlowgraph时以 vchannelName 为键并上报DataNodeNumFlowGraphs指标;
  • 每个 DataSyncService 内部持有 metacache,维护 channel 内各 segment 的实时状态(SegmentID / PartitionID / State / NumOfRows / BufferRows 等),这一点与设计中"segment 的当前 position 需要从内存元数据里收集"的要求一致,参见 flow_graph_manager.go。

4.3 DataNode 侧 RPC 的现状

需要特别指出:设计文档描述的WatchDmChannels行为属于旧版消费链路。在当前仓库中,DataNode 暴露的 gRPC 服务里该接口已被弃用——internal/datanode/services.go 的WatchDmChannels实现直接打印"DataNode WatchDmChannels is not in use"并返回成功;同文件的 FlushSegments 也标注deprecated after v2.6.0。也就是说,v2.6.0 之后 DataNode 的 DML 通道由新的流式/增量订阅链路接管,老 Watch 协议仅保留接口与数据模型层面的兼容。

因此在阅读本设计文档时,建议把它定位为一份理解 vchannel 恢复语义的权威入门材料:无论消费入口怎么迁移,"以 position 作为 checkpoint、flush 作为进度上报锚点、重启后从最早 checkpoint seek 并按已 flush segment 区间过滤重放数据"这套核心思想始终是 DataNode 数据一致性的地基。

5. 小结:可复用的恢复要点清单

关注点设计结论仓库印证位置
通道与流图一个 vchannel 由一个 flowgraph(DataSyncService)消费internal/flushcommon/pipeline/data_sync_service.go、flow_graph_manager.go
进度载体position(含 seek_position)是 checkpoint 基本单位pkg/proto/data_coord.proto
上报时机只在 flush 时更新 segment position;auto-flush 同样有效原文档常识 B/C
保存内容当前段 binlog 路径与 position + 其他段的当前/起始 position原文档 §3.A
失败语义gRPC 失败直接崩溃;普通失败重试 10 次后崩溃原文档 §3.A
恢复方式收集全量 segment position → 逆序排序求覆盖区间 → 从最早 seek → 按recovery_time阈值过滤 msgPack原文档 §3.B 与伪代码
DDLDDL 与 DML 同通道,flowgraph 内设 DDL 节点处理flow_graph_dd_node.go

对于希望深入 DataNode 数据链路(写入、flush、恢复、去重)的开发者,建议按此顺序阅读:先吃透本设计文档的 position/checkpoint 语义,再到 internal/flushcommon/pipeline 通读各节点与 DataSyncService,最后对照 pkg/proto/data_coord.proto 理解协议字段的演进,即可形成一条从"设计意图"到"现代实现"的完整认知链。

【免费下载链接】milvusMilvus is a high-performance, cloud-native vector database built for scalable vector ANN search项目地址: https://gitcode.com/GitHub_Trending/mi/milvus

创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考

版权声明: 本文来自互联网用户投稿,该文观点仅代表作者本人,不代表本站立场。本站仅提供信息存储空间服务,不拥有所有权,不承担相关法律责任。如若内容造成侵权/违法违规/事实不符,请联系邮箱:809451989@qq.com进行投诉反馈,一经查实,立即删除!
网站建设 2026/9/9 13:15:46

AI时代,Lisp为何是程序员保持独立思考的自留地?

最近在团队里聊AI编程&#xff0c;大家张口闭口都是Agent、Copilot、上下文窗口&#xff0c;我一个人在角落里用Emacs写着Lisp&#xff0c;被旁边的同事瞥了一眼&#xff0c;问了句&#xff1a;“你写这玩意儿&#xff0c;图啥&#xff1f;”我愣了一下&#xff0c;没接话。但这…

作者头像 李华
网站建设 2026/9/9 13:14:35

GESP C++ 5级备考指南:从数组指针到递归算法全解析

1. 为什么 5 级是 GESP 的分水岭——先看清考试定位1.1 5级到底考什么很多刚开始准备 CCF GESP 的同学&#xff0c;一上来就问“5级难不难”。我的答案是&#xff1a;它比 3 级、4 级难出一个明显的台阶&#xff0c;但还不至于像 7 级、8 级那样需要系统学完算法竞赛入门。说白…

作者头像 李华
网站建设 2026/9/9 13:14:08

理解Magnitude:从星等、震级到算法复杂度的量级思维

我第一次把 magnitude 这个词当回事&#xff0c;是在一台口径 25 厘米的望远镜前面。当时的想法很简单&#xff1a;为什么天文台的星表里&#xff0c;有的星星写 1.5 等&#xff0c;有的写 12.8 等&#xff0c;这些数字和“亮暗”到底是什么关系&#xff1f;后来搞数据处理&…

作者头像 李华
网站建设 2026/9/9 13:14:03

从收藏到掌握:用技能地图和刻意练习把知识变成能力

去年整理收藏夹和网盘时&#xff0c;我面对过一个尴尬的事实&#xff1a;攒了三百多个教程、买过十几门课&#xff0c;笔记软件里躺着上千条摘抄。但当别人问起“你擅长什么”的时候&#xff0c;我居然答不上来。收藏的东西很多&#xff0c;真正变成 skills 的却很少。这件事促…

作者头像 李华
网站建设 2026/9/9 13:13:43

AVM Triage Report for owner `{{owner_alias}}` - {{YYYY-MM-DD}}

AVM Triage Report for owner {{owner_alias}} - {{YYYY-MM-DD}} 【免费下载链接】awesome-copilot Community-contributed instructions, agents, skills, and configurations to help you make the most of GitHub Copilot. 项目地址: https://gitcode.com/GitHub_Trending…

作者头像 李华
网站建设 2026/9/9 13:13:18

四款降AI率工具横评:比AI检测分更低更关键的是保原意

1. 一个很容易被忽略的问题&#xff1a;AI检测高分不等于你的论文有救 1.1 我为什么突然开始系统性测降AI率工具 2026年这个时间点&#xff0c;论文写作里用AI辅助早就是常态了。我身边的研究生、青年老师&#xff0c;甚至一些高三学生写综述&#xff0c;都是先让大模型出框架…

作者头像 李华