Apache Pulsar 自定义 Schema 存储:实现 SchemaStorage 与 SchemaStorageFactory 接口指南
消息队列后端流处理【免费下载链接】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),仅供参考

相关新闻

从华为杯一等奖复盘看数学建模竞赛的流程管理与决策智慧

从华为杯一等奖复盘看数学建模竞赛的流程管理与决策智慧

/* MD / 富文本中的 .toc(含博客园搬家等嵌套结构);.toc-box 在侧栏,不受影响 */#content_views .toc,/* 编辑器常在目录前后插入空 p(:empty 仍占 20px),一并去掉避免顶空隙 */#content_views.markdown_views > p:empty:has(+ .toc),#content_views.markdown_views …

2026/9/25 6:53:20 阅读更多 →
【Coze】在Coze平台使用源码创建工作流

【Coze】在Coze平台使用源码创建工作流

Coze 提供了图形化的工作流搭建平台,适用于低代码构建自动化任务流程。通过资源管理、节点配置与流程连接,可实现多种业务逻辑的在线部署。 本文介绍如何在 Coze 中创建工作流资源、导入流程 JSON 配置,并完成起止节点的连接与字段设置,直至试运行与发布上线的全过程。 文…

2026/9/25 6:52:20 阅读更多 →
深度解析 Hypothesis 测试执行次数:`max_examples` 的完整运行语义与底层实现

深度解析 Hypothesis 测试执行次数:`max_examples` 的完整运行语义与底层实现

测试开发工具 【免费下载链接】hypothesis The property-based testing library for Python 项目地址: https://gitcode.com/gh_mirrors/hy/hypothesis 点击查看 免费下载 本指南聚焦 Hypothesis(Python 属性测试库)中一个看似简单实则微妙的…

2026/9/25 6:52:20 阅读更多 →

最新新闻

Hugo Blox Résumé 模板实战:用 blocks 组装在线简历 landing 页的完整配置指南

Hugo Blox Résumé 模板实战:用 blocks 组装在线简历 landing 页的完整配置指南

静态站点前端开发工具 【免费下载链接】kit 🧱 Describe your site, AI builds it, you own it as Markdown. Snap together Tailwind blocks like Lego — landing pages, blogs, portfolios, docs & more. No AI slop. Free to deploy anywhere 👇…

2026/9/25 8:53:13 阅读更多 →
有名的奢侈品名表回收品牌企业、服务不错的奢侈品名表回收企业、有名的奢侈品名表回收专业公司用户力荐

有名的奢侈品名表回收品牌企业、服务不错的奢侈品名表回收企业、有名的奢侈品名表回收专业公司用户力荐

有名专业的奢侈品名表回收,靠谱连锁品牌更安心很多想要出手闲置奢侈品名表的用户,都希望找到透明靠谱的专业平台,东莞市好岱贸易有限公司旗下品牌好岱中古汇,是一家深耕二手奢侈品回收行业的全国连锁直营平台,始终坚持…

2026/9/25 8:53:12 阅读更多 →
Sinon sandbox.replace() 完全指南:安全替换对象属性并自动还原

Sinon sandbox.replace() 完全指南:安全替换对象属性并自动还原

测试开发工具 【免费下载链接】sinon Test spies, stubs and mocks for JavaScript. 项目地址: https://gitcode.com/gh_mirrors/si/sinon 点击查看 免费下载 Sinon 的 sandbox.replace() 用于在测试中临时替换对象上的任意属性(方法、字符串、数值乃至…

2026/9/25 8:53:12 阅读更多 →
华为FTTR全光家庭网:光纤入室解决WiFi6卡顿与多终端抢带宽

华为FTTR全光家庭网:光纤入室解决WiFi6卡顿与多终端抢带宽

简介:本资源为华为FTTR全光家庭网络创新解决方案的完整技术白皮书PDF,面向通信工程师、宽带网络规划人员、智慧家庭方案集成商及运营商装维技术人员,聚焦解决大户型Wi-Fi覆盖弱、千兆宽带实际速率不足、多终端卡顿掉线等家庭网络核心痛点。文…

2026/9/25 8:53:12 阅读更多 →
河北欧米奇西点西餐学校行业口碑如何

河北欧米奇西点西餐学校行业口碑如何

核心定位河北欧米奇西点西餐学校是中国东方教育集团旗下经石家庄市批准设立的正规西式餐饮职业学校,作为河北省西式餐饮专业人才培养重点基地,核心面向15-18岁青少年提供技能学历双提升的西式餐饮技能培训服务,深耕行业三十余年,已…

2026/9/25 8:53:12 阅读更多 →
制造经理的高效:从救火队长到系统管理员

制造经理的高效:从救火队长到系统管理员

