首页
/ Milvus 分布式 QueryView 工作节点侧设计:QueryViewHandler 架构与实现解析

Milvus 分布式 QueryView 工作节点侧设计:QueryViewHandler 架构与实现解析

2026-09-09 23:59:43作者:田桥桑Industrious

本文基于 docs/design-docs/design_docs/qviews/query_view_handler.md 展开,并结合 Milvus 仓库内 internal/views/worknode/handlerinternal/querynodev2/qnviewinternal/streamingnode/server/wal/snview 等真实源码进行纵深解读。读完本文,你将掌握:QueryNode(QN)与 StreamingNode(SN)如何接收 Coord 推送的查询视图(QueryView)、如何通过统一的状态机与回调通道把本地状态上报回 Coord,以及崩溃恢复、持久化清理与 Liveness 契约背后的设计原理。

1. 背景:QueryViewHandler 在分布式 QueryView 体系中的位置

在 Milvus 的分布式查询视图(Distributed Query View)设计中,Coord(Coordinator)作为全局状态机领导者,负责计算并生成每个分片(Shard)的完整查询视图(QueryView),再通过 ReliableSyncer 可靠地推送给各个工作节点(Work Node,即 QueryNode 与 StreamingNode)。整套设计的全局背景可参考 分布式查询视图设计总览

QueryViewHandler 正是这条链路的工作节点侧落地组件

接收 Coord 推送的查询视图,并把本地状态变更上报回去;它是 Coord 侧 Syncer 在节点侧的对偶实现(Counterpart)。

其职责边界非常清晰:

  • Coord(ReliableSyncer):通过 gRPC 双向流把查询视图状态推送给工作节点,并接收节点的状态回报。
  • QueryViewHandler:节点侧单例,管理跨分片的多个状态机(State Machine,简称 SM)实例;它的生命周期长于任何一条 gRPC 流——流断开、重连都不会导致它被销毁。
  • ShardView:分片级内部组件,在分片级互斥锁下持有该分片的所有 SM 实例,负责把 Coord 推送和异步回调路由到正确的 SM。
  • 外部依赖注入:QN 注入 SegmentManager,SN 注入 ResourceManagerCatalog,通过异步回调驱动 SM 前进。

在源码中,这套设计对应三个包:

对应组件
internal/views/worknode/handler ApplyViewQueryViewHandler 接口、ViewSyncServerpendingReports(SN/QN 共享)
internal/querynodev2/qnview QNQueryViewHandlerQNQueryViewStateMachineSegmentManager 接口
internal/streamingnode/server/wal/snview SNQueryViewHandlerSNQueryViewStateMachineStreamingNodeResourceManager 接口、pchannel 维度绑定的 metastore.StreamingNodeCataLog

2. 总体架构与核心组件

2.1 组件清单

2.2 数据流全景

Coord 推送一条查询视图后,数据按如下路径流转:

  1. Coord 发送 SyncRequest,进入 ViewSyncServer 的 recv loop。
  2. ViewSyncServer 把 proto 转换成 ApplyView(每个 ApplyView 携带一个 OnReport 回调,回调内部调用 pendingReports.Update),再调用 handler.ApplyViews。这段转换逻辑对应 sync_server.gohandleViewsRequest
  3. QueryViewHandler 按 ShardID 把视图路由到对应的 ShardView。QN 侧实现见 qnview/handler.go,SN 侧实现见 snview/handler.go
  4. ShardView 创建/驱动 SM 实例,对外部依赖发起 Acquire/Release 调用。
  5. 外部依赖异步回调(OnReadyOnDropped 等)→ ShardView 驱动 SM → SM 产出报告 → 触发 OnReport 回调 → pendingReports.Update
  6. ViewSyncServer 的 send loop 排空 pendingReportsstream.Send(SyncResponse) → 回到 Coord。

