土法炼钢兴趣小组的算法知识备份

【向量检索引擎】Streaming Node 与 Woodpecker WAL:实时可搜的日志层

文章导航

分类入口
databasestorage
标签入口
#milvus#streaming-node#woodpecker#wal#tso#query-delegator#vector-engine

目录

第 3 篇 建立了 Growing / Sealed 状态机;第 4 篇 说明 DML 请求离开 Proxy 后直接打到 Streaming Node,而不是先进 Coordinator。本文补上中间缺的一环:写操作先以 Message 形式进入以 WAL 为中心的 Streaming Service,提交后才能在 Growing 路径上被搜到——这一步具体经过哪些组件、谁保证不丢、以及 2.6 引入的 Woodpecker 到底改变了什么。

本文按官方 Streaming ServiceWoodpecker 文档,钉住 Streaming 子系统的组件、Message/TSO 顺序、WAL 生命周期,以及 Woodpecker 相对 Kafka/Pulsar 的设计目标。不展开 Sealed 上 Knowhere 建索(第 8、10 篇)。

本文是「向量检索引擎」系列第 5 篇(共 19 篇)。→ 系列目录

版本锚定:Milvus v2.6.21 Streaming ServiceWoodpeckerArchitecture OverviewData ProcessingTime Synchronization;Woodpecker 运行时由 github.com/zilliztech/woodpecker 嵌入(见下文源码钉点)。


一、最小故事:insert 返回之后,数据到底在哪

client 调用 insert(vectors) 并拿到成功返回,这时候数据在哪?一个常见的错误直觉是:「返回成功 = 数据已经写进了某个数据库文件」。在 Milvus 2.6 里,更准确的说法是:返回成功 = 这批操作已经作为 Message 提交到了 WAL,而 WAL 提交后保证不丢,但数据本身还停留在内存里的 Growing Segment 中,没有落对象存储

把这句话拆成可验证的步骤(Data Processing · Data insertion):

  1. Proxy 校验请求,按 shard 规则拆包,送到负责该 vchannel 的 Streaming Node。
  2. Streaming Node 给这批操作分配 TSO,做一致性检查,写入底层 WAL。
  3. 一旦持久提交到 WAL,就保证不丢——即便这时候进程崩溃,重启后可以从 WAL replay 出所有 pending 操作。
  4. 提交后的条目被异步切成 Growing Segment;从这一刻起,它可以被增量查询路径搜到。
  5. 之后(可能是几秒到几分钟后)触发 flush,Growing 才变成落对象存储的不可变 Sealed(第 3、6 篇)。

也就是说,「实时可搜」精确指的是第 4 步之后立刻发生的事,而「数据已经安全落盘」在第 3 步就已经成立——这是两件不同时刻发生的事,容易被读者当成同一件事。第六节的常见误解会展开这一点。


二、Streaming Service 三件套:谁负责什么

官方定义 Streaming Service 为围绕 WAL 的内部流式模块,支撑:摄入/订阅、集群状态故障恢复、流式数据转历史数据、Growing 数据查询。架构上三部分:

组件 位置 职责
Streaming Coordinator Coordinator 内的逻辑组件 用 etcd 发现 Streaming Node;绑定 WAL 与节点;暴露 WAL 拓扑供客户端选路
Streaming Node 集群 Worker WAL append、状态恢复、Growing 查询等流式处理
Streaming Client 内部客户端 服务发现、就绪检查;发起写与订阅
flowchart LR
  client["Streaming Client<br/>service discovery, readiness check"]
  coord["Streaming Coordinator<br/>etcd-based WAL binding<br/>topology exposure"]
  cluster["Streaming Node Cluster<br/>WAL append, state recovery<br/>growing query"]
  client -->|"discover WAL location"| coord
  client -->|"write / subscribe"| cluster
  coord -->|"bind WAL to node"| cluster

