news 2026/9/24 6:52:22

Apache Pulsar 2.1.0 特性深度解析:Pulsar IO、分层存储、有状态函数与 Avro/Protobuf Schema

作者头像

张小明

前端开发工程师

1.2k 24
文章封面图
Apache Pulsar 2.1.0 特性深度解析:Pulsar IO、分层存储、有状态函数与 Avro/Protobuf Schema
  • 消息队列
  • 后端
  • 流处理

【免费下载链接】pulsar

Apache Pulsar - distributed pub-sub messaging system

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

本篇技术指南基于 Apache Pulsar 官方 2.1.0-incubating 发布说明展开,系统梳理该版本引入的五大核心能力:Pulsar IO 连接器框架、基于 BookKeeper 的分层存储(Tiered Storage)、Pulsar Functions 有状态函数、Avro/Protobuf Schema 原生支持,以及全新的 Go 客户端。通过结合当前仓库(gh_mirrors/pulsar28/pulsar)中的实际源码与模块结构,读者可以理解每个特性的设计动机、使用方式与底层实现位置,从而在真实业务中正确地引入与使用这些能力。

版本背景:从 2.0 到 2.1 的演进脉络

Apache Pulsar 2.1.0-incubating 是继 2.0 之后的重要里程碑,凝聚了约两个月的开发成果。2.0 版本确立了 Pulsar 的多租户架构、Segment 分段存储模型以及原生 Schema 支持;2.1 则在此基础之上,围绕"简化流处理与数据集成"这一主线,补齐了以下能力:

  • Pulsar IO:面向进出 Pulsar 数据流连接器框架,内置 6 个开箱即用的连接器;
  • Tiered Storage:把较旧的 topic 分段卸载到长期冷存储,让 topic 变成"无限"数据流;
  • Stateful Function:为 Pulsar Functions 提供状态管理 API(开发者预览特性);
  • Avro / Protobuf Schema:在 2.0 已有的 String、bytes、JSON 之外,新增两种主流结构化数据格式的原生 Schema;
  • Go Client:基于 C++ 客户端库的新语言客户端。

以下各节将逐项展开说明,并引用当前仓库中的实现文件作为佐证。

Pulsar IO:零代码接入外部数据系统的连接器框架

设计理念:延续 Pulsar Functions 的"极简优先"

自 2.0 引入的 Pulsar Functions 是一种受 serverless 启发的轻量级流内计算框架,开发者可以用最少的样板代码实现任意复杂度的流内处理逻辑。2.1 将这一"极简优先"原则延续到了数据集成领域:开发者不需要编写任何一行连接器代码,只需要准备一份描述目标系统连接信息的配置文件,再通过 Pulsar admin CLI 提交连接器,Pulsar 便会自动接管容错、负载均衡等底层事务。

2.1 内置的 6 个连接器

2.1 版本随包发布了 6 个内置连接器,在 pulsar-io 模块目录下均可找到对应的实现模块:

连接器仓库模块说明
Aerospike Connectorpulsar-io/aerospike将消息写入 Aerospike KV 数据库
Cassandra Connectorpulsar-io/cassandra将消息写入 Apache Cassandra
Kafka Connectorpulsar-io/kafka与 Kafka topic 之间双向桥接
Kinesis Connectorpulsar-io/kinesis对接 AWS Kinesis 流
RabbitMQ Connectorpulsar-io/rabbitmq与 RabbitMQ 队列桥接
Twitter Firehose Connectorpulsar-io/twitter接入 Twitter Firehose 数据源

从仓库结构看,除上述 6 个外,pulsar-io 还持续演进出更多连接器(如 hdfs2、elastic-search、redis、mongo、influxdb、jdbc、debezium 等),印证了发布说明中"更多连接器将在后续版本推出"的规划。

使用流程:配置 + 提交

按发布说明给出的使用方式,接入一个外部系统只需两步:

  1. 准备配置文件:描述要连接的外部系统(如 Cassandra 集群地址、认证信息、目标 keyspace/table 等);
  2. 提交连接器:通过 Pulsar admin CLI 将连接器提交到 Pulsar 集群,由 Pulsar 负责后续的调度、容错与负载均衡。