做了十几年制造经理,我对“高效”这个词越来越警惕。面试的时候老板爱问“你怎么理解效率”,很多人脱口就是“结果导向、执行力强、按计划走”,听着都对,但放到现场基本没用。真正的高效不是你自己一天处理了多少件事,…

2026/9/25 8:52:12 阅读更多 →

日新闻

AI元人文:从工具使用到思维重构的深度探索

AI元人文:从工具使用到思维重构的深度探索

最近半年我一直在琢磨一件事:AI元人文到底是什么?说白了,就是“用元视角重新审视人与AI的关系”,也在“探索AI如何反向逼着我们发现自己的思考边界”。标题里的“元探索”,在我看就是一层套一层的追问——当你用AI解决…

2026/9/25 0:00:41 阅读更多 →
Python+CNN车牌识别实战:从数据预处理到模型训练与部署

Python+CNN车牌识别实战:从数据预处理到模型训练与部署

简介:基于Python与卷积神经网络的车牌识别项目,面向计算机视觉初学者及智能交通开发者,目标是帮助用户掌握从数据预处理、模型构建到实际部署的完整流程。压缩包共25个文件,包含jpg/png图像样本、py训练脚本、md说明文档、dat数据…

2026/9/25 0:00:41 阅读更多 →
Vim基础操作全攻略:保存退出、模式切换与高频命令实战

Vim基础操作全攻略:保存退出、模式切换与高频命令实战

1. 项目概述1.1 核心需求解析今天聊聊Vim。写这个题目的原因是:几乎每个后端开发者、运维人员、数据工程师某天都会遇到一个场景——深夜加班,服务器登录界面只有黑底白字,编辑器只有vi/vim,你必须在五分钟内完成一次配置修改并保…

2026/9/25 0:00:41 阅读更多 →

周新闻

Flutter for OpenHarmony游戏卡片渐变背景实战:从原理到性能优化

Flutter for OpenHarmony游戏卡片渐变背景实战:从原理到性能优化

直接铺开项目本身吧。这几个月我一直在折腾一件事:用Flutter给OpenHarmony做一款游戏集合类的App,说白了就是把若干小游戏塞进一个壳里,用统一入口分发。这个方向本身不算新鲜,真正让我花了不少心思的,是首页那堆游戏卡…

2026/9/24 14:34:13 阅读更多 →
Word表格编号全攻略:从列表编号到题注交叉引用

Word表格编号全攻略:从列表编号到题注交叉引用

写Word文档,最让人头疼的往往是那些“看起来不起眼”的小问题。比如表格编号这事:今天在表后面多加了两个空白行,明天给客户交稿前发现整个章节的编号全部错位,光是挨个改序号就能耗掉大半个下午。我前阵子帮人整理一份上百页的技…

2026/9/24 9:10:42 阅读更多 →
从第一个站到第二个站:独立开发者的静态网站选型与落地实践

从第一个站到第二个站:独立开发者的静态网站选型与落地实践

1. 项目概述1.1 核心需求解析做独立开发者这几年,说实话,第一个网站上线的那天晚上我兴奋得没睡着。但等它跑了半年,流量惨淡、功能臃肿、代码自己都懒得看第二遍之后,我才慢慢琢磨明白一个道理:第一个网站是练手&…

2026/9/24 14:33:56 阅读更多 →

月新闻

持续集成 流水线自动化与 声明式交付 实践:原型怎样变成可用功能

持续集成 流水线自动化与 声明式交付 实践:原型怎样变成可用功能

持续集成 流水线自动化与 声明式交付 实践:原型怎样变成可用功能分类:[AI/大模型]细分主题:AI 增强型 CI/CD 流水线自动化与 GitOps 实践:Agent 工作流、工具调用与任务拆解:从原型到生产的验收清单很多团队在尝试用大…

2026/9/24 12:50:34 阅读更多 →
容器编排 生产环境运维与排障实战:复盘记录怎样真正派上用场

容器编排 生产环境运维与排障实战:复盘记录怎样真正派上用场

容器编排 生产环境运维与排障实战:复盘记录怎样真正派上用场分类:[工程技术]细分主题:Kubernetes 生产环境运维与排障实战:可复制的项目复盘模板与决策记录大部分团队的事故复盘报告,最后都变成了躺在 Confluence 或钉…

2026/9/24 14:33:48 阅读更多 →
容器 容器化技术与镜像安全管理:核心链路应该先拆哪一步

容器 容器化技术与镜像安全管理:核心链路应该先拆哪一步

容器 容器化技术与镜像安全管理:核心链路应该先拆哪一步分类:[工程技术]细分主题:Docker 容器化技术与镜像安全管理:核心链路的逐步实现与关键代码取舍面对一个积累了五六年历史包袱的单体架构应用(包含 Web 接口、后台…

2026/9/24 12:49:17 阅读更多 →