Access Layer 的 Proxy 内部封装的正是这套 Streaming Client 逻辑:先问 Streaming Coordinator「这个 vchannel 的 WAL 现在绑在哪个 Streaming Node 上」,再把 DML 送到正确的 Streaming Node(第 4 篇 API 路径表)。三个组件各自的职责边界很窄:Coordinator 只做绑定与发现,不摸数据;Node 只处理自己被绑定到的 WAL;Client 只做发现与转发——这与第 4 篇「Worker 是 dumb executor,权威视图只在 Coordinator」的结论一致。


三、Message 与 TSO:日志即写顺序

Streaming Service 是 日志驱动:Milvus 中的写操作(DML 与 DDL)被抽象为 Message

Data Processing 补充:Streaming Node 在写入底层 WAL 前对 payload 做一致性检查;一旦持久提交到 WAL,即保证不丢——崩溃时通过 replay 恢复 pending 操作。

这与经典数据库「先日志后数据」同构,但作用域是 VChannel / shard,不是单机一个 redo 文件。多 WAL、水平扩展的含义见下一节。

3.1 一条 Message 从 Proxy 到 Growing 的时序

把第一节的故事换成时序图,标出每一步谁在做什么:

sequenceDiagram
  participant P as Proxy
  participant SN as Streaming Node
  participant WAL as WAL storage
  participant G as Growing segment
  P->>SN: send InsertMsg (vchannel, payload, TSO)
  SN->>SN: consistency check on payload
  SN->>WAL: append(Message)
  WAL-->>SN: durably committed
  Note over SN,WAL: crash between here and next step<br/>can be recovered by WAL replay
  SN->>G: apply committed entry to growing segment
  G-->>SN: entry now visible to incremental query
  SN-->>P: ack insert success

图中两个时间点值得记住:「WAL 提交」(不丢的分界线)和 「应用到 Growing」(可被搜到的分界线)。两者之间有一个短暂窗口——已提交但尚未应用到 Growing 的数据既不会丢,但也还搜不到,这个窗口通常很短,但在压测或故障恢复场景下会被放大,是排障时需要单独确认的一环(第 17 篇)。


四、WAL 组件:多日志、单所有者、段生命周期

为支撑大规模水平扩展,官方明确:Milvus 的 WAL 不是单一日志文件,而是多个日志的复合;每个日志可独立服务多个 VChannel。约束:

WAL 组件的其它能力(官方列举):

  1. Segment 生命周期管理:按内存、segment 大小、空闲时间等策略管理(与第 3 篇 flush 衔接)。
  2. 基本事务支持:单条 Message 有大小限制时,提供 VChannel 级原子写。
  3. 高并发远程日志写:支持第三方远程消息队列作 WAL 存储;用并发写提升吞吐,靠 TSO 与 TSO 同步维持顺序,按 TSO 顺序读回。
  4. Write-Ahead Buffer:写入 WAL 后暂存在缓冲中,支持尾读而无需每次打远程存储。
  5. 多种 WAL 后端Woodpecker、Pulsar、Kafka;Woodpecker 零本地盘模式可去掉对远程消息队列的依赖。

4.1 Recovery Storage

Recovery Storage 总是运行在持有对应 WAL 的那个 Streaming Node 上。职责:

这是 Growing → 对象存储上 Sealed 的桥梁之一(细节与 flush 策略见第 3、6 篇)。

4.2 Query Delegator

每个 Streaming Node 上的 Query Delegator 负责 单个 shard 上的增量查询

Delegator 与 WAL 组件共存于同一 Streaming Node。若 Collection 配置多副本,则另有 \(N-1\) 个 Delegator 部署在其它 Streaming Node 上(官方 Streaming Service)。

这解释了第 3 篇的查询落点:Proxy → Streaming Node(Delegator)→ 本地 Growing + 远程 Sealed。

4.3 WAL 迁移与 Wait for Ready

存算分离使 WAL 可从一节点迁到另一节点。迁移时:旧节点会拒绝部分请求;新节点恢复 WAL;客户端 Wait for Ready,直到新节点上的 WAL 可服务。排障「写入短暂失败 / 查询抖动」时应考虑是否与 WAL 迁移窗口重合(第 14、17 篇)。