官方提供的快速入门教程以连接 Apache Cassandra 为例演示完整流程。对于想贡献自定义连接器的开发者,编写一个连接器的复杂度与编写一个 Pulsar Function 相当,这也是该框架的核心卖点之一——底层复用 Pulsar Functions 的运行时能力。

Tiered Storage:把 topic 变成"无限"数据流

为什么需要分层存储

Apache Pulsar 的核心优势之一是其基于 Apache BookKeeper 的 Segment 分段存储架构:topic backlog 可以按需增长,集群空间不足时只需追加存储节点,系统会自动接管新节点而无需对已有分区做 rebalancing。然而,随着数据规模持续增长,长期保留全部热数据在 BookKeeper 中的成本会越来越高。

分层存储正是为了解决这一"成本 vs. 容量"的权衡而设计的:它将较早的 Segment 从 BookKeeper 卸载到面向冷数据设计的长期存储(如 AWS S3),在 2.1 版本中首先支持 S3,后续版本再陆续补齐 GCS、Azure Blobstore、HDFS 等 offloader(当前仓库的 tiered-storage 目录下已能看到 file-system 与 jcloud 等模块的实现雏形,其中 jcloud 提供了对接 S3、GCS、Azure 等对象存储的统一能力)。

对上层应用的透明性

分层存储对终端用户完全透明:消费者读取数据时,无论是数据仍然位于 BookKeeper 中,还是已被卸载到长期存储,体验上没有可感知的差异。所有底层的卸载机制与元数据管理都由 Pulsar 内部完成,应用无需感知数据物理存储位置的变化。

这意味着开发者可以放心地把 topic 当作真正的"无限流"来使用:不需要预先规划存储上限,历史数据自动分层归档,新数据始终享受 BookKeeper 的低延迟写入与读取。

Stateful Function:为 Pulsar Functions 引入状态管理

状态:流处理引擎的最大挑战

状态管理是流处理引擎面临的最大挑战,Pulsar Functions 也不例外。Pulsar Functions 的目标是简化流式处理逻辑的开发,因此为函数提供易用的状态管理 API 成为自然延伸。

State API 与 BookKeeper Table Service 集成

2.1 为 Pulsar Functions Java SDK 引入了一套 State API,用于持久化函数状态。该 API 与 Apache BookKeeper 中的 Table Service 集成,状态存储由 BookKeeper 负责。该特性在 2.1 中以**开发者预览(developer preview)**形态发布,官方希望收集社区反馈以在后续版本中持续改进。

从当前仓库的 pulsar-functions/api-java 可以看到这套状态抽象已经沉淀为标准接口:

  • StateStore.java:函数状态存储的顶层接口,函数通过Context按名称访问对应的 StateStore;
  • CounterStateStore.java:内置的分布式计数器能力,提供incrCounter(key, amount)/getCounter(key)同步方法及对应的incrCounterAsync/getCounterAsync异步方法,适用于词频统计、事件计数等典型场景;
  • ByteBufferStateStore.java:以字节缓冲为载体的键值状态读写接口。

借助这些 API,函数可以在多次调用之间保持并累积状态(例如统计某个窗口内出现的单词次数),而无需自行对接外部存储。

Schemas:Avro 与 Protobuf 原生支持

2.0 的 Schema 基础

Pulsar 2.0 引入了 Schema 原生支持:开发者可以声明消息数据的结构,由 Pulsar 强制校验——只有符合声明结构的生产者才能向对应 topic 发布合法数据。2.0 仅支持StringbytesJSON三种 Schema,2.1 在此基础上新增了AvroProtobuf两种主流序列化格式的支持。

从源码看 AvroSchema 的实现

当前仓库中,Avro Schema 的实现位于 AvroSchema.java。从源码可以看出几个关键设计:

  • 继承自AvroBaseStructSchema,通过AvroReader/AvroWriter完成序列化与反序列化;
  • of(Class<T> pojo)系列静态工厂方法让开发者可以用一个 POJO 直接构造 Schema;
  • 支持supportSchemaVersioning()返回true,即配合MultiVersionAvroReader支持按 Schema 版本解码历史消息,这是结构化 Schema 相比原始字节流的显著优势;
  • 内置了 Avro Logical Type 的转换支持(如 decimal、date、time、timestamp、uuid 等),并可在jsr310ConversionEnabled开关下在 Joda-Time 与 Java 8 时间类型之间切换。

