- 消息队列
- 后端
- 流处理
【免费下载链接】pulsar
Apache Pulsar - distributed pub-sub messaging system
导读
本文聚焦 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 却试图把数据解析成湿度传感器读数,就会立刻出错。
围绕消息的类型安全,应用通常采用两种基本思路:
- “客户端侧(client-side)”方案:Producer 和 Consumer 不仅要负责消息(原始字节)的序列化与反序列化,还要自行“知晓”哪个 topic 传输哪种类型。这种方案把所有类型安全的维护负担都交给了应用层,即所谓的 "out-of-band"(带外)管理。
- “服务端侧(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 上不存在任何 Schema | Producer 以给定 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 编码的字符串 |
| JSON | JSON 对象编码与校验 |
| Protobuf | Protocol 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
相关推荐
Apache Pulsar Schema Registry 完整指南:主题级类型安全、Schema 版本管理与多格式支持
Apache Pulsar Schema Registry 完整指南:主题级类型安全、Schema 版本管理与多格式支持 Apache Pulsar 内置的 S
消息队列后端流处理Apache Pulsar Schema 入门:理解 Schema Registry、类型安全与生产消费实战
Apache Pulsar Schema 入门:理解 Schema Registry、类型安全与生产消费实战 本指南以 Apache Pulsar 的 Sche
消息队列后端流处理Apache Pulsar Schema 入门指南:Schema Registry、类型安全与首个带 Schema 的 Java 客户端
Apache Pulsar Schema 入门指南:Schema Registry、类型安全与首个带 Schema 的 Java 客户端 Apache Puls
消息队列后端流处理
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考