五、Woodpecker:云原生 WAL

Milvus 2.6 引入 Woodpecker 作为专用云原生 WAL,替代直接依赖 Kafka/Pulsar,设计目标(官方 Woodpecker):云环境高吞吐、可靠 append-only 恢复、无本地盘 / 无外部 broker 的运维面。

5.1 零本地盘

组件:Client(读写协议层)、LogStore(高速写缓冲、异步上传、日志 compaction)、Storage backend(S3/GCS/文件系统等)、etcd(元数据与协调)。

5.2 两种部署模式

模式 机制(官方) 适用直觉 官方给出的延迟量级
MemoryBuffer 客户端内嵌缓冲,周期性刷到对象存储;etcd 管元数据 偏批、要简单、可接受较高写延迟 写延迟一般 200–500 ms
QuorumBuffer 与三副本 quorum 交互;至少写成功到 2/3 后视为成功,再异步刷对象存储 低延迟、高耐久 典型 个位数毫秒 完成 quorum 确认

把两种模式画成同一张图,能更直接看出它们在「谁先确认成功」这件事上的差别:

flowchart TB
  subgraph mem ["MemoryBuffer mode"]
    c1["Woodpecker client<br/>embedded buffer"]
    etcd1["etcd<br/>metadata"]
    obj1["Object storage<br/>periodic flush, 200-500 ms"]
    c1 --> etcd1
    c1 -->|"batched flush"| obj1
  end
  subgraph quorum ["QuorumBuffer mode"]
    c2["Woodpecker client"]
    q1["Quorum node 1"]
    q2["Quorum node 2"]
    q3["Quorum node 3"]
    obj2["Object storage<br/>async flush"]
    c2 --> q1
    c2 --> q2
    c2 --> q3
    c2 -.->|"ack after 2 of 3, single-digit ms"| c2
    q1 -.-> obj2
    q2 -.-> obj2
  end

MemoryBuffer 的「成功」定义在 对象存储刷盘完成 之后,所以写延迟直接绑定刷盘周期;QuorumBuffer 的「成功」定义在 quorum 多数副本确认 之后,对象存储刷盘被挪到异步路径,写延迟因此从「百毫秒级」降到「个位数毫秒级」,但多了一套quorum 节点需要运维。以上延迟数字来自 官方文档对模式的描述,不是本机实测。选型时应以自身负载复测为准。