ProtobufSchema 与 JSONSchema 的对照

  • ProtobufSchema.java 面向com.google.protobuf.GeneratedMessageV3的 protobuf 消息类:通过ProtobufData.get().getSchema(pojo)把 protobuf descriptor 转换为 Avro Schema 表示并注册到SchemaInfo中,同时把字段的解析信息(字段号、名称、类型、label)序列化为属性__PARSING_INFO__随 Schema 一起发布,方便消费者侧还原消息结构;
  • JSONSchema.java 则展示了 2.0 时代 JSON Schema 的延续:基于 Jackson 实现读写,并保留了向后兼容的 JSON Schema(非 Avro)生成逻辑。

使用价值

Schema 的意义在于把"数据结构契约"纳入消息系统管理:生产者只能发布符合声明的数据,消费者可以按 Schema 安全解码,topic 的历史 Schema 版本被系统记录。在 2.1 中引入 Avro/Protobuf 后,基于强类型的跨语言数据交换(例如 Java 生产者发布 Avro 消息、Go 消费者消费)成为可能,这也与同一版本推出的 Go 客户端形成配合。

Go Client:新语言客户端

2.1 引入了全新的 Go Client,这是 Pulsar 官方客户端家族的新成员。值得说明的是,Go 客户端库基于 C++ 客户端库构建(通过 CGO 绑定 pulsar-client-cpp 实现),因此两者共享底层的协议实现与行为语义。开发者可以按照官方安装指引在自己的 Go 应用中引入并使用该客户端进行生产/消费。

需要补充的是,当前仓库主线聚焦于 Java 生态与 C++ 客户端(pulsar-client-cpp),Go 客户端的独立代码库自 2.1 之后单独演进;从仓库结构看,pulsar-client-cpp 中 Python 绑定(python/pulsar)与 C API 的存在,也印证了 C++ 客户端作为多语言客户端共同底层的事实。

总结

Apache Pulsar 2.1.0-incubating 以"简化数据接入与流内处理"为核心主题,交出了五项重要成果:

  1. Pulsar IO把数据集成从"写代码"变成"写配置 + 提交";
  2. Tiered Storage借助 BookKeeper 的 Segment 模型实现了透明、可扩展的冷热分层;
  3. Stateful Function以开发者预览形式把状态能力注入 Pulsar Functions,其 State API 在当前仓库的 pulsar-functions/api-java 中已沉淀为稳定的接口抽象;
  4. Avro / Protobuf Schema扩展了 Pulsar 的结构化数据契约能力,实现位于 pulsar-client;
  5. Go Client让 Go 开发者拥有了基于 C++ 客户端的官方接入途径。

对于希望深入研究的读者,推荐从 pulsar-io 的连接器实现、tiered-storage 的 offloader 模块、AvroSchema.java 与 ProtobufSchema.java 等关键文件入手,结合本文的脉络逐一验证各特性的实际行为。

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

【免费下载链接】pulsar

Apache Pulsar - distributed pub-sub messaging system

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

相关推荐

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

版权声明: 本文来自互联网用户投稿,该文观点仅代表作者本人,不代表本站立场。本站仅提供信息存储空间服务,不拥有所有权,不承担相关法律责任。如若内容造成侵权/违法违规/事实不符,请联系邮箱:809451989@qq.com进行投诉反馈,一经查实,立即删除!
网站建设 2026/9/24 6:50:38

TMS320F280049开发选型指南:C2000Ware与MotorControl SDK深度对比

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

作者头像 李华
网站建设 2026/9/24 6:40:13

用 AI 处理敏感数据前,先分清这三层的边界

问题的本质 「AI 会不会泄露我的数据」这个问题问得太笼统。把它拆成三层&#xff0c;答案就清楚了&#xff1a;数据在哪一层&#xff0c;决定了它有没有出网。 第一层&#xff1a;模型层&#xff08;生成建议&#xff09; 你把数据贴进对话框&#xff0c;让 AI 帮你写方法、…

作者头像 李华
网站建设 2026/9/24 6:36:24

操作系统期末复习:PV操作、银行家算法与页面置换高频考点详解

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

作者头像 李华
网站建设 2026/9/24 6:21:10

STM32G4 ADC硬件过采样与软件滤波实战:从配置到代码落地

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

作者头像 李华