☰
Apache Pulsar Schema Registry 概念与实践指南:从类型安全到版本化 Schema 管理
2026/9/25 21:03:29 网站建设 项目流程
  • 消息队列
  • 后端
  • 流处理

【免费下载链接】pulsar

Apache Pulsar - distributed pub-sub messaging system

项目地址:https://gitcode.com/gh_mirrors/pulsar28/pulsar
点击查看免费下载

导读

本文聚焦 Apache Pulsar 内置的Schema Registry(Schema 注册中心),系统讲解它在 Pulsar 消息系统中解决“类型安全(Type Safety)”问题的核心机制:Schema 如何在 topic 级别上传、存储、校验与版本化,客户端如何基于 Schema 进行序列化/反序列化,以及如何通过 REST API 与pulsar-admin命令行工具管理 Schema。读完本文,你将掌握 Pulsar Schema Registry 的架构与工作流程、受支持的全部 Schema 格式、Schema 版本演进规则,以及从 Java 客户端到管理命令的完整实战用法。

适用版本说明:本文以当前仓库site2/website-next/versioned_docs/version-2.2.0/下的 2.2.0 版本文档为主体,并结合仓库源码中 Schema Registry 相关实现进行深化,文中涉及的 API 与命令行为以仓库当前实际实现为准。

为什么需要 Schema Registry:消息系统中的类型安全

在围绕 Pulsar 这类消息总线构建的任何应用里,类型安全都极其重要。Producer 与 Consumer 需要在 topic 级别协调数据类型,否则会引发一系列潜在问题——最典型的就是序列化/反序列化错误。消息在 Pulsar 中本质上是原始字节流,如果 Producer 在topic-1上发送温度传感器数据,而该 topic 的 Consumer 却试图把数据解析成湿度传感器读数,就会立刻出错。

围绕消息的类型安全,应用通常采用两种基本思路:

  1. “客户端侧(client-side)”方案:Producer 和 Consumer 不仅要负责消息(原始字节)的序列化与反序列化,还要自行“知晓”哪个 topic 传输哪种类型。这种方案把所有类型安全的维护负担都交给了应用层,即所谓的 "out-of-band"(带外)管理。
  2. “服务端侧(server-side)”方案:Producer 和 Consumer 主动告知系统某个 topic 可以传输哪些数据类型。消息系统负责强制类型安全,确保 Producer 与 Consumer 始终保持同步。

Pulsar 对两种方案都支持,你可以自由选择其中一种,也可以在逐 topic 粒度上混用:

  • 采用“客户端侧”方案时,Producer 与 Consumer 可以发送/接收由原始字节数组构成的消息,把全部类型安全交给应用在带外处理;
  • 采用“服务端侧”方案时,Pulsar 内置的Schema Registry允许客户端按 topic 上传数据 Schema,这些 Schema 决定了该 topic 上哪些数据类型被视为合法。

注意:在 2.2.0 版本中,Pulsar Schema Registry 仅对 Java 客户端、CGo 客户端、Python 客户端 和 C++ 客户端 可用。

基本架构:Schema 如何进入注册中心

从架构上看,Pulsar Schema Registry 的写入与读取遵循以下路径:

  • 当你使用 Schema 创建带类型的 Producer时,Schema 会被自动上传;
  • 你也可以通过 Pulsar 的 REST API(/admin/v2/schemas系列端点)手动上传、获取和更新Schema。

在存储层面,Pulsar 开箱即用地使用Apache BookKeeper日志存储系统作为 Schema 的持久化后端(对应 Pulsar 架构中的 persistent storage 层)。如果你愿意,也可以接入其他存储后端——2.2.0 文档说明自定义 Schema 存储逻辑的文档“即将推出”,因此本文以默认的 BookKeeper 后端为准。

从源码看,这一架构在 broker 侧由SchemaRegistryService接口统一抽象,其实现通过工厂方法创建:SchemaRegistryService.create(SchemaStorage, Set<String> checkerClasses)会先构建SchemaType -> SchemaCompatibilityCheck的兼容性检查器映射(并为KEY_VALUE类型装配KeyValueSchemaCompatibilityCheck),随后用SchemaRegistryServiceWithSchemaDataValidator包装真实的SchemaRegistryServiceImpl;当 Schema 存储创建失败时,则退化为DefaultSchemaRegistryService(空实现),相关逻辑见 SchemaRegistryService.java。

Schema 如何工作:数据模型与应用范围

Pulsar Schema 是相当简单的数据结构,由以下部分组成:

组成部分说明
Name(名称)在 Pulsar 中,Schema 的名称就是它被应用的那个 topic
Payload(负载)Schema 的二进制表示
Type(类型)Schema 的格式类型,见下文“支持的 Schema 格式”
Properties(属性)用户自定义的 string/string 映射。其用法完全取决于具体应用,常见的例子包括:关联的 Git hash、环境标识(如dev、prod)等

两个关键约束必须牢记:

  • Schema 只作用于 topic 级别,不能应用到 namespace 或 tenant 级别;
  • Producer 和 Consumer 都是把 Schema上传给 Pulsar broker,由 broker 统一存储与校验。

