多彩编程 多彩编程MZPH · CODE BLOG
ARTICLE DETAIL

文章详情

深耕前端与后端开发技术的一线实战笔记与踩坑复盘。

Apache Pulsar 可插拔元数据存储接口(PIP-45)深度解析:从 ZooKeeper 单一绑定到统一 MetadataStore 架构

Apache Pulsar 可插拔元数据存储接口(PIP-45)深度解析:从 ZooKeeper 单一绑定到统一 MetadataStore 架构 消息队列流处理后端微服务消息路由【免费下载链接】pulsarApache Pulsar - distributed pub-sub messaging system项目地址https://gitcode.com/gh_mirrors/pu/pulsar点击查看免费下载本文以 PIP-45: Pluggable metadata interface 为骨架结合当前仓库pulsar-metadata模块的完整落地实现系统讲解 Apache Pulsar 如何将散落在各处的 ZooKeeper 访问收敛为统一、可插拔的元数据接口。你将掌握MetadataStore、CoordinationService、MetadataCache三层 API 的设计动机、核心方法语义、后端实现原理以及metadataStoreUrl等配置的实战用法并了解这一设计如何在保留 ZooKeeper 数据格式兼容性的前提下支撑 memory、RocksDB、Oxia 等多种元数据后端。背景为什么 Pulsar 需要一个统一的元数据接口在 PIP-45 提出之前Pulsar 的元数据访问存在两个突出问题强绑定 ZooKeeperBroker 与部分管理 CLI 工具都直接通过 ZooKeeper 完成元数据读写和协调leader 选举、负载上报、bundle 所有权等代码中ZooKeeperCache、ZooKeeperDataCache、ZooKeeperChildrenCache等类被广泛使用任何一处更换后端都要大动干戈。访问入口分散ZooKeeper 客户端 API 在代码库多个位置被直接调用缺乏统一抽象的约束难以替换、难以测试。同时Pulsar 的姊妹系统 BookKeeper 早已支持可插拔的元数据存储pluggable metadata store这为 Pulsar 提供了可借鉴的成熟先例。PIP-45 的核心理念因此非常直接定义一个统一的MetadataStore接口把 Pulsar 与元数据的所有交互抽象到这一层让后端实现ZooKeeper、Etcd、内存、本地磁盘等可插拔。该提案提出了明确的兼容性约束重构完成后默认实现仍然基于 ZooKeeper与既有元数据 100% 兼容——数据存放在相同位置、使用完全相同的格式因此对现有集群而言是无感升级无需迁移数据。值得注意的是PIP-45 在提出时即将 API 定性为 Beta在积累足够多的具体实现之前接口允许以破坏性方式演进且至少在初期它属于 Pulsar 内部 API不对用户插件开放。这一点在今天的源码中依然可见——MetadataStore.java 类注释明确标注Beta并说明 This API is still evolving and will be refactored as needed until all the metadata usages are converted into using it。重构路线五个渐进步骤PIP-45 将整体重构拆解为五个步骤每一层都解决一个独立问题定义元数据存储 APIMetadataStore最底层的键值接口定义协调接口CoordinationService在键值之上抽象分布式锁、leader 选举与计数器将 ManagedLedger 移植到MetadataStoreManagedLedger 已有自己的元数据抽象MetaStore转换成本低定义元数据缓存 APIMetadataCache把ZooKeeperCache系迁移到新接口并引入原子读-改-写能力定义更高层的元数据抽象如ConfigurationStore避免业务代码直接面向裸getChildren。下面逐一深入每一层的设计与当前仓库中的实现形态。第一步MetadataStore——基于版本 CAS 的统一键值接口PIP-45 明确将元数据存储建模为基于 Key-Value 的接口核心操作语义是带版本校验的compareAndSet()写入时可以通过期望版本做原子性校验版本不匹配即失败。提案给出的接口原型如下public interface MetadataStore extends AutoCloseable { CompletableFutureOptionalGetResult get(String path); CompletableFutureListString getChildren(String path); CompletableFutureBoolean exists(String path); CompletableFutureStat put(String path, byte[] value, OptionalLong expectedVersion); CompletableFutureVoid delete(String path, OptionalLong expectedVersion); }几个关键语义值得展开异步贯穿始终所有操作均返回CompletableFuture与 Pulsar 全异步的执行模型一致get的语义按 path 读取值成功时返回GetResult内含值 关联的Stat键不存在时返回空OptionalgetChildren的语义返回按字典序排序的子节点列表路径本身不存在时返回空列表exists的注意点对于多级路径如/a/b/c检查父级如/a不一定返回 true除非该父键被显式创建过put/delete与expectedVersion这是 CAS 的核心。期望版本存在时必须与当前存储版本一致操作才成功-1用于强制要求键不存在即仅创建语义异常语义版本不匹配抛BadVersionException路径不存在抛NotFoundException。当前仓库中的完整实现PIP-45 落地后MetadataStore接口被置于独立的pulsar-metadata模块中路径为 MetadataStore.java。与提案原型相比当前接口在保持核心语义不变的前提下做了系统性扩展Option机制每个操作都支持传入SetOption例如Option.PartitionKey分片后端的路由键、Option.Ephemeral会话结束后自动过期、Option.Sequential创建唯一递增后缀键、Option.SecondaryIndex二级索引等并通过Set.of()提供无选项的默认重载方法getChildrenFromStore类似getChildren但绕过本地缓存、直接读底层存储deleteIfExists删除时若路径不存在则不视为错误日志记录后正常完成deleteRecursive删除键值对及其全部子节点——注意此操作可能不是原子的失败时可能只完成部分删除registerListener注册ConsumerNotification监听底层存储变更——这正是 PIP-45 中创建MetadataStore时指定观察者函数子树变化即触发回调无需轮询即可保持本地缓存新鲜设计的最终形态。当前实现中NotificationType枚举包含Created、Modified、ChildrenChanged、Deleted四类事件见 NotificationType.javasync(path)确保下一次本地读取能拿到与所有其他客户端一致的、最新版本的值getMetadataCache直接从 store 派生指定类型的元数据缓存见第四步扩展扫描与订阅能力scanByIndex按二级索引做点查或范围扫描支持在原生索引后端一次扫描或在 ZooKeeper 等后端回退为列子节点 逐条读取 过滤器、scanChildren带值流式列出直接子节点、subscribeSequence订阅序列表前缀上的更新ZooKeeper 后端从父路径变更通知合成事件流。此外还有MetadataStoreExtended接口见 extended 包其中CreateOption枚举补充了Ephemeral会话结束自动过期与Sequential创建顺序唯一键两类创建选项——这两者在 ZooKeeper 场景正是提案提到的ephemeral 节点与顺序节点的抽象。GetResultGetResult.java与StatStat.java分别封装读取结果与版本/时间戳等元信息异常层次则统一收敛在 MetadataStoreException.java内含BadVersionException、NotFoundException、AlreadyExistsException、InvalidImplementationException等子类。第二步CoordinationService——分布式锁与 leader 选举的抽象键值接口足以表达协调语义例如用 ZooKeeper 的 ephemeral 节点实现锁但每个后端系统可能有更直接的原生实现如 Etcd 的租约、Oxia 的锁服务。PIP-45 因此单独定义了一层协调接口public interface CoordinationService extends AutoCloseable { CompletableFutureOptionalbyte[] readLock(String path); CompletableFutureResourceLock acquireLock(String path, byte[] content); CompletableFutureListString listLocks(String path); CompletableFutureResourceLock becomeLeader(String path, byte[] content); CompletableFutureLong getNextCounterValue(String path); }提案中还定义了锁句柄ResourceLock其能力包括getContent()读取锁携带的负载、release()释放锁、getLockExpiredFuture()获取锁失效通知注意主动释放不会触发该 future它只用于会话过期等被动失效场景。提案明确列出了 Pulsar broker 使用协调的典型场景这些场景在今天仍然成立活跃 broker 列表及其负载数据上报获取 namespace 下某段 topicbundle的所有权leader 选举load manager用于生成唯一前缀标识符的计数器。接口文档中还给出了一条重要的工程警告由于锁的分布式本质即使成功获得锁也永远无法强保证没有其他进程认为自己拥有同一资源调用方必须在使用资源时自行处理竞态例如借助compareAndSet()或 fencing 机制。这条警示被原样保留在当前 CoordinationService.java 的实现中。当前仓库中的协调实现落地后的协调层位于 coordination 包接口演化为LeaderElectionT通过getLeaderElection(ClassT clazz, String path, ConsumerLeaderElectionState stateChangesListener)创建以泛型反序列化负载并通过状态监听器感知LeaderElectionState变迁——这正是提案中becomeLeader语义的面向对象化封装LockManagerT通过getLockManager(ClassT clazz)或getLockManager(MetadataSerdeT serde)获取用于管理指定类型负载的分布式锁对应提案的acquireLock/readLock/listLocksgetNextCounterValue(String path)保证在 path 上下文内唯一递增的计数器。当前实现增加了自动重试语义遇MetadataStoreException会重试达到最大重试次数仍失败才以异常完成 future。其默认实现位于 coordination/impl包含CoordinationServiceImpl、LeaderElectionImpl、LockManagerImpl、ResourceLockImpl等类ResourceLock接口定义见 ResourceLock.java。第三步ManagedLedger 迁移到 MetadataStore提案指出ManagedLedger 在 PIP-45 之前就已经拥有自己的元数据访问抽象MetaStore位于managed-ledger模块因此将其转换到统一的MetadataStoreAPI 属于低成本改造——这正是 PIP-45 能够以渐进式重构推进的基础每一层改动都有既有抽象可对照、可单测。当前仓库中managed-ledger/src/main已全面围绕新的元数据层工作而pulsar-metadata模块还为此提供了 BookKeeper 元数据驱动适配层见 bookkeeper 包含PulsarLedgerManager、PulsarLedgerIdGenerator、PulsarLayoutManager、PulsarRegistrationClient等说明元数据抽象不仅覆盖 broker 侧也覆盖了 BookKeeper 侧的 ledger 元数据验证了 PIP-45 与 BookKeeper 既有可插拔元数据设计对齐的方向。第四步MetadataCache——类型化缓存与原子读-改-写提案指出重构前绝大多数元数据读操作都经过ZooKeeperCache及其派生类ZooKeeperDataCache、ZooKeeperChildrenCache。这些缓存需要移植到MetadataStoreAPI 之上同时保留通知接收与失效条目清理的能力。更重要的是提案为缓存层新增了一个关键概念原子读-改-写read-modify-update操作。此前代码库中多处各自实现这一逻辑应当收敛为单一实现public interface TypedMetadataCacheT { // ... CompletableFutureVoid readModifyUpdate(String path, FunctionT, T modifyFunction); }modifyFunction在存在并发更新时可能被调用多次这正是乐观重试语义先读取当前值应用修改函数再以期望版本写回若版本冲突则重读重试。当前仓库中的缓存实现今天的 MetadataCache.java 完整继承了这一设计并进一步泛化get(path)/getWithStats(path)/getIfCached(path)先查缓存、未命中再回源存储的读取模式getWithStats额外返回CacheGetResultT以携带缓存命中统计getChildren(path)/exists(path)缓存化的子树列表与存在性查询readModifyUpdate(path, FunctionT, T)原子读-改-写冲突自动重试readModifyUpdateOrCreate(path, FunctionOptionalT, T)变体——对象尚不存在时修改函数接收Optional.empty()从而把创建或更新统一进同一原子语义。缓存由MetadataStore.getMetadataCache(...)直接创建支持按ClassT、TypeReferenceT或自定义MetadataSerdeT反序列化并可通过MetadataCacheConfig调节缓存行为。类型化意味着缓存中直接存放已反序列化的对象避免每次读取都做序列化/反序列化——这与提案ZooKeeperCache 已经是类型化的描述一脉相承。第五步更高层的元数据抽象ConfigurationStore提案认为仅靠MetadataStore.getChildren()直接获取业务数据如集群租户列表依然太底层容易在代码库多处重复出现。因此需要一层面向业务的抽象例如public interface ConfigurationStore { CompletableFutureListString getTenants(); CompletableFutureListString getNamespaces(String tenant); // .... }这层抽象的价值在于业务代码只面向语义化方法给我租户列表不关心路径约定、存储后端差异与缓存策略也便于在引入新后端时只改实现、不改调用方。从当前仓库的模块划分看这种分层思想已经渗透到整个pulsar-metadataAPI 设计中底层MetadataStore、缓存层MetadataCache、协调层CoordinationService/LeaderElection/LockManager以及 BookKeeper 侧的管理类各自职责单一、可独立替换。实现落地MetadataStore 的四种后端与工厂机制PIP-45 在 Goals 中展望了四类后端ZooKeeper、Etcd、内存供单元测试、本地磁盘供 standalone 使用。当前仓库的 MetadataStoreFactoryImpl.java 展示了这一愿景的最终落地与演进URL 前缀实现类用途zk:ZKMetadataStore.java默认后端兼容既有 ZK 数据内部基于 BookKeeper 的 ZooKeeper 客户端采用批处理策略见 batching 包合并异步操作memory:LocalMemoryMetadataStore.java内存实现基于NavigableMap 原子版本号主要用于单元测试rocksdb:RocksdbMetadataStore.java本地磁盘实现用于 Pulsar standalone 等场景oxia:OxiaMetadataStore.java面向 Apache Oxia 的分布式元数据后端工厂的加载逻辑值得注意URL 协议分发MetadataStoreFactory.create(String metadataURL, MetadataStoreConfig)解析 URL 前缀并匹配到对应MetadataStoreProvider未指定协议如my-zk-1:2181时默认回退到 ZooKeeper保证老配置零改动可用可扩展性支持通过系统属性pulsar.metadatastore.providers注册自定义MetadataStoreProvider类按逗号分隔类名实现真正的可插拔历史沿革源码中保留了etcd:协议的显式拒绝分支——Etcd 后端已在 Pulsar 5.0 中移除PIP-462并提示改用zk:或oxia:。这恰好印证了 PIP-45 API 初期按 Beta 演进、后端可增删的设计预期。另外DualMetadataStore.java 提供了双写/双读的迁移桥接实现可用于集群从一种元数据后端平滑迁移到另一种是兼容性与可演进性工程思想的又一落地。实战配置standalone 与 broker 中的 metadataStoreUrlPIP-45 抽象带来的最直接收益是部署时只需改一个 URL 即可切换元数据后端。conf/standalone.confstandalone.conf中给出了完整注释示例# The metadata store URL # Examples: # * zk:my-zk-1:2181,my-zk-2:2181,my-zk-3:2181 # * my-zk-1:2181,my-zk-2:2181,my-zk-3:2181 (will default to ZooKeeper when the schema is not specified) # * zk:my-zk-1:2181,my-zk-2:2181,my-zk-3:2181/my-chroot-path (to add a ZK chroot path) metadataStoreUrl # Configuration file path for metadata store. Its supported by RocksdbMetadataStore and EtcdMetadataStore for now metadataStoreConfigPath # The metadata store URL for the configuration data. If empty, we fall back to use metadataStoreUrl configurationMetadataStoreUrl三个配置项的用法metadataStoreUrl主元数据存储地址。生产集群多写为zk:host1:2181,host2:2181,host3:2181standalone 可用rocksdb:/path/to/dir配合metadataStoreConfigPath指向default_rocksdb.conf、entry_location_rocksdb.conf等 RocksDB 配置测试环境可直用memory:metadataStoreConfigPath元数据存储的配置文件路径当前主要服务于 RocksDB 等后端configurationMetadataStoreUrl配置类元数据的独立存储地址留空时回退使用metadataStoreUrl这也是configuration 与 coordination 元数据可分离部署的体现conf/broker.conf中另有已废弃的configurationStoreServers旧参数。切换后端的约束前提是所有 Pulsar 组件broker、proxy、工具使用同一套元数据后端配置且需满足 PIP-45 承诺的格式兼容约定——默认 ZooKeeper 路径布局与数据格式不变因此存量集群迁移风险可控。演进观察从提案到现实的差距与延续将 PIP-45 的原始设计与当前源码对照可以清晰看到一条持续的演进主线接口扩充而非推翻MetadataStore的核心五个方法get/getChildren/exists/put/delete语义原样保留只是在之上叠加了Option、sync、递归删除、监听注册、扫描订阅与缓存工厂等能力协调层从扁平方法演化为LeaderElection/LockManager对象模型但提案中的锁、leader 选举、计数器三类能力全部保留后端比预期更丰富除了提案展望的 ZooKeeper、Etcd、内存、本地磁盘四种还引入了 Oxia 后端而 Etcd 因维护原因被移除——说明抽象层确实让后端可增可删Beta 约束的兑现API 至今标注BetaMetadataStore注释明确声明将持续重构直到所有元数据使用方完成迁移且初期不向用户插件开放自定义 provider 需通过系统属性显式注册与提案约定完全一致。总结PIP-45 为 Apache Pulsar 奠定了一层承上启下的元数据抽象底层MetadataStore解决存储可替换中层CoordinationService解决协调语义统一缓存与高层抽象解决业务代码与存储解耦而默认 ZooKeeper、格式 100% 兼容的约束保证了迁移安全。对于希望深入 Pulsar 内部机制或二次开发元数据层的工程师建议按以下路径继续阅读源码接口全貌MetadataStore.java、CoordinationService.java、MetadataCache.java工厂与后端MetadataStoreFactoryImpl.java、ZKMetadataStore.java、LocalMemoryMetadataStore.java、RocksdbMetadataStore.java协调实现coordination/impl部署配置standalone.conf 中metadataStoreUrl相关段落。这一设计不仅让 Pulsar 摆脱了对单一 ZooKeeper 的硬绑定也为集群元数据的高可用、多租户隔离与运维可观测性提供了统一的扩展点——这正是 PIP-45 从提案走向架构基石的价值所在。赞分享消息队列流处理后端微服务消息路由【免费下载链接】pulsarApache Pulsar - distributed pub-sub messaging system项目地址https://gitcode.com/gh_mirrors/pu/pulsar点击查看免费下载相关推荐Apache Pulsar PIP-335 深度解析Oxia 元数据存储插件与去 ZooKeeper 化路线Apache Pulsar PIP 335 深度解析Oxia 元数据存储插件与去 ZooKeeper 化路线 本文基于 Pulsar 官方提案文档 pip/p消息队列流处理后端微服务消息路由Apache Pulsar PIP-378 深度解析ServiceUnitStateTableView 抽象与可插拔的 Bundle 归属事件源存储Apache Pulsar PIP 378 深度解析ServiceUnitStateTableView 抽象与可插拔的 Bundle 归属事件源存储 导读 本消息队列流处理后端微服务消息路由Apache Pulsar 元数据存储迁移框架PIP-454从 ZooKeeper 到 Oxia 的零停机迁移完整指南Apache Pulsar 元数据存储迁移框架PIP 454从 ZooKeeper 到 Oxia 的零停机迁移完整指南 导读 Apache Pulsar消息队列流处理后端微服务消息路由创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
返回列表