Milvus 多集群 WAL 复制与 CDC 机制详解:ReplicateConfig、Replicate Interceptor 与拓扑切换全解析
本文基于当前 Milvus 开源仓库中 docs/agent_guides/streaming-system/replication/replicate.md 的核心内容,结合 replicateutil、WAL Replicate Interceptor 与 CDC ChannelReplicator 的源码实现展开。Milvus 通过"星形拓扑 + 按 PChannel 复制"的架构,将一个 PRIMARY 集群的 WAL 消息实时转发到多个 SECONDARY 集群,实现跨集群的 WAL 复制(Replication)与变更数据捕获(CDC)。读完本文,你将掌握 ReplicateConfig 的配置模型与校验规则、PRIMARY/SECONDARY 角色如何由 Replicate Interceptor 强制实施、消息逐条复制与 Checkpoint 推进的一致性语义,以及 AddNewMember / SwitchOver / FailOver / RemoveMember 四类拓扑变更的底层触发路径与恢复流程。
整体架构:星形拓扑的跨集群 WAL 复制
Milvus 的多集群复制建立在 WAL(Write-Ahead Log,预写日志) 之上,核心思路是:把一次写入产生在 PRIMARY 上的 WAL 消息,原样转发并重放到 SECONDARY 集群的本地 WAL,从而使副本集群在消息层面保持一致。这里有几个关键定语:
- 跨集群(cross-cluster):复制发生在多个独立 Milvus 集群之间,而不是单个集群内部的多副本之间;
- 基于 WAL(WAL-based):被复制的最小单元是 WAL 中的消息(message),消息携带原始的
MessageID、Properties与Payload; - 按 PChannel(per-PChannel):复制任务以 PChannel 为粒度并行推进,每个 PChannel 独立维护检查点;
- 星形拓扑(star topology):所有写流量只进 PRIMARY 中心集群,其余 SECONDARY 作为叶子节点单向接收。
关联阅读:流式系统总览、WAL 恢复存储 WALCheckpoint。
拓扑在配置上是有向边集合:SourceClusterID → TargetClusterID 表示"源集群 → 目标集群"的复制方向。对任意合法拓扑,存在唯一一个 PRIMARY 中心节点(出度 = N-1,入度 = 0,负责接收全部客户端写),其余 N-1 个节点全部是 SECONDARY 叶子节点(入度 = 1,出度 = 0),不允许出现链条式、环状或多中心的结构。
ReplicateConfig:复制拓扑的配置模型
配置结构
复制拓扑由 ReplicateConfiguration(protobuf)描述,它保存了两个部分:
| 组成部分 | 字段 | 说明 |
|---|---|---|
| Clusters 列表 | ClusterID |
参与复制的每个集群的唯一标识,集群内不允许重复、不允许含空白字符 |
PChannels |
有序的 PChannel 列表,集群内不重复、不允许为空 | |
ConnectionParam |
连接参数,包含 uri 与 token,用于建立到目标集群的复制流 |
|
| CrossClusterTopology 边列表 | SourceClusterID → TargetClusterID |
声明复制方向的有向边,边不可重复,端点必须存在于 Clusters 列表中 |
从源码看,这份配置的持久化位置是 WALCheckpoint(详见 recovery-storage.md):WALCheckpoint 中除了常规消费进度(MessageID = LastConfirmedMessageID、TimeTick),还专门保存了面向 SECONDARY 集群的 ReplicateCheckpoint。拓扑的任何变更都不是本地直改,而是通过一条 AlterReplicateConfig 广播消息 在全集群范围内原子更新(参见 集群级消息语义)。
跨集群 PChannel 按索引位置一一对应
所有集群必须配置相同数量的 PChannel,跨集群的通道映射不依赖通道名,而是按索引位置对齐:
Source.PChannels[i] → Target.PChannels[i]
pkg/util/replicateutil/config_helper.go 中的 NewConfigHelper 负责把 protobuf 配置构造成一张内存图:
vs[cluster.GetClusterId()] = &MilvusCluster{ ... }
// 为每个集群建立 pchannel -> 下标 的索引,供后续按位置映射
for i, pchannel := range cluster.Pchannels {
vs[cluster.GetClusterId()].idxMap[pchannel] = i
}
建图后立刻做三类结构校验,校验失败会直接报 ErrWrongConfiguration:
- 拓扑边端点必须真实存在,且不能形成环(一个集群只能有一个 source),否则报错;
- PRIMARY 数量必须恰好为 1(
primaryCount != 1即报错); - 所有集群的 PChannel 数量必须相等,否则报
"pchannel count is not equal for cluster ..."。
在这张图中,每个 MilvusCluster 节点携带 role(RolePrimary/RoleSecondary)、source 与 targets 三个派生属性:
Role()返回当前集群角色;SourceCluster()返回 PRIMARY 节点返回nil,SECONDARY 节点返回其复制源集群;TargetClusters()返回该集群作为源所指向的全部下游集群;MustGetSourceChannel(pchannel)借助idxMap通过"目标集群通道 → 同一下标 → 源集群通道"完成通道名反向换算;GetTargetChannel(currentPChannel, targetClusterID)完成正向换算。
这也解释了为什么源、目标集群的 PChannel 必须数量一致且顺序稳定——位置即语义,任何一侧通道列表的乱序都会导致复制到错误的通道上。
ConfigValidator:变更时的业务规则校验
pkg/util/replicateutil/config_validator.go 中的 ReplicateConfigValidator.Validate() 在提交新配置前逐条执行检查,覆盖几个层次:
- 集群基础格式(
validateClusterBasic):ClusterID非空且不含空白;ConnectionParam.uri非空且必须能通过url.ParseRequestURI,并且 URI 全局唯一(不允许两个集群指向同一地址);每个集群pchannels非空、内部无重复;所有集群 PChannel 数量一致。 - 相关性(
validateRelevance):当前集群必须出现在 clusters 列表中,且配置中的 PChannel 集合必须与当前集群实际的 PChannel 一致。 - 边唯一性(
validateTopologyEdgeUniqueness):不允许出现重复的source->target边,边的两个端点都必须存在。 - 星形约束(
validateTopologyTypeConstraint):通过统计每个节点的入度/出度验证——必须存在恰好一个中心节点(出度 = 集群数-1、入度 = 0),其余节点一律入度 = 1、出度 = 0,否则报"no center node found, star topology must have exactly one center node"或类似错误。 - 变更对比(
validateConfigComparison/validateClusterConsistency):与当前配置比较时,已有集群的属性不可变——已有 PChannel 必须保持原顺序原位置(只允许在尾部追加扩容,禁止删减或调序),connection_param.uri与connection_param.token一经配置不可修改。若检测到 PChannel 数量增长,会置位IsPChannelIncreasing(),供上层感知。
这些校验与 NewConfigHelper 在建图时的校验共同构成"准入"防线:前者在提交/广播前把关,后者在加载/应用时兜底。
角色与准入控制:Replicate Interceptor 如何强制 PRIMARY/SECONDARY
复制功能通过 WAL 写入链路上的 Replicate Interceptor 施加角色约束。相关实现位于 internal/streamingnode/server/wal/interceptors/replicate/replicate_interceptor.go,核心是 DoAppend。
| 角色 | 允许的写入 | 被拒绝的写入 |
|---|---|---|
| PRIMARY | 接受客户端写入(DML/DDL/DCL) | 任何携带 replicate header 的复制消息(副本消息不允许在 PRIMARY 上二次写入) |
| SECONDARY | 只接受从 PRIMARY 转发过来的、携带 replicate header 的消息 | 任何不携带 replicate header 的消息 |
需要特别说明的是 SECONDARY 侧的例外:WAL 自控消息(self-controlled messages),例如 TimeTick、CreateSegment、Flush,它们是各角色本地生成、与复制来源无关的控制消息,整条绕过 Replicate Interceptor,不受上述角色约束。这一点非常重要——它保证了 SECONDARY 即使不接收任何复制数据,其本地的 WAL 心跳、段分配与刷盘控制流仍能正常推进。
replicates.ReplicatesManager(manager.go)抽象了单个 WAL 上的两种运行模式:
- primary:WAL 只会收到非复制消息;
- secondary:WAL 只会收到复制消息。
角色切换通过 SwitchReplicateMode 完成(由 AlterReplicateConfig 消息触发),其注释精确列举了五种情形:
primary → secondary:进入 replicating 模式,此后不携带 replicate header 的消息将被拒绝;primary → primary:无操作;secondary → primary:退出 replicating 模式,丢弃旧的远端复制检查点(secondary state);secondary → secondary且源集群发生变更:丢弃指向旧源集群的检查点;secondary → secondary且源集群未变:无操作。
从角色语义可以推出:复制配置本质上决定了"谁拥有写权"——拓扑变更(SwitchOver / FailOver)的核心就是通过 AlterReplicateConfig 消息原子性地搬动这把写锁。
消息级复制跳过:Unreplicable(_ur)消息属性
并非所有 WAL 消息都适合立刻被 SECONDARY 重放。某些 DDL/控制消息在"跨集群重放契约"尚未确定化之前,直接重放可能导致两个集群的状态分叉。因此 Milvus 采用了一种软跳过机制:
- 生产者侧:产生这些消息的上游组件给具体 WAL 消息打上
Unreplicable消息属性(内部 key 为_ur)。相关定义见 pkg/streaming/util/message/properties.go,构建入口见 builder.go 的WithUnreplicable(),读取侧见 message_impl.go 的IsUnreplicable()。 - CDC 发送侧:
ChannelReplicator对待_ur消息如同被忽略的消息——不发送但推进复制进度。startConsumeLoop中对返回ErrReplicateIgnored的消息直接continue,见 channel_replicator.go。 - SECONDARY 接收侧:Replicate interceptor 同样忽略携带
_ur的复制消息,保护混合版本滚动升级期或已转发的流量不被重复处理。
关键的语义边界是:这是"消息属性"而非静态的 MessageType 规则。也就是说,一旦未来某个 DDL 的重放契约被确定化,只要停止在新产生的消息上设置 _ur,即可自动开启该消息的复制;而旧 WAL 中已经带 _ur 的历史消息在滚动升级期间仍保持被跳过,从而保证向后兼容。
数据流:一条写入如何变成 SECONDARY 上的一条 WAL 记录
复制是按 PChannel 独立并行进行的。每条 PChannel 在 PRIMARY 侧对应一个由 streamingcoord 的 ChannelManager 编排的 CDC ChannelReplicator 任务(相关代码在 internal/streamingcoord/server/balancer/ 与 internal/cdc/replication/ 下)。完整链路分三步:
① PRIMARY WAL → CDC ChannelReplicator
ChannelReplicator 的 init()(channel_replicator.go)会做三件初始化工作:
- 建立到目标(target)集群的 Milvus client;
- 从 SECONDARY 的
ReplicateCheckpoint开始读 PRIMARY WAL。具体地,它使用DeliverPolicyStartFrom(cp.MessageID)从检查点记录的MessageID之后投递,并用DeliverFilterTimeTickGT(cp.TimeTick)过滤掉不大于检查点TimeTick的历史消息——这保证了复制流在 PRIMARY 侧只发送 SECONDARY 尚未消费的部分; - 建立 ReplicateStreamClient。
在消费循环中,TimeTick/CreateSegment/Flush 等自控消息以及携带 _ur 属性的消息被跳过;其余消息交给 ReplicateStreamClient 发送。
② ChannelReplicator → SECONDARY Proxy(gRPC 双向流)
通过 CreateReplicateStream 这个 gRPC 双向流式 RPC(运行在 SECONDARY 的 Proxy 上),PRIMARY 把每条消息连同其原始的 MessageID、Properties、Payload 以及 SourceClusterID 一并发送过去。复制流保持长连接,用于持续推送并同时承载确认信息。
③ SECONDARY Proxy → SECONDARY WAL
SECONDARY 的 Proxy 收到后,将 VChannel 名称重映射到本地通道(这正是上面"PChannel 按索引位置映射"的实际落地环节),然后追加写入本地 WAL。写入时 Replicate Interceptor 做准入校验——验证来源集群 ID 是否匹配、按 TimeTick 去重,并推进本地检查点。
另外值得注意的实现细节:ChannelReplicator 会维护复制延迟(lag)指标序列。为避免 CDC Pod 重启或目标不可达期间监控面板上 lag 读数"消失"而误显示为 0,init() 会先用**初始化检查点(创建时刻的保守下界)**预置指标,等目标可达后再用确认过的检查点覆盖(channel_replicator.go)。这个"预置保守下界、只可能高估延迟"的设计,对复制可观测性很有参考价值。
Checkpoint 与一致性语义
SECONDARY 上每个 PChannel 独立维护一个 ReplicateCheckpoint,其结构为:
ReplicateCheckpoint { ClusterID, PChannel, MessageID, TimeTick }
secondaryState(secondary_state.go)把这个检查点与事务复制辅助结构(replicateTxnHelper)绑定在一起。检查点推进规则直接决定了一致性强度:
- 非事务消息:成功 append 后立即推进检查点;
- 事务消息:检查点只在
CommitTxn时推进,BeginTxn与事务体(body)消息均不推进。这样当 SECONDARY 崩溃恢复时,未提交事务的所有消息仍会从头重放,不会因"提交了一半"造成数据不完整; - 去重规则:
TimeTick ≤ checkpoint.TimeTick的消息一律忽略(见 secondary_state.go 中PushForwardCheckpoint对旧 TimeTick 的短路逻辑)。对于当前进行中事务的事务体消息,因为同一事务内所有消息共享同一个 TimeTick,会出现TimeTick == checkpoint.TimeTick的"相等"情形,此时不能简单按 TimeTick 丢弃,而是交由事务辅助器按 MessageID 去重,从而保证在断点重放时不会重复写入同一事务的消息。
检查点持久化在 WALCheckpoint 中(写入 catalog/etcd,参见 recovery-storage.md)。PRIMARY 侧可通过 GetReplicateInfo 主动查询 SECONDARY 的复制进度,以便在 CDC 任务重启后从正确位置续传。
恢复流程:SECONDARY 如何无缝续传
WAL 打开时,RecoverReplicateManager 从 RecoveryStorage 快照中加载两类状态:
ReplicateConfig:当前集群的复制角色与完整拓扑;ReplicateCheckpoint:每个 PChannel 已确认复制到的位置。
对于 SECONDARY 集群,还会从快照的 TxnBuffer 中恢复未提交的复制事务:recoverSecondaryState(secondary_state.go)遍历 TxnBuffer.GetUncommittedMessageBuilder() 取出的未提交事务,过滤出 ReplicateHeader.ClusterID == sourceClusterID(即确实从当前源集群复制而来)的那一个进行中事务,重建 BeginTxn → body 消息 的辅助状态。之后当该事务后续的 body/commit 消息到达时,SECONDARY 可以正确接续而非误判为非法写入。
正是因为"事务检查点只在 CommitTxn 推进 + TxnBuffer 恢复未提交事务"这两个机制叠加,SECONDARY 才能做到断点续传且不丢不重。
拓扑变更:AlterReplicateConfig 广播与四类操作
所有拓扑变更都通过 AlterReplicateConfig 广播消息触发。正如 集群级消息语义 所述,该消息与 FlushAll、AlterWAL 一样属于全局屏障(global barrier):它要求 ExclusiveCluster 资源锁,保证变更期间任何 PChannel 上都没有其他消息在途(相关锁语义见 广播与协调)。这为"在多个 PChannel 之间一致地切换角色"提供了强同步保证。
Replicate Interceptor 在处理 AlterReplicateConfig 时会按以下顺序执行(replicate_interceptor.go):
- 若消息头带
Ignore标记(用于"强制晋升后忽略不完整的切换消息"),则跳过角色切换与事务回滚,仅原样 append; - 否则先调用
SwitchReplicateMode在 append 之前切换复制模式——注释明确说明该消息受 WAL 级锁保护,因此此时切换是安全的; - append 消息使配置变更持久化;
- 若消息头带
ForcePromote(强制晋升),append 成功后回滚所有在途事务(RollbackAllInFlightTransactions),保证"只有配置变更真正落 WAL 后才会回滚事务"的原子性。强制晋升过程中捕获的 salvage checkpoints 可通过GetSalvageCheckpoint()查询。
文档定义的四类拓扑变更及其语义如下:
| 操作 | 语义 | 底层行为 |
|---|---|---|
| AddNewMember | 新增一个集群及对应的拓扑边 | 复制从新到达的 AlterReplicateConfig 消息所处 WAL 位置开始;已有集群的属性不可变(由 ConfigValidator 保证) |
| AddNewPChannel | 通过配置变更新增 PChannel | 不支持——所有集群的 PChannel 数量必须在初始配置时一致,后续只能整体规划(validator 只放行尾部追加,但复制任务按初始等量通道建立) |
| SwitchOver | 更新拓扑边以反转角色(例如 PRIMARY A → SECONDARY B 变为 PRIMARY B → SECONDARY A) | 在旧 PRIMARY 上,SwitchReplicateMode 丢弃 secondary state;在新 PRIMARY 上创建指向新源集群的新 secondary state |
| FailOver | 从拓扑边中移除故障的 PRIMARY,将某个 SECONDARY 提升为新的 PRIMARY | 旧 PRIMARY 上的 CDC ChannelReplicator 一旦检测到自身拓扑边被移除就停止;提升侧通过 ForcePromote 语义回滚在途事务并接管写权 |
| RemoveMember | 移除指向目标集群的拓扑边 | CDC ChannelReplicator 通过 AlterReplicateConfig 消息感知边被移除,从 etcd 清理该复制 PChannel 的元数据后退出(在 startConsumeLoop 中即"检测到 replication removed → stop consume loop"的分支) |
在 ChannelReplicator 的消费循环中,每个消息复制成功后都会检查它是否是 AlterReplicateConfig 类型,一旦发现该消息使当前复制关系被移除(util.IsReplicationRemovedByAlterReplicateConfigMessage),就会 BlockUntilFinish() 等待在途确认全部完成、删除 lag 指标序列后优雅退出(channel_replicator.go)——这保证了 CDC 任务不会在拓扑拆除后继续向已脱离的集群发送数据。
核心代码包地图
复制的实现横跨工具层、协调层、WAL 层与 CDC 层,以下是按文档归纳的关键包及其职责,方便深入阅读:
- pkg/util/replicateutil/ —
ConfigHelper(拓扑图构建与角色/通道查询)、ConfigValidator(变更校验)、角色定义RolePrimary/RoleSecondary; internal/streamingcoord/server/balancer/—ChannelManager对复制配置的持久化编排、AvailableInReplication判定、CDC 任务创建;- internal/streamingnode/server/wal/interceptors/replicate/ — Replicate Interceptor(
replicate_interceptor.go)、ReplicatesManager、secondaryState与事务辅助器(replicates/子目录); - internal/cdc/replication/ — CDC
ChannelReplicator(replicatemanager/)与ReplicateStreamClient(replicatestream/,含复制延迟指标与消息队列实现)。
结语与进一步阅读
Milvus 的多集群 WAL 复制是一套围绕"消息、检查点、角色、屏障"四个要素设计的机制:消息是复制的最小单元,检查点是断点续传的锚点,PRIMARY/SECONDARY 角色由 Replicate Interceptor 在写入路径上强制实施,而所有拓扑变更都必须以 AlterReplicateConfig 全局屏障的方式原子落地。理解这条链路后,可以沿以下仓库内文档继续深入:
- 流式系统架构总览
- WAL 恢复存储与 WALCheckpoint
- 集群级消息与全局屏障语义
- 事务消息语义(BeginTxn/CommitTxn 与去重)
- 广播机制与 ExclusiveCluster 资源锁
需要说明的是,跨集群复制在 Milvus 中作为独立部署能力存在,具体启用方式与运维边界请结合官方发布渠道的部署文档与版本说明确认,以上机制描述均以当前仓库源码与内部文档为准。
atomcodeClaude Code 的开源替代方案。连接任意大模型,编辑代码,运行命令,自动验证 — 全自动执行。用 Rust 构建,极致性能。 | An open-source alternative to Claude Code. Connect any LLM, edit code, run commands, and verify changes — autonomously. Built in Rust for speed. Get StartedRust0629
MiniCPM5-2BMiniCPM5-2B 是一款面向端侧、本地部署和资源受限场景的 2B 稠密 Transformer,能够达到同尺寸开源模型 SOTA 水平。Markdown00
GLM-5.3GLM-5.3 与 GLM-5.2 使用相同的基座模型——所有提升均来自后训练。与 GLM-5.2 相比,它在复杂编程和长程任务上的表现显著提升。Jinja00
HivisionIDPhotos⚡️HivisionIDPhotos: a lightweight and efficient AI ID photos tools. 一个轻量级的AI证件照制作算法。Python07
DragonOSDragonOS is an operating system developed from scratch using Rust, with Linux compatibility. It is designed for **Serverless** scenarios. 使用Rust从0自研内核,具有Linux兼容性的操作系统,面向云计算Serverless场景而设计。Rust00
Spark-X2.5-1.7BSpark-X2.5-1.7B 旨在让强大的 AI 更加实用、高效且易于获取。这些模型在广泛的日常任务中表现出色,涵盖对话、写作、翻译、推理、编程、工具调用和智能体工作流,并在同等规模的开源模型中取得领先结果。Spark-X2.5 将面向效率的架构与最高 1M tokens 的原生上下文窗口相结合,并支持 200 多种语言。Python00