2.3 关键不变量(Key Invariants)

  • 报告路径统一:所有报告(无论是 Coord 推送后的即时报告,还是异步回调产生的报告)都走同一条路径:OnReport → pendingReports → send loop → stream.Send。这保证了没有"旁路"状态上报,Coord 永远通过单一通道获知节点状态。
  • 回调替换:当流重连、Coord 重新推送时,ApplyViews 会用新 OnReport 回调替换旧回调。旧回调写入的是已停止(stopped)的 pendingReports——写入被静默忽略,不会 panic。这一点在 pending_reports.goUpdate 中体现:if p.stopped { return }
  • 分片粒度加锁:外层互斥锁保护分片 map,分片级互斥锁串行化分片内 SM 操作;不同分片的视图可以并发应用。两层锁设计在 qnview/handler.go 的注释中明确说明。
  • 可达分片所有权:空分片在触发 onEmpty 回调前会先把自己标记为 detached;handler 只删除"同一个 map 实例",若某批次被已 detach 的分片拒绝,则在当前替换分片上重试。这正是 qnview/handler.gofor { shard := h.getOrCreateShard(shardID); if shard.ApplyViews(shardViews) { break } } 循环的语义。

3. ViewSyncServer:按流的 gRPC 双向流处理

ViewSyncServer 实现 SyncQueryView 双向流 RPC(proto 定义见 pkg/proto/view.proto)。它是按流创建、随流销毁的:流建立时创建,流结束时销毁。SN 与 QN 共享同一实现(sync_server.go)。

3.1 recv loop(主 goroutine)

  1. 从 Coord 接收 SyncRequest
  2. 若是 SyncQueryViewsRequest:把 proto 转换为 ApplyView(每个携带 OnReportpendingReports.Update),调用 handler.ApplyViews
  3. 若是 SyncCloseRequest:通过 pendingReports.SetCloseResponse() 入队关闭信号,然后返回。

对应实现见 sync_server.go

3.2 send loop(后台 goroutine,stream.Send 的唯一调用者)

  1. 等待 pendingReports.Ready()(容量为 1 的通知 channel)。
  2. 排空所有 pending 报告,批量打包成 SyncResponse 发送。
  3. 若设置了关闭标记,发送 SyncCloseResponse 后退出。

stream.Send 收敛到唯一 goroutine,从根本上规避了 gRPC 流并发发送的线程安全问题(见 sync_server.go 的注释与实现)。

3.3 流的生命周期

阶段 行为
建立 创建 pendingReports,启动 send loop,进入 recv loop
重连 Coord 重新推送所有视图;ApplyViews 替换旧回调;SM 重新上报当前状态(fast-forward)
优雅关闭 Coord 发送 close 请求 → send loop 排空剩余报告 → 发送 close 响应 → 退出
流中断 recv loop 返回错误 → pendingReports.Close() → send loop 排空后退出;旧回调变成 stale(对已关闭的 pendingReports 是 no-op)

pendingReports.Close() 会关闭 notify channel,从而唤醒 send loop 做最后一次排空后退出(见 pending_reports.go)。

4. QueryViewHandler:接口、SM 生命周期与幂等性

QueryViewHandler 接口只有一个方法(handler.go):

type QueryViewHandler interface {
    ApplyViews(views []ApplyView)
}

每个 ApplyView 携带一个 Coord 推送的 View 与一个 OnReport 回调。所有状态报告——无论是 Coord 推送处理的即时报告,还是外部依赖回调产生的异步报告——都只通过 OnReport 交付handler.go)。

4.1 SM 生命周期

  • 自动创建:未知 QueryViewKey + Preparing 状态 → 创建新 SM + 发起资源获取。
  • 自动销毁:SM 到达 Dropped → 从分片 map 移除条目 → onEmpty 回调在分片为空时移除该分片。
  • 回调替换:同一 QueryViewKey 被重新应用时替换 OnReport;替换后旧回调永远不会再被调用。
  • 操作幂等性:Coord 对同一 QueryViewKey 的重复推送会复用已有 handler 条目并替换其回调;SM 消费推送的状态,只有当 SM 产生新的资源操作时才调用外部依赖,从而避免重复 Acquire/Release。