这个“topic 级”数据模型在 broker 实现中体现得非常直接:SchemaRegistryServiceImpl的所有读写方法都以schemaId为键(而 schemaId 正是 topic 名),putSchemaIfAbsent、getSchema、getAllSchemas、deleteSchema等操作全部围绕schemaId展开,参见 SchemaRegistryServiceImpl.java。

Schema 版本机制:三种连接场景逐项剖析

Schema 版本化是 Pulsar Schema Registry 最核心的机制。为了说明其工作原理,先看一个 Java 客户端创建带 Schema 的 Producer 的示例:

PulsarClient client = PulsarClient.builder() .serviceUrl("pulsar://localhost:6650") .build(); Producer<SensorReading> producer = client.newProducer(JSONSchema.of(SensorReading.class)) .topic("sensor-data") .sendTimeout(3, TimeUnit.SECONDS) .create();

当这个 Producer 尝试连接 broker 时,可能出现以下三种场景及其对应处理:

场景会发生什么
topic 上不存在任何 SchemaProducer 以给定 Schema 创建。Schema 被传送到 broker 并存储(因为没有已有 Schema 与SensorReading兼容)。任何使用相同 Schema/topic 创建的 Consumer 都可以消费sensor-datatopic 上的消息
已存在Schema,Producer 使用相同Schema 连接Schema 被传送到 broker。broker 判定该 Schema 是兼容的,尝试将其存储到 BookKeeper,但随后发现它已经存储过,于是该 Schema 被用来为生产的消息打上标签(version)
已存在Schema,Producer 使用新的、兼容的Schema 连接Producer 将 Schema 传送给 broker。broker 判定其兼容,将新 Schema 存储为当前版本(获得新的版本号)

版本分配规则:Schema 按先后顺序进行版本化。Schema 的存储发生在处理对应 topic 的 broker 上,以便分配版本号。一旦某个 Schema 被分配/获取了版本,该 Producer 后续生产的所有消息都会被标记上相应的版本号。

源码层面印证了这套流程:putSchemaIfAbsent会先查询已存在的 Schema 列表,若发现相同 Schema(通过 SHA-256 哈希比对,见SchemaRegistryServiceImpl中的hashFunction = Hashing.sha256()与SchemaHash.of(...)),则直接返回已有版本;否则按兼容性策略校验,通过后把新 Schema 连同类型、用户、时间戳、属性等写入存储并生成新版本。兼容性校验失败的路径会抛出IncompatibleSchemaException(例如“已存在 schema 类型 X,新 schema 类型 Y”),参见 SchemaRegistryServiceImpl.java。

支持的 Schema 格式

Pulsar Schema Registry 支持以下格式:

Schema 类型说明
None若 topic 未指定 Schema,Producer 和 Consumer 直接处理原始字节
String用于 UTF-8 编码的字符串
JSONJSON 对象编码与校验
ProtobufProtocol Buffers 消息编码与解码
Avro通过 Avro 进行序列化/反序列化

对应地,客户端 API 的SchemaType枚举定义了这些类型的序号(NONE=0, STRING=1, JSON=2, PROTOBUF=3, AVRO=4),并注释说明新增需要记录进 Schema Registry 的类型时,应同步修改pulsar-common/src/main/proto/PulsarApi.proto与pulsar-broker/src/main/proto/SchemaRegistryFormat.proto两个 proto 文件,见 SchemaType.java。

2.2.0 文档同时指出:其他 Schema 格式的支持将在未来版本中陆续加入。从当前仓库的SchemaType枚举看,后续版本已扩展了BOOLEAN、INT8/16/32/64、FLOAT、DOUBLE、DATE、TIME、TIMESTAMP、BYTES等更多内建类型(标记为@since 2.3.0等),这印证了文档的“未来扩展”方向。

实战示例:用 RecordSchemaBuilder 构建 Avro Schema 并消费 GenericRecord

下面的例子演示如何用RecordSchemaBuilder定义 Avro Schema、用GenericRecordBuilder生成通用 Avro 记录,并把消息消费为GenericRecord。

第 1 步:用RecordSchemaBuilder构建 Schema

RecordSchemaBuilder recordSchemaBuilder = SchemaBuilder.record("schemaName"); recordSchemaBuilder.field("intField").type(SchemaType.INT32); SchemaInfo schemaInfo = recordSchemaBuilder.build(SchemaType.AVRO); Producer<GenericRecord> producer = client.newProducer(Schema.generic(schemaInfo)).create();

第 2 步:用GenericRecordBuilder构建并发送通用记录

producer.newMessage().value(schema.newRecordBuilder() .set("intField", 32) .build()).send();

这段代码中涉及的核心 API 都定义在 pulsar-client-api 模块下:RecordSchemaBuilder用于按字段描述 Schema(支持field(name).type(SchemaType)链式声明),SchemaBuilder.record(name)创建记录型 Schema 构建器,Schema.generic(schemaInfo)生成可处理GenericRecord的通用 Schema,GenericRecordBuilder则用于按字段名设置值并产出GenericRecord。这是“带类型 Producer + 服务端 Schema 校验”组合的典型用法:Producer 端定义并注册 Schema,broker 端依据该 Schema 校验后续消息。