5.2.1 源码落点:Milvus 如何挂上 Woodpecker(v2.6.21

MemoryBuffer / QuorumBuffer 的缓冲区实现落在独立仓库 zilliztech/woodpecker(含 woodpecker/quorum 等包);Milvus 侧通过 WAL 插件适配器接入。pkg/streaming/walimpls/impls/wpinit 时注册 WALNameWoodpeckerBuild 里用 etcd + 对象存储构造嵌入式 client:

// milvus-io/milvus v2.6.21
// pkg/streaming/walimpls/impls/wp/builder.go
func init() {
    registry.RegisterBuilder(&builderImpl{})
    message.RegisterMessageIDUnmsarshaler(message.WALNameWoodpecker, UnmarshalMessageID)
}

func (b *builderImpl) Name() message.WALName {
    return message.WALNameWoodpecker
}

func (b *builderImpl) Build() (walimpls.OpenerImpls, error) {
    // ... ObjectStorage / etcd ...
    wpClient, err := woodpecker.NewEmbedClient(ctx, cfg, etcdCli, storageClient, true)
    // ...
}

写路径在 walImpl.Append:把 Message 交给 Woodpecker LogWriter.Write;若返回 ErrLogWriterLockLost,则标记为 walimpls.ErrFenced(fencing),避免双写:

// pkg/streaming/walimpls/impls/wp/wal.go
func (w *walImpl) Append(ctx context.Context, msg message.MutableMessage) (message.MessageID, error) {
    r := w.p.Write(ctx, &wp.WriteMessage{Payload: pb.Payload, Properties: pb.Properties})
    if r.Err != nil {
        if werr.ErrLogWriterLockLost.Is(r.Err) {
            return nil, errors.Mark(r.Err, walimpls.ErrFenced)
        }
        return nil, r.Err
    }
    return wpID{r.LogMessageId}, nil
}

因此:文档里的两种 buffer 模式是 Woodpecker 库内策略;Milvus 负责把 Streaming Message 接到 wp WAL 实现,并用 etcd/对象存储参数(paramtable.WoodpeckerCfg)注入。读源码时不要在 internal/streamingnode 里找 MemoryBuffer 类型名,应先看 walimpls/impls/wp,再跟进 woodpecker 依赖版本。

5.3 官方吞吐对照(引用,非本机复现)

Woodpecker 文档给出单节点、单客户端、单日志流设置下的对照表(单位与条件以原文为准):

System Kafka Pulsar WP Minio WP Local WP S3
Throughput 129.96 MB/s 107 MB/s 71 MB/s 450 MB/s 750 MB/s
latency 58 ms 35 ms 184 ms 1.8 ms 166 ms

文档同时给出测试机上后端理论吞吐上限(MinIO ~110 MB/s、本地盘 600–750 MB/s、单 EC2 上 S3 最高约 1.1 GB/s),并称 Woodpecker 达到各后端最大可能吞吐的约 60–80%。

引用规则:这些是 厂商文档基准,硬件、对象存储实现、并发模型与你的集群不同则不可直接当容量规划。本系列不在未复现前改写或「约等于」这些数字。

5.4 争论:自研 WAL vs Kafka/Pulsar

立场 主张 代价
沿用 Kafka/Pulsar 运维栈成熟、生态工具多 Milvus 额外依赖消息系统;资源与版本耦合
Woodpecker 零盘 降外部依赖、贴对象存储扩缩 新组件成熟度、延迟模式(MemoryBuffer vs Quorum)要单独理解

Architecture Overview 仍列出 Kafka、Pulsar、Woodpecker 为常见 WAL 实现——不是强制唯一后端。生产选择应写清:一致性需求、可接受的写延迟、是否已有消息队列团队。


六、常见误解

6.1 「insert 返回成功 = 数据已经落对象存储」

不对。返回成功意味着 Message 已经持久提交到 WAL(不会因崩溃丢失),并很快应用到内存里的 Growing Segment(可被搜到)。数据落对象存储要等到 flush 之后(第 3、6 篇)——这中间可能是几秒到几分钟。混淆这两个时刻,会让人错误地以为「数据已经安全归档」,实际它此时还完全依赖 WAL 与 Streaming Node 的存活。

6.2 「Woodpecker 让写延迟总是变成个位数毫秒」

不对。这只在 QuorumBuffer 模式下成立。MemoryBuffer 模式官方写明写延迟一般在 200–500 ms,这与「毫秒级 RAG」的营销话术可能直接冲突。选哪种模式是一次显式的架构决策,不是 Woodpecker 的默认行为。

6.3 「Milvus 的 WAL 就是一份全局日志,类似单机数据库的 redo log」

不对。Milvus 的 WAL 是多个独立日志的复合,每个日志服务一批 VChannel,且任一时刻只能在一个 Streaming Node 上运行。全局顺序不是靠「一份日志」维持,而是靠每条 Message 的 TSO 与 timetick 汇总(第四篇)维持跨 shard 的因果关系;shard 内部顺序才是真正意义上的「日志顺序」。


七、学术谱系与工程间隙

7.1 谱系

阶段 概念 在本篇的落点
ARIES / 数据库 WAL 先日志后数据、崩溃恢复 Message 提交到 WAL 后不丢;replay
分布式日志 Kafka 等共享日志、多副本复制 Milvus 可插拔远程 WAL;多日志水平扩展
Quorum 复制 多数派确认即可对外承诺成功(经典 quorum-based 复制思想) Woodpecker QuorumBuffer 的 2/3 确认策略
云原生 WAL Woodpecker 零本地盘 日志直接(或经 quorum)面向对象存储

Wang et al. (SIGMOD 2021) 强调动态更新与查询并存;2.6 的 Streaming Node + Woodpecker 是该目标在工程上的一轮重构——以官方 2.6 文档为准,不要把 2021 论文组件名硬套过来。QuorumBuffer「2/3 确认即成功、剩余异步补齐」的策略,与经典分布式存储里「多数派写入即可对外承诺持久」的思路同构,区别在于 Woodpecker 把最终归档目标定为对象存储而非本地磁盘副本集。

7.2 工程间隙

7.3 开放问题

  1. 多租户下 WAL 迁移频率与尾延迟的可接受上界如何制度化测量?
  2. Woodpecker 相对 Kafka 的运维成本,在「已有 Kafka 平台」的组织里是否总是净收益?
  3. Delete 经 Delegator 广播到 Query Node 与 Growing 可见性之间的窗口,如何在一致性级别下解释(第 12–13 篇)?
  4. QuorumBuffer 的三副本 quorum 与对象存储之间的一致性窗口,在极端网络分区下是否可能出现「quorum 已确认、对象存储永久落后」的边界情形,官方文档未展开,需要结合具体 release 源码核实。

八、小结

三句话小结

  1. insert 返回成功 意味着 Message 已提交 WAL 并可进入 Growing,不等于已落对象存储。
  2. Streaming 三件套(Coordinator / Node / Client)加上 Delegator,把本地 Growing 与远程 Sealed 拼成一次 shard 查询。
  3. Woodpecker 降低对外部 broker 的硬依赖,但 MemoryBuffer 与 QuorumBuffer 的延迟语义不同,必须显式选型。

下一篇看对象布局:对象存储上的 Segment 布局;或进入执行层:Query Node 与 Segcore

排障时可以按本文的三条误解反向定位问题:查「写完读不到」先看是否卡在 WAL 提交与 Growing 应用之间的窗口;查「写延迟忽高忽低」先确认集群用的是 MemoryBuffer 还是 QuorumBuffer;查「顺序错乱」先确认问题出现在单个 VChannel 内部(理论上不该发生)还是跨 VChannel 之间(本身没有顺序保证,需要靠上层语义处理)。


参考资料

  1. Milvus Documentation v2.6.x, Streaming Service(Coordinator、Node、Client、Message、WAL、Recovery Storage、Query Delegator、Wait for Ready)。
  2. Milvus Documentation v2.6.x, Woodpecker(零盘、架构组件、MemoryBuffer/QuorumBuffer、官方基准表)。
  3. Milvus Documentation v2.6.x, Architecture OverviewData Processing(insertion 流程、WAL 提交语义)。
  4. Milvus Documentation v2.6.x, Time Synchronization(TSO 与 Message 顺序;与第 4 篇共用来源)。
  5. milvus-io/milvus v2.6.21pkg/streaming/walimpls/impls/wp/builder.goWALNameWoodpecker 注册、NewEmbedClient);wal.goAppend / fencing)。
  6. zilliztech/woodpecker(MemoryBuffer/Quorum 等库内实现;经 Go module 引入,版本随 Milvus go.mod)。
  7. Wang et al., Milvus: A Purpose-Built Vector Data Management System, SIGMOD 2021。
  8. 第 3 篇 Segment 状态机第 4 篇 Proxy 与 Coordinator系列 index

返回 系列目录 | 上一篇:Proxy 与 Coordinator | 下一篇:对象存储布局

同主题继续阅读

把当前热点继续串成多页阅读,而不是停在单篇消费。

2026-07-12 · database / storage

【向量检索引擎】一致性模型:四级 GuaranteeTs 与 PACELC 的延迟账

按官方 Consistency Level 与 Timestamp 文档拆解 Strong/Bounded/Session/Eventually 如何映射到 GuaranteeTs,用最小故事、四级时间轴与 Strong 等待时序图说明「一致性」在 Milvus 里首先是一笔延迟账;对照 Abadi PACELC 定理与 Bailis PBS,说明 Bounded 是定性旋钮而非概率保证。


By .