在 QN 侧,qnShardView.applyOneLocked 的"先应用 Preparing/Up,再应用 teardown 状态"的两遍处理,保证了新 serving 候选先被安装、旧视图后释放(shard_view.go);SN 侧 snShardView.ApplyViews 采用完全相同的两遍策略(snview/shard_view.go)。

4.2 未知视图处理

当节点上不存在某个被推送的视图时(例如节点重启后),行为取决于被推送的状态:

被推送状态 行为
Preparing 创建新 SM,开始资源获取
Down 立即上报 Dropped:SN 已经丢失了这个 teardown 视图,Coord 可以 fast-forward 清理
Dropped 立即上报 Dropped(让 Coord 完成清理)
其他 上报 Unrecoverable(状态已丢失,Coord 生成替代视图)

该逻辑在 QN 侧 qnview/shard_view.go 与 SN 侧 snview/shard_view.go 中均有对应实现。值得注意的是,SN 对"缺失的 Down/Dropped 视图"的处理略有差异:缺失 Down 也直接上报 Dropped,因为 SN 已丢失本地状态与资源,Coord 需要据此通知所有 QN 完成清理。

5. SN 与 QN 的实现差异

维度 QN SN
外部依赖 SegmentManager ResourceManagerCatalog
SM 状态集 Preparing → Ready → Dropping → Dropped Preparing → Ready → Up → Down → Dropping → Dropped
恢复机制 无(无状态,重启后 Coord 重新推送全部 Preparing) UpRecovering(从持久化的 Up 视图恢复)
持久化 Up → 保存;Down/Dropped → 删除
ApplyViews 顺序 先 Preparing/Up,后 teardown 状态 先 Preparing/Up,后 teardown 状态

对应到源码:QN 侧 QNQueryViewHandler 注释明确"QN is stateless: no persistence, no recovery. On restart, Coord re-pushes all Preparing views.";SN 侧 SNQueryViewHandler 则持有 pchannelcatalogresMgr,并通过 recoverSNQueryViewHandler 从持久化视图重建(snview/handler.go)。

6. QN 侧详解:SegmentManager 交互

正常流程(Preparing → Ready → Dropped):

  1. Coord 推送 Preparing → handler 创建 SM → 调用 segMgr.Acquire(OnReady, OnUnrecoverable)。QN 侧 SM 构造函数只把状态置为 Preparing,不产生 pendingReport——回答要等 SegmentManager 回调(见 qnview/state_machine.go 的注释)。
  2. SegmentManager 异步加载 segment:
    • 成功:调用 OnReady(readySegments)(可多次调用,表示增量进度)→ SM 从 Preparing 前进到 Ready → 上报 Ready 给 Coord。
    • 致命错误:对该 Acquire 调用触发 OnUnrecoverable → SM 从 Preparing 前进到 Unrecoverable → 上报 Unrecoverable 给 Coord。
  3. Coord 推送 Dropped → SM 进入 Dropping → 调用 segMgr.Release(OnDropped)
  4. SegmentManager 异步释放 segment → 调用 OnDropped → SM 从 Dropping 前进到 Dropped → 上报 Dropped 给 Coord → 清理条目。

关键约束:所有回调必须异步(不能在 Acquire/Release 调用期间同步触发),否则会死锁分片互斥锁。这一约束被写进了 SegmentManager 接口的 Liveness Contracts 文档注释中。

另外一个重要边界:重复 QueryViewKey 的处理由 handler/SM 对负责,而不是 SegmentManager——SegmentManager 只负责按调用语义执行 Acquire/Release,不做视图级去重。

6.1 QN 各场景的应答路径(源码注释摘要)

  • 视图不存在:Preparing → 创建 SM + Acquire,无即时应答;Dropped → 立即应答 Dropped(QN 重启场景);其他状态 → 立即应答 Unrecoverable(重启后状态丢失)。
  • 视图已存在:Preparing 且 SM 仍在 Preparing → 无即时应答,等回调;Preparing 且 SM 已越过 Preparing → 立即上报当前状态供 Coord fast-forward;Dropped 且 SM 在 Preparing/Ready/Unrecoverable → 进入 Dropping 并发起 Release;Dropped 且 SM 在 Dropping → 忽略(Release 已在进行);Dropped 且 SM 已 Dropped → 立即重报 Dropped(实际不可达,条目到达 Dropped 即被删除)。

