☰
Apache Pulsar 自定义 Schema 存储:实现 SchemaStorage 与 SchemaStorageFactory 接口指南
2026/9/25 6:53:07 网站建设 项目流程
  • 消息队列
  • 后端
  • 流处理

【免费下载链接】pulsar

Apache 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 行)中也都显式写出了同样的取值:
schemaRegistryStorageClassName=org.apache.pulsar.broker.service.schema.BookkeeperSchemaStorageFactory

因此,要使用非 BookKeeper 的存储系统,就必须提供自己的实现并替换上述配置。根据官方文档,需要实现两个 Java 接口:SchemaStorage(存储客户端)和SchemaStorageFactory(存储工厂)。

SchemaStorage 接口:存储客户端的契约

SchemaStorage接口定义在pulsar-common模块中:SchemaStorage.java。文档给出的最小方法集如下:

public interface SchemaStorage { // How schemas are updated CompletableFuture<SchemaVersion> put(String key, byte[] value, byte[] hash); // How schemas are fetched from storage CompletableFuture<StoredSchema> get(String key, SchemaVersion version); // How schemas are deleted CompletableFuture<SchemaVersion> 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 { CompletableFuture<SchemaVersion> 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 CompletableFuture<SchemaVersion> put(String key, Function<CompletableFuture<List<CompletableFuture<StoredSchema>>>, CompletableFuture<Pair<byte[], byte[]>>> fn) { return fn.apply(getAll(key)).thenCompose(pair -> put(key, pair.getLeft(), pair.getRight())); } CompletableFuture<StoredSchema> get(String key, SchemaVersion version); CompletableFuture<List<CompletableFuture<StoredSchema>>> getAll(String key); CompletableFuture<SchemaVersion> delete(String key, boolean forcefully); CompletableFuture<SchemaVersion> delete(String key); SchemaVersion versionFromBytes(byte[] version); void start() throws Exception; void close() throws Exception; }

各方法的设计意图与实现要点如下:

方法作用实现注意点
put(String key, byte[] value, byte[] hash)写入/更新一个 schema,key为 schema 标识(通常是主题名),value为序列化后的 schema 定义,hash为其内容哈希;返回本次写入产生的版本号返回值必须是CompletableFuture<SchemaVersion>,异步语义是契约的一部分;同一key多次写入应产生递增/可区分的版本
put(key, fn)(default 方法)在写入前先拿到该 key 下所有历史版本 schema 用于兼容性检查,由调用方回调fn计算最终的value与hash再落盘有默认实现:先调getAll(key)把全部历史 schema 传给fn,再用其结果调用三参put。自定义实现如果只实现了三参put和getAll,即可复用该逻辑
get(String key, SchemaVersion version)按版本读取一个已存储的 schema,返回CompletableFuture<StoredSchema>需要能根据版本定位到具体一条StoredSchema记录
getAll(String key)取回某 key 下的全部历史 schema 版本,是默认put(key, fn)做兼容性检查的数据来源注意返回类型是CompletableFuture<List<CompletableFuture<StoredSchema>>>——外层一个 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.java:SchemaVersion的通用字节数组实现,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.java:

public 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; }

从这段源码可以确认四条硬性要求:

  1. schemaRegistryStorageClassName配置的是工厂类(SchemaStorageFactory实现),而不是SchemaStorage实现本身——Class.forName加载它,再反射调用其无参构造函数;
  2. 工厂类必须有可访问的无参构造器(getDeclaredConstructor().newInstance());
  3. 工厂类必须存在签名为create(PulsarService)的方法,且返回值可转换为SchemaStorage;
  4. 返回的SchemaStorage会立即被调用start(),因此你的启动逻辑(建连、初始化缓存等)必须能在 Broker 启动上下文中完成,失败会直接导致 Broker 启动失败。

部署:让自定义存储在集群中生效

根据官方文档的 Deployment 章节,完整部署步骤为:

  1. 打包:把你的SchemaStorage与SchemaStorageFactory实现(含第三方依赖)打包成一个 JAR 文件;

  2. 放置:将该 JAR 放入 Pulsar 二进制或源码发行版的lib目录中,使 Broker 类路径可以加载到你的类;

  3. 配置:修改broker.conf中的schemaRegistryStorageClassName,指向你的工厂类全限定名(再次强调:是SchemaStorageFactory实现类,不是SchemaStorage实现类):

    schemaRegistryStorageClassName=com.example.pulsar.MyCustomSchemaStorageFactory
  4. 启动:重启/启动 Pulsar,Broker 会在启动流程中经由createAndStartSchemaStorage反射加载你的工厂并启动存储客户端。

验证方式:启动日志中不应出现ClassNotFoundException或反射调用异常;之后创建主题并写入带 schema 验证配置的消息,即可在自定义后端中观察到 schema 记录的产生。

编写自定义实现时的实践要点

  • 异步语义要贯穿到底:所有读写方法都返回CompletableFuture,不要在put/get内部做阻塞式 I/O 而不切换到异步,否则可能耗尽 Broker 的事件循环线程。
  • getAll是兼容性检查的数据源:Pulsar 的 schema 兼容性校验依赖某 key 下全部历史 schema,getAll返回的“版本索引 + 逐版本 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 持久化落到任意你已有的分布式存储系统上。

  • 消息队列
  • 后端
  • 流处理

【免费下载链接】pulsar

Apache Pulsar - distributed pub-sub messaging system

项目地址:https://gitcode.com/gh_mirrors/pulsar28/pulsar
点击查看免费下载
上一篇:如何快速入门大语言模型评估:HuggingFace evaluation-guidebook新手必读
下一篇:Vial-QMK 项目常见问题解决方案

创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考

需要专业的网站建设服务?

联系我们获取免费的网站建设咨询和方案报价,让我们帮助您实现业务目标

立即咨询