Milvus StreamingCoord Broadcaster 源码级解析:跨 PChannel 的原子 DDL/DCL 广播与幂等保证
【免费下载链接】milvusMilvus is a high-performance, cloud-native vector database built for scalable vector ANN search项目地址: https://gitcode.com/GitHub_Trending/mi/milvus
本文围绕 Milvus 流式协调服务(StreamingCoord)中运行在单例里的Broadcaster组件展开,它是所有 DDL/DCL 消息(建集合、加索引、鉴权变更、导入任务等)以“跨 PChannel 原子广播”方式下发到 WAL 的统一执行引擎。你将看到其广播 API 的调用约定、六阶段执行流程、基于 ResourceKey 的加锁模型、带 Ack 追踪与回调的任务状态机,以及一套以对象ID为身份、以_ik幂等键为线索的去重设计——读完可以掌握如何正确、安全地接入该广播体系(例如实现类似 BulkImport 的客户端重试语义)。
Broadcaster 在 Milvus 中的定位
Milvus 的流式架构中,一份 DDL/DCL 通常需要投递到多个 PChannel 下的所有目标 VChannel(例如一个集合的消息可能分布在多个虚拟通道上),并要求这些写入是“原子提交、整体可见、可恢复”的。这个职责由StreamingCoord 内部唯一运行的 Broadcaster 单例承担。
从源码接口看,Broadcaster 的核心能力定义在 broadcaster.go:
WithResourceKeys(ctx, resourceKeys...):按资源键加锁并返回BroadcastAPI;若当前集群不是主集群(non-primary),直接返回ErrNotPrimary;WithSecondaryClusterResourceKey(ctx):用于 force promote,仅在次集群(secondary)可用,获取一个排他性集群级资源键;Ack(ctx, msg)/LegacyAck(...):接收各 VChannel 侧的消费确认(旧版 2.6.0 导入消息走 Legacy 接口);GetPendingSchemaFileResources():恢复期重建文件资源引用计数用;Close():优雅关闭。
其包级门面broadcast.StartBroadcastWithResourceKeys(...)见 singleton.go,会先等待 WAL-based DDL 就绪(WaitUntilWALbasedDDLReady)再进入加锁路径。BroadcastAppendResult 与 AppendOperator(真正写 WAL 的streaming.WAL()实现)共同构成了向多通道追加消息的抽象。
值得注意的错误语义:ErrNotPrimary(“cluster is not primary, cannot do any DDL/DCL”)与ErrNotSecondary在 broadcaster.go 中定义,非主集群的任何广播请求都会被主从角色检查拦截——这是跨集群复制(replication)场景下防止次集群产生新 DDL 的硬约束。
Broadcast API:调用方的完整契约
一个调用方使用广播的标准姿势是:
// 1. 拿锁 + 获取广播句柄(内部已分配 broadcastID,并等待 WAL based DDL 就绪) api, err := broadcast.StartBroadcastWithResourceKeys(ctx, resourceKeys...) if err != nil { /* 非主集群会得到 ErrNotPrimary */ } defer api.Close() // 若未发起 Broadcast,Close 负责释放已持有的资源键锁 // 2. 构造 BroadcastMutableMessage:目标 VChannels 必须包含 CChannel(控制通道) msg := message.NewBroadcastMutableMessage(...) // 3. 执行广播 result, err := api.Broadcast(ctx, msg)关键约定(来自 broadcaster_with_rk.go):
Broadcast()一旦执行,句柄上的锁守卫会被立即“消费”(所有权移交):要么由新注册的任务持有并在其 ack 回调结束时释放,要么在去重命中、关闭等路径上由broadcast()内部立即释放。此后Close()必须是无操作的,包括 panic 路径在内;- 调用方必须在加锁获取句柄之后才可调用
Broadcast();若只拿锁而不广播(例如前置校验失败),必须调用Close()释放锁,避免 DDL 死锁; Broadcast()会为消息盖上广播头(OverwriteBroadcastHeader(id, resourceKeys...)),并注入 trace context,保证原始调用方的追踪上下文在 ack 回调执行时(原始调用早已结束)仍然可还原。
Broadcast Flow:六阶段执行流程
Broadcaster 将一个 DDL/DCL 从提交到回收划分为六个阶段,各阶段都能在源码里找到对应的实现:
- 加锁(Lock):资源键按
(Domain, Key)排序后依序加锁(防止死锁),并自动追加一个 SharedCluster 共享集群键。实现见 resource_key_locker.go 与 broadcast_manager.go。 - 持久化(Persist):在加锁临界区内先分配
broadcastID(全局 ID 分配器),创建处于PENDING状态的任务并把消息落盘到 catalog(etcd)。一旦持久化成功,即使进程崩溃,该广播也保证最终完成——恢复时会从 catalog 重建。相关逻辑集中在 broadcast_task.go(构造 PENDING 任务、标记 dirty)和saveTaskIfDirty(写SaveBroadcastTask)。 - 追加(Append):任务交给
broadcastScheduler。调度器按硬件 CPU 数 × streaming.WALBroadcasterConcurrencyRatio启动一批 worker(见 broadcast_scheduler.go),worker 调用streaming.WAL().AppendMessages()把消息按 VChannel 拆条写入目标 PChannel(pending_broadcast_task.go)。部分追加失败的消息会被保留并进入指数退避重试队列,直到全部成功。 - 快速确认(FastAck):追加全部成功后触发
FastAck(broadcast_task.go)。若广播头的AckSyncUp未置位,Broadcaster 直接依据 append 结果“自确认”所有 VChannel,不必等待消费端确认,从而显著降低 DDL 延迟;若AckSyncUp置位,则必须等待各 StreamingNode 消费者真正 ACK 每个 VChannel。 - Ack 回调(AckCallback):CChannel(控制通道)被 ACK 后,任务进入
ackCallbackScheduler。只有全部 VChannel 都被确认后回调才会执行;存在资源键冲突的任务,回调严格按CChannel TimeTick 顺序执行(恢复场景下按ControlChannelTimeTick排序,见 ack_callback_scheduler.go),保证与 WAL 顺序一致。回调失败采用指数退避重试直至成功(初始 10ms、上限 10s,见 ack_callback_scheduler.go)。 - 墓碑与回收(Tombstone & GC):回调全部完成后任务转为
TOMBSTONE状态,tombstoneScheduler按maxLifetime/maxCount策略把过期任务从 catalog 清除(详见后文配置小节)。
幂等广播(Idempotent Broadcast)
Broadcaster 的幂等设计是文档中着墨最深的部分,也是整个组件的精髓。它解决的是“客户端因超时/网络抖动重试同一个 DDL,是否会执行两次”的问题。
幂等键的编码与作用域
携带_ik幂等键属性(属性常量见 properties.go)的广播消息,会被额外按该键建立索引。幂等键不是一个裸字符串,而是一个自描述的三段编码:
<domain>:<scopeID>:<clientKey>其中 domain 取值为集群(1)/数据库(2)/集合(3),scopeID 是对象的数字 ID。编码构造在 idempotency.go 的三个构造函数中完成:
NewClusterScopedIdempotencyKey(clientKey)—— 全集群唯一;NewDatabaseScopedIdempotencyKey(dbID, clientKey)—— 同一数据库内唯一;NewCollectionScopedIdempotencyKey(collectionID, clientKey)—— 同一集合内唯一。
编码是前缀可解析、无需定界技巧的:domain 与 scopeID 都是十进制数字,因此前两个冒号天然完成分段,客户端键作为无界尾部任意包含冒号都不会与别的 scope 冲突(有单测专门验证这一点,见 idempotency_test.go)。
不存在“未定界”的裸键:WithIdempotencyKey只接受IdempotencyKey类型,而该类型只能由上面三个带作用域的构造函数产出,空 clientKey 会得到空值(“非幂等写入”)。因此调用方不可能因为“忘记”选作用域而造出一个会跨集群静默去重的键。
去重作用域:messageType + 调用方选择的对象身份
Broadcaster 内部把去重身份压缩成一个 map 键(idempotency_index.go):
scope = strconv(MessageType) + "/" + IdempotencyKeymessageType由Broadcaster追加:否则同一集合上的 CreateIndex 与 DropIndex 会共享同一作用域,后者会被静默吞掉;- 对象作用域来自调用方:只有调用方知道自己的操作作用于哪个对象。
索引与任务注册同处 manager 锁下的一个临界区(getOrAddBroadcastTask,见 broadcast_manager.go),所以两个并发的同键请求无论各自持有何种资源键,都不可能同时 miss。
命中后的“等待原广播完成”
后到的同键请求命中去重后:不创建任何任务,而是等待原广播的 ack 回调执行完毕,再把原始广播的结果返回,并把原消息放进返回结果的Duplicated字段。等待通过BlockUntilDone实现(broadcast_task.go),等待有界于请求 context——调用方 context 过期时拿到的是自己的超时错误,而不是一个“无主”的 broadcastID。
这个等待是必要的,原因有两个:
- 原任务并非总能在重试请求前面被资源锁串行化——例如重试与改名(rename)竞态时,重试持有的是旧名字的锁;又如从复制 WAL 恢复出来的任务根本不持锁;
- 若不等原广播完成就返回,调用方会拿到一个 broadcastID,但该 ID 的“效应”(导入场景即 ack 回调里创建的 job)还不存在。
等待结束后,AppendResults会由原广播持久化的每个 VChannel 的 checkpoint重建,且绝不会为 nil——因此只看 append 结果的调用方无法区分重复请求与新广播。
“名字 vs 身份”:重命名与删建同名
设计上以 ID(而非名字)作为作用域,正是为了同时解决经典的名字-身份问题:
- 改名(Rename):
RenameCollection保持 collectionID 不变而只改名字。因为作用域是 ID,键始终绑定在该集合上,用新名字重试仍然能命中原广播。但要注意,这说的是作用域语义而非名字解析:像 import 的importTask.PreExecute这类先“名字→ID”解析的入口(经由 proxy 元数据缓存),会在到达 Broadcaster 之前就拒绝仍带着旧名字的重试——该请求直接失败,而不是重复导入。 - 删除后同名重建(Drop & recreate):重建的集合获得全新 ID,作用域随之变化,查找必然 miss,于是创建全新任务——这正是正确结果:两个请求针对的是不同的集合。此时导入侧仍会对比解码出的
collectionID,但那是针对编码或作用域 bug 的不变量检查,不再是语义守卫。
使用方的对等义务与边界
文档与源码同时强调了一个“以文档记载、而非强制”的义务,以及两个坑:
- 锁轴必须覆盖作用域对象:幂等去重的串行化保证,只有在广播持有的排他锁确实覆盖键所作用对象的条件下才成立。Import 满足这一点(集合作用域 + 对该集合的
ExclusiveCollectionName锁);若采用者给对象 A 打键却给对象 B 加锁,两个并发的同键请求就可能同时 miss。 - Broadcaster 自身不做准入检查:调用方在校验阶段做的所有限制都发生在查重之前。如果调用方执行的限额把自己原始请求也算在内(import 的
dataCoord.import.maxImportJobNum正是这种限额),那么它会拒绝原始请求的重试,导致重试无法找回原 broadcastID。对客户端的契约是:限额释放后用同一键重试;若在被拒绝时重新铸造新键,才会真正造成重复工作。 - 唯一例外:原始任务已失败。去重分支把键解析到原 ID 时不会检查该 job 的状态;若原任务以
Failed收场,则窗口期内每次同键重试都会拿回同一个失败 ID,客户端永远无法推进。这符合普通幂等语义(键命名了一次确实发生过的尝试),但却是唯一应该“换新键”而非“复用键”的情形。ImportV2会在每次去重命中时记录原 job 状态日志,便于运维人员区分“卡在失败任务上”的键与“等待健康 job”的键。
幂等窗口 = 墓碑保留期
索引的生命周期与任务条目完全绑定,因此客户端观察到的幂等窗口恰好等于墓碑保留期:maxLifetime或maxCount,先到者为准。count 上限是硬的——繁忙集群可能在maxLifetime远未到达时就提前驱逐墓碑、提前结束窗口。任何对外承诺该保证的子系统(目前是 BulkImport)必须让自己的元数据保留期至少长于maxLifetime。
仅仅“恰好相等”也不够:tombstoneScheduler.Initialize会给每个恢复出来的墓碑盖上time.Now()(tombstone_scheduler.go),即墓碑的年龄从最近一次 StreamingCoord 启动算起,每次重启都会延长其剩余寿命;而子系统自己的保留期是从原始事件持续计数的。务必留出余量。
复制场景的索引
getOrCreateBroadcastTask(broadcast_manager.go)会对从复制 WAL 恢复的 REPLICATED 任务同样建立幂等索引:查询路径在次集群上不可达(WithResourceKeys拒绝非主集群),但在次集群上建索引,能让故障转移后被提升的主集群继续兑现故障前的幂等键。
资源键锁定模型(Resource Key Locking)
每个ResourceKey由三部分组成(proto 定义于messagespb.ResourceKey,构造器见 resource_key.go):
| 字段 | 含义 |
|---|---|
| Domain | 资源类型(Cluster、DBName、CollectionName、Privilege、SnapshotName) |
| Key | 实体标识符(如集合名db/collection) |
| Shared | 读共享锁 vs 写排他锁 |
每个广播都会自动追加NewSharedClusterResourceKey()——除非调用方自己已经带了一个集群域键(见 broadcast_manager.go 的appendSharedClusterRK)。
底层基于 KeyLock,注释明言这是“低性能实现但可接受”——因为 Broadcaster 只在低频 DDL 上使用。实现要点:
Lock(阻塞版)先uniqueSortResourceKeys去重并按(Domain, Key)排序再加锁,从根上避免死锁;Unlock则按逆序释放;FastLock(非阻塞版)用于 ack 回调调度:任一键被占用立即失败返回。调度器在后台循环里用FastLock抢占(见 ack_callback_scheduler.go),抢不到的冲突任务留在队列延迟重试,以此把冲突任务严格串行化、保住 WAL 顺序。
各消息类型对资源键的具体用法,可继续阅读 Message 语义文档 中的各消息说明页(原文档以 链接 引用)。
BroadcastTask 状态机与恢复
文档给出了清晰的状态机:
PENDING → TOMBSTONE → DONE(从 catalog 移除) REPLICATED → TOMBSTONE → DONE(从 catalog 移除)各状态含义与源码对应如下:
- PENDING:任务创建,等待 WAL 追加与 ACK。追加完成后若无
AckSyncUp立即 FastAck 自确认全部 VChannel;否则等待消费端 ACK。 - REPLICATED:次集群上从复制的
ImmutableMessage重建的任务,不持有任何资源锁;其执行顺序由ackCallbackScheduler中 CChannel TimeTick 序保证。构造见 broadcast_task.go。 - TOMBSTONE:所有 ack 回调完成、资源锁已释放,等待 GC。
MarkAckCallbackDone(broadcast_task.go)负责状态迁移、关闭done通道并释放锁守卫。 - DONE:已从 catalog 删除(
DropTombstone,见 broadcast_task.go)。注意其注释警告:墓碑一旦被丢弃,幂等与去重保证即失效。
崩溃恢复路径由RecoverBroadcaster(broadcast_manager.go)承担:从 catalogListBroadcastTask读出全部任务,PENDING/WAIT_ACK 任务用FastLock立即重占资源键(拿不到即 panic,等待运维介入),未完成追加的进入广播调度器续写,其余按状态归入 ack 调度或墓碑列表;幂等索引对包括墓碑在内的每一种状态都重建——因为迟到的重试必须命中的恰恰是已墓碑化的任务。
此外,主从复制还引出一条force promote的完整链路:force promote 消息先等全部 VChannel ACK(复制消息自此被 fence),随后fixIncompleteBroadcastsForForcePromote找出所有未完成任务、对残留的AlterReplicateConfig打上ignore标记(内存级操作,见 broadcast_task.go 的MarkIgnore),再委托广播调度器补写 WAL,最后才执行 ack 回调、释放 RPC。
关键配置参数
Broadcaster 的墓碑 GC 与并发行为由streaming.walBroadcaster.tombstone.*与streaming.WALBroadcasterConcurrencyRatio控制,参数注册见 component_param.go,默认值可从 component_param_test.go 得到印证:
| 参数 Key | 默认值 | 作用 |
|---|---|---|
streaming.walBroadcaster.tombstone.checkInternal | 5m | 墓碑 GC 巡检间隔(可热更新) |
streaming.walBroadcaster.tombstone.maxCount | 8192 | 墓碑数量上限,先到先清(硬边界) |
streaming.walBroadcaster.tombstone.maxLifetime | 24h | 墓碑最长存活时长 |
streaming.WALBroadcasterConcurrencyRatio | 与 CPU 数相乘 | 广播追加 worker 数 |
正如前文所述,这三个墓碑参数共同决定了外部子系统(如 BulkImport)幂等窗口的真实长度;tombstoneScheduler.background的实现见 tombstone_scheduler.go。
测试与可深入阅读的代码路径
Broadcaster 的健壮性高度依赖并发与故障场景测试,仓库中已有成体系的用例可作行为级文档:
- broadcaster_test.go 与 broadcaster_with_rk_guards_test.go——广播主流程与锁守卫所有权;
- idempotency_index_test.go——幂等索引的命中/丢失/移除竞态;
- resource_key_locker_test.go——多键排序加锁与 FastLock 失败路径;
- force_promote_failover_test.go——主从故障转移后的 force promote 修复流程;
- ack_callback_scheduler_trace_test.go 与 pending_broadcast_task_trace_test.go——回调/追加的追踪上下文贯穿。
组件代码集中于 internal/streamingcoord/server/broadcaster/,其中broadcast/子目录是供外部模块引用的单例门面,registry/维护各类消息的 ack 回调注册表;消息模型本身可回溯 Message 模型文档(含 BroadcastMutableMessage 的三阶段生命周期)。
小结
Broadcaster 是 Milvus DDL 数据面的“单点枢纽”:用排序加锁保证资源互斥,用 catalog 先持久化保证崩溃后可续,用 FastAck 与 CChannel TimeTick 序平衡延迟与顺序,用“ID 身份 + 消息类型”的幂等索引把客户端重试收敛为一次真实执行,最后用墓碑机制划定幂等窗口。理解它的状态机与义务边界,是正确实现高可靠 DDL 客户端(如 BulkImport)和排查复制/导入重复问题的前提。
【免费下载链接】milvusMilvus is a high-performance, cloud-native vector database built for scalable vector ANN search项目地址: https://gitcode.com/GitHub_Trending/mi/milvus
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考