完整逻辑见 qnview/handler.go 的 Response Guarantee 注释。

7. SN 侧详解:ResourceManager + Catalog 交互

正常流程(Preparing → Ready → Up → Down → Dropped):

  1. Coord 推送 Preparing → handler 创建 SM(立即生成 Preparing 报告——与 QN 不同,SN 的 SM 构造函数即产生 pendingReport,见 snview/state_machine.go)→ 调用 resMgr.Acquire(OnReady, OnUnrecoverable)
  2. ResourceManager 异步准备资源:OnReady 使 SM 从 Preparing 前进到 Ready;OnUnrecoverable 使 SM 前进到 Unrecoverable 并上报失败给 Coord。
  3. Coord 推送 Up → SM 从 Ready 前进到 Up → 持久化 Up → 上报 Up 给 Coord。对应 handleCoordUp:同时设置 pendingReport 与 pendingPersist。
  4. Coord 推送 Down → SM 从 Up 前进到 Down → 持久化 Down(= 删除恢复信息) → 上报 Down 给 Coord。对应 handleCoordDown
  5. Coord 推送 Dropped → SM 进入 Dropping → 持久化 Dropped(= 删除恢复信息) → 调用 resMgr.Release(OnDropped)。对应 handleCoordDropped
  6. ResourceManager 异步释放资源 → 调用 OnDropped → SM 从 Dropping 前进到 Dropped → 上报 Dropped 给 Coord → 清理条目。

7.1 Persist-before-report 不变量

持久化永远先于报告执行。如果 SN 在报告之后、持久化之前崩溃,Coord 会认为状态已推进而 SN 实际已丢失该状态,产生不可恢复的认知分裂。因此 consumeReportPersistAndCleanup 严格按照 "先 consumeAndPersist,再 consumeReport" 的顺序执行(见 snview/shard_view.go)。

7.2 可靠写入与取消语义

StreamingNode 的 catalog 把其元数据 KV 用 ReliableWriteMetaKv 包装(见 internal/metastore/kv/streamingnode/kv_catalog.go),使瞬时重试与"结果不确定的写入"处理集中在 metastore 层。handler 把 WAL 生命周期上下文传给 catalog 写入,使关闭(shutdown)时能取消进行中的可靠写入:

  • 被取消的写入不会推进对应的报告或 Release;
  • 其他写失败在这一层是终态的——handler 不在 ReliableWriteMetaKv 之上实现第二层持久化重试机制。

7.3 全视图持久化不变量

SN 持久化的 Up 视图是 Coord 推送的完整 QueryViewOfShard,而不仅是 QueryViewOfStreamingNode。StreamingNode 本地的资源管理器只消费其中的 SN 部分;保留完整拓扑使恢复元数据自包含,供后续消费者使用,而无需在此变更中引入查询执行。在 snview/state_machine.go 的注释中明确:"The SN stores the complete shard view it receives from Coord... QueryNode topology must be retained for query planning after SN crash recovery."

7.4 SN 本地持久化 key 格式

streamingnode-meta/wal/{pchannel}/qv/{collectionID}/{replicaID}/{vchannelIndex}/{streamingVersion}/{compactVersion}/{queryVersion}

值即为 QueryViewOfShard proto。常量定义在 internal/metastore/kv/streamingnode/constant.goMetaPrefix = "streamingnode-meta"DirectoryWAL = "wal"DirectoryQueryView = "qv"),key 拼接逻辑见 kv_catalog.go 的 buildQueryViewKey

设计细节:

  • pchannel 已经存在于父路径中,因此紧凑 key 只存储规范化的 vchannel 索引加上 QueryView/DataView 版本元组(streamingVersion/compactVersion/queryVersion)。
  • 恢复时会校验重建出的 key 身份与持久化 proto 一致(见 kv_catalog.go 附近的 expectedKey 校验)。
  • SaveQueryViews 根据状态决定保存或删除:Up → 保存;其他所有状态(Down/Dropped/Unrecoverable)→ 删除恢复记录(kv_catalog.go)。