管理 Schema:REST API 与 pulsar-admin 命令

你可以使用 Pulsar 管理工具(pulsar-admin)或 REST API 管理 topic 的 Schema。

REST API 端点

Broker 侧的v2管理端点定义在 SchemasResource.java(基类实现见 SchemasResourceBase.java),路径前缀为/schemas:

方法路径功能
GET/schemas/{tenant}/{namespace}/{topic}/schema获取 topic 的最新 Schema
GET/schemas/{tenant}/{namespace}/{topic}/schema/{version}获取指定版本的 Schema
GET/schemas/{tenant}/{namespace}/{topic}/schemas获取 topic 的全部(各版本)Schema
POST/schemas/{tenant}/{namespace}/{topic}/schema上传/更新 Schema(请求体为PostSchemaPayload)
DELETE/schemas/{tenant}/{namespace}/{topic}/schema删除 topic 的最新 Schema

pulsar-admin schemas 子命令

命令行入口位于 CmdSchemas.java,命令组为pulsar-admin schemas,支持以下四个子命令:

1.get:获取 topic 的 Schema

pulsar-admin schemas get persistent://tenant/namespace/topic
  • 参数:persistent://tenant/namespace/topic(必填)
  • 可选参数:
    • -v, --version:指定版本号,必须大于 0;
    • -a, --all-version:列出全部版本;
    • --version与--all-version不能同时指定。
  • 不带任何可选参数时,默认输出最新 Schema 及其版本信息。

2.delete:删除 topic 的最新 Schema

pulsar-admin schemas delete persistent://tenant/namespace/topic

对应getAdmin().schemas().deleteSchema(topic),删除操作在 broker 端会写入一条标记deleted=true的 Schema 记录,因此删除后该 topic 的旧版本 Schema 无法再被获取(参见SchemaRegistryServiceImpl.getSchema对isDeleted()的过滤逻辑)。

3.upload:为 topic 上传/更新 Schema

pulsar-admin schemas upload persistent://tenant/namespace/topic -f /path/to/schema.json
  • -f, --filename(必填):包含PostSchemaPayload(即 type、schema、properties 等字段)的 JSON 文件路径;
  • 内部实现是把文件内容解析为PostSchemaPayload后调用createSchema(topic, input)。

4.extract:从 JAR 中提取 POJO 的 Schema 并上传

pulsar-admin schemas extract persistent://tenant/namespace/topic \ -j /path/to/pojo.jar -t avro -c com.example.SensorReading
  • -j, --jar(必填):包含 POJO 类的 JAR 文件路径;
  • -t, --type(必填):提取类型,仅支持avro或json;
  • -c, --classname(必填):POJO 类全限定名;
  • --always-allow-null:设置 Schema 是否始终允许 null(默认true);
  • -n, --dry-run:只打印将上传的PostSchemaPayload(JSON 形式),不真正应用到 Schema Registry。

该命令通过URLClassLoader加载 JAR 中的 POJO,再借助SchemaExtractor生成 Avro/JSON Schema,是“从现有 Java 类快速注册 Schema”的实用工具。

与消息生产/消费流程的联动

理解 Schema Registry 后,还需要把它放回消息生命周期中看待。从 broker 侧源码看,Schema 的注册与校验深度嵌入在消息处理链路中:

  • 生产者(Producer)侧:ServerCnx在处理 producer 创建请求时调用 SchemaRegistryService 完成 Schema 注册与兼容性检查;
  • 消费者(Consumer)侧:checkConsumerCompatibility依据兼容性策略校验消费者携带的 Schema 与已存 Schema 是否匹配(ALWAYS_COMPATIBLE策略直接放行,其余策略按 BACKWARD/FORWARD/FULL 等规则与最新版或全部历史版本比对),见 SchemaRegistryServiceImpl.java。

这样,Producer 上传并版本化 Schema、消息被打上 Schema 版本标签、Consumer 按兼容策略校验,三者构成闭环:类型安全由 broker 强制执行,Producer 与 Consumer 无需在带外自行约定类型——这正是本文开头“服务端侧方案”的完整落地。

小结

Pulsar Schema Registry 为消息总线场景下的类型安全提供了一条服务端强制的路径:Schema 在 topic 级别注册,由 broker 统一存储于 BookKeeper 并按顺序版本化;Producer 上传 Schema 后,其生产的消息被标记上对应版本,Consumer 通过兼容性策略与已存 Schema 对齐。配合pulsar-admin schemas命令与 REST API,你可以完整地管理 Schema 的生命周期(查询、上传、更新、删除)。在实际项目中,可结合本文示例从 Java 客户端侧开启 Schema(如JSONSchema.of(Class)或RecordSchemaBuilder+Schema.generic(...)),并在后续演进 Schema 时充分理解版本兼容规则,避免不兼容变更导致生产链路中断。

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

【免费下载链接】pulsar

Apache Pulsar - distributed pub-sub messaging system

项目地址:https://gitcode.com/gh_mirrors/pulsar28/pulsar
点击查看免费下载

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

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

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

立即咨询