Apache Pulsar 自定义 Schema 存储:实现 SchemaStorage 与 SchemaStorageFactory 接口指南
发布时间:2026/9/25 6:53:13来源:尧图网络
消息队列后端流处理【免费下载链接】pulsarApache Pulsar - distributed pub-sub messaging system项目地址https://gitcode.com/gh_mirrors/pulsar28/pulsar点击查看免费下载本文基于 Pulsar 官方开发文档《Custom schema storage》讲解如何为 Pulsar 的 Schema 注册中心Schema Registry替换默认存储后端如何设计并实现SchemaStorage与SchemaStorageFactory两个 Java 接口如何将其打包部署到 Pulsar 发行版中以及结合仓库源码说明 Broker 启动时是如何通过反射加载并启动自定义 Schema 存储的。读完后你可以为自己的部署环境例如已有 Redis、etcd 或分布式数据库编写一套完整的 Schema 持久化方案并理解put/get/delete各方法与版本号语义背后的实现约束。默认情况下Schema 存放在 BookKeeper 中Pulsar 的 Schema 注册中心负责保存每个主题topic上消息的数据类型定义schema。默认实现中这些 schema 定义被存储在与 Pulsar 一起部署的 Apache BookKeeper 上。这一默认行为由 Broker 配置项schemaRegistryStorageClassName控制在 ServiceConfiguration.java 中该配置项的默认值被定义为FieldContext( category CATEGORY_SCHEMA, doc The schema storage implementation used by this broker ) private String schemaRegistryStorageClassName org.apache.pulsar.broker.service.schema .BookkeeperSchemaStorageFactory;发行版默认配置文件 conf/broker.conf第 1350 行与 conf/standalone.conf第 938 行中也都显式写出了同样的取值schemaRegistryStorageClassNameorg.apache.pulsar.broker.service.schema.BookkeeperSchemaStorageFactory因此要使用非 BookKeeper 的存储系统就必须提供自己的实现并替换上述配置。根据官方文档需要实现两个 Java 接口SchemaStorage存储客户端和SchemaStorageFactory存储工厂。SchemaStorage 接口存储客户端的契约SchemaStorage接口定义在pulsar-common模块中SchemaStorage.java。文档给出的最小方法集如下public interface SchemaStorage { // How schemas are updated CompletableFutureSchemaVersion put(String key, byte[] value, byte[] hash); // How schemas are fetched from storage CompletableFutureStoredSchema get(String key, SchemaVersion version); // How schemas are deleted CompletableFutureSchemaVersion delete(String key); // Utility method for converting a schema version byte array to a SchemaVersion object SchemaVersion versionFromBytes(byte[] version); // Startup behavior for the schema storage client void start() throws Exception; // Shutdown behavior for the schema storage client void close() throws Exception; }当前仓库中的接口比文档版本略有扩展文档对应 2.3.0 版本接口会随版本演进完整定义还包括以下内容public interface SchemaStorage { CompletableFutureSchemaVersion put(String key, byte[] value, byte[] hash); /** * Put the schema to the schema storage. * param key The schema ID * param fn The function to calculate the value and hash that need to put to the schema storage * The input of the function is all the existing schemas that used to do the schemas compatibility check */ default CompletableFutureSchemaVersion put(String key, FunctionCompletableFutureListCompletableFutureStoredSchema, CompletableFuturePairbyte[], byte[] fn) { return fn.apply(getAll(key)).thenCompose(pair - put(key, pair.getLeft(), pair.getRight())); } CompletableFutureStoredSchema get(String key, SchemaVersion version); CompletableFutureListCompletableFutureStoredSchema getAll(String key); CompletableFutureSchemaVersion delete(String key, boolean forcefully); CompletableFutureSchemaVersion delete(String key); SchemaVersion versionFromBytes(byte[] version); void start() throws Exception; void close() throws Exception; }各方法的设计意图与实现要点如下方法作用实现注意点put(String key, byte[] value, byte[] hash)写入/更新一个 schemakey为 schema 标识通常是主题名value为序列化后的 schema 定义hash为其内容哈希返回本次写入产生的版本号返回值必须是CompletableFutureSchemaVersion异步语义是契约的一部分同一key多次写入应产生递增/可区分的版本put(key, fn)default 方法在写入前先拿到该 key 下所有历史版本 schema 用于兼容性检查由调用方回调fn计算最终的value与hash再落盘有默认实现先调getAll(key)把全部历史 schema 传给fn再用其结果调用三参put。自定义实现如果只实现了三参put和getAll即可复用该逻辑get(String key, SchemaVersion version)按版本读取一个已存储的 schema返回CompletableFutureStoredSchema需要能根据版本定位到具体一条StoredSchema记录getAll(String key)取回某 key 下的全部历史 schema 版本是默认put(key, fn)做兼容性检查的数据来源注意返回类型是CompletableFutureListCompletableFutureStoredSchema——外层一个 future 包装“版本清单”内层每个 future 才是真正读取某个版本的 schema实现时通常先取版本索引再并发拉取各版本内容delete(String key, boolean forcefully)/delete(String key)删除某 key 下的 schema返回删除后的版本状态两个重载并存简单实现可以让delete(key)委托给delete(key, forcefully)versionFromBytes(byte[] version)把版本号的字节数组还原为SchemaVersion对象与SchemaVersion.bytes()互逆版本号在存储层就是一段byte[]start()客户端启动逻辑建立连接、初始化内部状态等Broker 启动流程中会被显式调用close()客户端关闭逻辑释放连接与资源与start()对称异常需向外抛出以便上层感知支撑类型SchemaVersion 与 StoredSchema接口中的关键值类型同在 pulsar-common 包下SchemaVersion.java版本号抽象只有一个方法byte[] bytes()并内置两个特殊常量public interface SchemaVersion { SchemaVersion Latest new LatestVersion(); SchemaVersion Empty new EmptyVersion(); byte[] bytes(); }即Latest表示“最新版本”、Empty表示“无版本”你的存储层需要理解这两种非具体版本的语义例如get(key, Latest)要能解析出该 key 的最新版 schema。StoredSchema.java一次存储结果的载体包含 schema 内容、其哈希与对应版本号get方法的返回值就是它。BytesSchemaVersion.javaSchemaVersion的通用字节数组实现versionFromBytes可以直接构造它返回。参考实现BookkeeperSchemaStorage官方文档建议以 BookKeeper 版实现作为完整范例BookkeeperSchemaStorage.java位于pulsar-broker模块的org.apache.pulsar.broker.service.schema包中。编写自定义实现前建议通读该类的put/get/getAll/delete实现观察它如何组织版本键version key与 schema 键schema key、如何处理Latest与Empty版本、以及versionFromBytes与bytes()的互逆关系。SchemaStorageFactory 接口Broker 加载存储的入口SchemaStorageFactory接口定义在 Broker 侧SchemaStorageFactory.javapublic interface SchemaStorageFactory { NotNull SchemaStorage create(PulsarService pulsar) throws Exception; }工厂的职责只有一个接收正在初始化的PulsarService实例创建并返回一个SchemaStorage客户端。之所以需要工厂层而不是直接实例化SchemaStorage是因为存储客户端通常依赖 Broker 运行时提供的资源例如 BookKeeper 客户端从PulsarService中获取工厂模式把“依赖注入”这一步标准化了。官方默认工厂 BookkeeperSchemaStorageFactory.java 展示了最简形态SuppressWarnings(unused) public class BookkeeperSchemaStorageFactory implements SchemaStorageFactory { Override NotNull public SchemaStorage create(PulsarService pulsar) { return new BookkeeperSchemaStorage(pulsar); } }自定义实现照此办理工厂负责持有你的存储配置连接串、地址等create里把PulsarService和你的配置一起传入存储客户端构造函数。Broker 如何加载你的自定义存储反射调用链理解部署机制的关键在于 Broker 启动时的加载代码。PulsarService.java 中的createAndStartSchemaStorage约第 1303–1312 行完整揭示了加载契约private SchemaStorage createAndStartSchemaStorage() throws Exception { final Class? storageClass Class.forName(config.getSchemaRegistryStorageClassName()); Object factoryInstance storageClass.getDeclaredConstructor().newInstance(); Method createMethod storageClass.getMethod(create, PulsarService.class); SchemaStorage schemaStorage (SchemaStorage) createMethod.invoke(factoryInstance, this); schemaStorage.start(); return schemaStorage; }从这段源码可以确认四条硬性要求schemaRegistryStorageClassName配置的是工厂类SchemaStorageFactory实现而不是SchemaStorage实现本身——Class.forName加载它再反射调用其无参构造函数工厂类必须有可访问的无参构造器getDeclaredConstructor().newInstance()工厂类必须存在签名为create(PulsarService)的方法且返回值可转换为SchemaStorage返回的SchemaStorage会立即被调用start()因此你的启动逻辑建连、初始化缓存等必须能在 Broker 启动上下文中完成失败会直接导致 Broker 启动失败。部署让自定义存储在集群中生效根据官方文档的 Deployment 章节完整部署步骤为打包把你的SchemaStorage与SchemaStorageFactory实现含第三方依赖打包成一个 JAR 文件放置将该 JAR 放入 Pulsar 二进制或源码发行版的lib目录中使 Broker 类路径可以加载到你的类配置修改broker.conf中的schemaRegistryStorageClassName指向你的工厂类全限定名再次强调是SchemaStorageFactory实现类不是SchemaStorage实现类schemaRegistryStorageClassNamecom.example.pulsar.MyCustomSchemaStorageFactory启动重启/启动 PulsarBroker 会在启动流程中经由createAndStartSchemaStorage反射加载你的工厂并启动存储客户端。验证方式启动日志中不应出现ClassNotFoundException或反射调用异常之后创建主题并写入带 schema 验证配置的消息即可在自定义后端中观察到 schema 记录的产生。编写自定义实现时的实践要点异步语义要贯穿到底所有读写方法都返回CompletableFuture不要在put/get内部做阻塞式 I/O 而不切换到异步否则可能耗尽 Broker 的事件循环线程。getAll是兼容性检查的数据源Pulsar 的 schema 兼容性校验依赖某 key 下全部历史 schemagetAll返回的“版本索引 逐版本 future”结构是默认put(key, fn)的前提实现不当会直接破坏 schema 升级流程。版本号的字节化必须可逆SchemaVersion.bytes()与versionFromBytes要构成严格互逆映射且你的存储中按版本检索get(key, version)要能高效定位Latest/Empty这类逻辑版本不应被当作普通字节版本原样落盘。start/close的异常要如实抛出Broker 启动与优雅停机流程依赖这两个回调的异常信号。保持与版本对齐本文依据当前仓库源码说明接口含getAll、delete(key, forcefully)与默认put(key, fn)2.3.0 文档所示接口为较早的最小集若你的目标部署版本较低请以对应版本的SchemaStorage定义为准实现避免引用旧版本不存在的方法。小结Pulsar 将 Schema 注册中心的存储层抽象为“工厂 存储客户端”两级接口SchemaStorageFactory负责在 Broker 启动时被反射实例化并产出存储客户端SchemaStorage定义 schema 的异步增删查与版本转换契约配置项schemaRegistryStorageClassName是唯一入口。以 BookkeeperSchemaStorage.java 为参照实现自己的后端打包进发行版lib目录并更新 Broker 配置后重启即可让 Pulsar 的 Schema 持久化落到任意你已有的分布式存储系统上。赞分享消息队列后端流处理【免费下载链接】pulsarApache Pulsar - distributed pub-sub messaging system项目地址https://gitcode.com/gh_mirrors/pulsar28/pulsar点击查看免费下载相关推荐Apache Pulsar 自定义 Schema 存储开发指南实现 SchemaStorage 与 SchemaStorageFactory 接口Apache Pulsar 自定义 Schema 存储开发指南实现 SchemaStorage 与 SchemaStorageFactory 接口 在 Apa消息队列后端流处理Apache Pulsar 自定义 Schema 存储Custom Schema Storage开发实战指南Apache Pulsar 自定义 Schema 存储Custom Schema Storage开发实战指南 Apache Pulsar 默认将 Topic消息队列后端流处理Apache Pulsar Schema 管理完全指南AutoUpdate、手动注册与自定义存储实战Apache Pulsar Schema 管理完全指南AutoUpdate、手动注册与自定义存储实战 本指南以 Apache Pulsar 的 Schema消息队列后端流处理上一篇如何快速入门大语言模型评估HuggingFace evaluation-guidebook新手必读下一篇Vial-QMK 项目常见问题解决方案创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
网站建设高端定制企业官网