8. SN 崩溃恢复:UpRecovering 状态

SN 只持久化 Up 状态。崩溃恢复流程如下:

  1. Catalog 加载持久化的完整 Up 分片视图。
  2. 为每个视图创建 UpRecovering 状态的 SM(Coord 视角仍视为 Up)。
  3. 构造每个恢复出的分片,安装"身份校验过的空回调",并发布到 handler map。
  4. 为每个视图调用注入的 resMgr.Acquire(OnReady, OnUnrecoverable)
  5. OnReady 驱动 UpRecovering → Up → 上报 Up;OnUnrecoverable 驱动 UpRecovering → 本地 Unrecoverable(不向 Coord 报告),并保留持久化恢复元数据,直到 Coord 后续推送 Dropped。

对应实现:

  • SM 恢复构造 recoverSNQueryViewStateMachine:状态为 UpRecovering,无 pendingReport(WAL 追平前不报告),无 pendingPersist(已以 Up 持久化)。
  • 分片恢复与资源获取 startRecovery:在 handler 发布分片并安装空回调之后才发起 Acquire。
  • OnUnrecoverable 在 UpRecovering 下不产生报告(snview/state_machine.go),其设计动机在注释中说明:避免依赖恢复期间可能未设置的 OnReport 回调,同时为可能只是瞬时性的失败保留重试机会。
  • OnRecoveringDone 驱动 UpRecovering → Up(snview/state_machine.go)。
  • coordVisibleState 把 UpRecovering 映射为 Up 供报告使用(snview/state_machine.go)。

注意:资源接口与状态机故障接线属于本变更范围,但具体的资源准备实现(如 WAL 回放、growing segment 重建)在本文档范围内不展开。

9. SN 的 handleCoordDropped 与持久化清理

当 Coord 推送 Dropped 时,SM 进入 Dropping。持久化行为取决于先前状态

先前状态 pendingPersist
Up、UpRecovering Delete(存在已持久化的恢复信息)
Unrecoverable Delete(可能从 UpRecovering 进入,磁盘上有过期恢复信息)
Preparing、Ready、Down None(无已持久化的恢复信息)

对应实现见 handleCoordDropped

  • Up/UpRecovering/Unrecoverable 分支设置 pendingPersist = sm.buildDroppedPersist()(构造 Dropped 状态的持久化 proto,buildDroppedPersist),触发 catalog 删除;同时注释说明"Catalog delete is idempotent, so safe for Preparing→Unrecoverable too"。
  • Preparing/Ready/Down 分支只设置 pendingRelease = true,不设置 pendingPersist。

如果未来的资源接线把 UpRecovering 驱动到 Unrecoverable,状态机会保留持久化的 Up 元数据直到 Coord 推送 Dropped;删除被推迟到那次 Dropping 转换。这保证了在恢复视图失效、但 Coord 尚未清理期间,重启仍有机会重试。

10. SN 的 Handoff Release 所有权

正常 Dropping 与 CloseForHandoff 共享每个视图条目的一条 release 记录:

  • 第一条路径启动 ResourceManager.Release 并拥有一个完成 channel;
  • 后续路径只复用并等待该 channel;
  • Handoff 在其互斥锁下 detach 并清空分片,然后不带锁等待每个已存在的 release 回调完成。

这保证了每个视图恰好一次 Release 调用,同时仍能等待已在途的清理。对应实现见 startReleaseLockedreleaseStarted/releaseDone 守卫)与 CloseForHandoff(在锁外等待所有 releaseDone)。handler 层入口是 SNQueryViewHandler.CloseForHandoff

11. Liveness 契约:应答保证的基石

handler 的"最终必应答"保证依赖于外部依赖履行回调义务:

SegmentManager(QN)

