ARTICLE DETAIL

资讯详情

深耕编程入门与网站建设的一线实战洞察。

Milvus StreamingCoord Broadcaster 源码级解析:跨 PChannel 的原子 DDL/DCL 广播与幂等保证

Milvus StreamingCoord Broadcaster 源码级解析:跨 PChannel 的原子 DDL/DCL 广播与幂等保证 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.goWithResourceKeys(ctx, resourceKeys...)按资源键加锁并返回BroadcastAPI若当前集群不是主集群non-primary直接返回ErrNotPrimaryWithSecondaryClusterResourceKey(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() // 若未发起 BroadcastClose 负责释放已持有的资源键锁 // 2. 构造 BroadcastMutableMessage目标 VChannels 必须包含 CChannel控制通道 msg : message.NewBroadcastMutableMessage(...) // 3. 执行广播 result, err : api.Broadcast(ctx, msg)关键约定来自 broadcaster_with_rk.goBroadcast()一旦执行句柄上的锁守卫会被立即“消费”所有权移交要么由新注册的任务持有并在其 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状态的任务并把消息落盘到 catalogetcd。一旦持久化成功即使进程崩溃该广播也保证最终完成——恢复时会从 catalog 重建。相关逻辑集中在 broadcast_task.go构造 PENDING 任务、标记 dirty和saveTaskIfDirty写SaveBroadcastTask。追加Append任务交给broadcastScheduler。调度器按硬件 CPU 数 × streaming.WALBroadcasterConcurrencyRatio启动一批 worker见 broadcast_scheduler.goworker 调用streaming.WAL().AppendMessages()把消息按 VChannel 拆条写入目标 PChannelpending_broadcast_task.go。部分追加失败的消息会被保留并进入指数退避重试队列直到全部成功。快速确认FastAck追加全部成功后触发FastAckbroadcast_task.go。若广播头的AckSyncUp未置位Broadcaster 直接依据 append 结果“自确认”所有 VChannel不必等待消费端确认从而显著降低 DDL 延迟若AckSyncUp置位则必须等待各 StreamingNode 消费者真正 ACK 每个 VChannel。Ack 回调AckCallbackCChannel控制通道被 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 BroadcastBroadcaster 的幂等设计是文档中着墨最深的部分也是整个组件的精髓。它解决的是“客户端因超时/网络抖动重试同一个 DDL是否会执行两次”的问题。幂等键的编码与作用域携带_ik幂等键属性属性常量见 properties.go的广播消息会被额外按该键建立索引。幂等键不是一个裸字符串而是一个自描述的三段编码domain:scopeID:clientKey其中 domain 取值为集群1/数据库2/集合3scopeID 是对象的数字 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.goscope 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而非名字作为作用域正是为了同时解决经典的名字-身份问题改名RenameRenameCollection保持 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 启动算起每次重启都会延长其剩余寿命而子系统自己的保留期是从原始事件持续计数的。务必留出余量。复制场景的索引getOrCreateBroadcastTaskbroadcast_manager.go会对从复制 WAL 恢复的 REPLICATED 任务同样建立幂等索引查询路径在次集群上不可达WithResourceKeys拒绝非主集群但在次集群上建索引能让故障转移后被提升的主集群继续兑现故障前的幂等键。资源键锁定模型Resource Key Locking每个ResourceKey由三部分组成proto 定义于messagespb.ResourceKey构造器见 resource_key.go字段含义Domain资源类型Cluster、DBName、CollectionName、Privilege、SnapshotNameKey实体标识符如集合名db/collectionShared读共享锁 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。MarkAckCallbackDonebroadcast_task.go负责状态迁移、关闭done通道并释放锁守卫。DONE已从 catalog 删除DropTombstone见 broadcast_task.go。注意其注释警告墓碑一旦被丢弃幂等与去重保证即失效。崩溃恢复路径由RecoverBroadcasterbroadcast_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.checkInternal5m墓碑 GC 巡检间隔可热更新streaming.walBroadcaster.tombstone.maxCount8192墓碑数量上限先到先清硬边界streaming.walBroadcaster.tombstone.maxLifetime24h墓碑最长存活时长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),仅供参考
返回列表