操作 义务
Acquire SM 发起的每次调用,最终必须至少调用该次调用对应的 OnReadyOnUnrecoverable 之一(OnReady 可多次调用表示增量进度,OnUnrecoverable 至多一次表示终态失败)
Release SM 发起的每次调用,最终必须恰好一次调用该次调用的 OnDropped

ResourceManager(SN)

操作 义务
Acquire 最终必须恰好调用 OnReadyOnUnrecoverable 之一
Release 最终必须恰好调用 OnDropped 一次

所有回调必须异步——在 Acquire/Release 期间同步调用会死锁分片互斥锁。违反任何契约都会使对应视图卡在 Preparing/UpRecovering/Dropping,永远无法向 Coord 报告。

这些契约被完整写入了两份接口的文档注释:segment_manager.goresource_manager.go,并且 handler 的 Response Guarantee 注释(qnview/handler.gosnview/handler.go)都以"Provided the ... fulfills its liveness contracts"为前提。

12. 包组织与源码导航

内容
internal/views/worknode/handler ApplyViewQueryViewHandler 接口、ViewSyncServerpendingReports(QN/SN 共享)
internal/querynodev2/qnview QNQueryViewHandlerQNQueryViewStateMachineSegmentManager 接口
internal/streamingnode/server/wal/snview SNQueryViewHandlerSNQueryViewStateMachineStreamingNodeResourceManager 接口、pchannel 绑定的 metastore.StreamingNodeCataLog 使用

配套的测试用例可帮助理解行为语义:

  • 传输层:internal/views/worknode/handler/sync_server_test.go
  • QN 状态机与 handler:internal/querynodev2/qnview/state_machine_test.gointernal/querynodev2/qnview/handler_test.go
  • SN 状态机:internal/streamingnode/server/wal/snview/state_machine_test.go(覆盖 UpRecovering 恢复、UpdateView、报告/持久化/释放语义等)、internal/streamingnode/server/wal/snview/handler_test.go

13. 关联设计文档

QueryViewHandler 是 分布式查询视图设计 中的节点侧组件,与以下文档配套阅读可获得完整图景:

登录后查看全文
热门项目推荐
相关项目推荐

项目优选

收起
kernelkernel
deepin linux kernel
C
33
18
docsdocs
暂无描述
Markdown
899
5.83 K
ops-transformerops-transformer
本项目是CANN提供的transformer类大模型算子库,实现网络在NPU上加速计算。
C++
1.14 K
2.76 K
pytorchpytorch
作为 Ascend for PyTorch 社区的核心组件,TorchNPU 是昇腾专为 PyTorch 打造的深度学习适配插件,使 PyTorch 框架能够直接调用昇腾 NPU,为开发者提供昇腾 AI 处理器的超强算力。
Python
860
1.35 K
ops-nnops-nn
本项目是CANN提供的神经网络类计算算子库,实现网络在NPU上加速计算。
C++
925
1.85 K
jiuwenswarmjiuwenswarm
JiuwenSwarm 是一款基于openJiuwen开发的智能AI Agent,它能够将大语言模型的强大能力,通过你日常使用的各类通讯应用,直接延伸至你的指尖。
Python
3.84 K
1.02 K
kernelkernel
openEuler内核是openEuler操作系统的核心,既是系统性能与稳定性的基石,也是连接处理器、设备与服务的桥梁。
C
533
601
ops-mathops-math
本项目是CANN提供的数学类基础计算算子库,实现网络在NPU上加速计算。
C++
1.37 K
1.46 K
AscendNPU-IRAscendNPU-IR
AscendNPU-IR是基于MLIR(Multi-Level Intermediate Representation)构建的,面向昇腾亲和算子编译时使用的中间表示,提供昇腾完备表达能力,通过编译优化提升昇腾AI处理器计算效率,支持通过生态框架使能昇腾AI处理器与深度调优
C++
548
395
cann-learning-hubcann-learning-hub
CANN 学习中心仓,支持在线互动运行、边学边练,提供教程、示例与优化方案,一站式助力昇腾开发者快速上手。
Jupyter Notebook
1.04 K
525