- 消息队列
- 后端
- 流处理
【免费下载链接】pulsar
Apache Pulsar - distributed pub-sub messaging system
本篇技术指南基于 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 Connector | pulsar-io/aerospike | 将消息写入 Aerospike KV 数据库 |
| Cassandra Connector | pulsar-io/cassandra | 将消息写入 Apache Cassandra |
| Kafka Connector | pulsar-io/kafka | 与 Kafka topic 之间双向桥接 |
| Kinesis Connector | pulsar-io/kinesis | 对接 AWS Kinesis 流 |
| RabbitMQ Connector | pulsar-io/rabbitmq | 与 RabbitMQ 队列桥接 |
| Twitter Firehose Connector | pulsar-io/twitter | 接入 Twitter Firehose 数据源 |
从仓库结构看,除上述 6 个外,pulsar-io 还持续演进出更多连接器(如 hdfs2、elastic-search、redis、mongo、influxdb、jdbc、debezium 等),印证了发布说明中"更多连接器将在后续版本推出"的规划。
使用流程:配置 + 提交
按发布说明给出的使用方式,接入一个外部系统只需两步:
- 准备配置文件:描述要连接的外部系统(如 Cassandra 集群地址、认证信息、目标 keyspace/table 等);
- 提交连接器:通过 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 仅支持String、bytes和JSON三种 Schema,2.1 在此基础上新增了Avro与Protobuf两种主流序列化格式的支持。
从源码看 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 以"简化数据接入与流内处理"为核心主题,交出了五项重要成果:
- Pulsar IO把数据集成从"写代码"变成"写配置 + 提交";
- Tiered Storage借助 BookKeeper 的 Segment 模型实现了透明、可扩展的冷热分层;
- Stateful Function以开发者预览形式把状态能力注入 Pulsar Functions,其 State API 在当前仓库的 pulsar-functions/api-java 中已沉淀为稳定的接口抽象;
- Avro / Protobuf Schema扩展了 Pulsar 的结构化数据契约能力,实现位于 pulsar-client;
- Go Client让 Go 开发者拥有了基于 C++ 客户端的官方接入途径。
对于希望深入研究的读者,推荐从 pulsar-io 的连接器实现、tiered-storage 的 offloader 模块、AvroSchema.java 与 ProtobufSchema.java 等关键文件入手,结合本文的脉络逐一验证各特性的实际行为。
- 消息队列
- 后端
- 流处理
【免费下载链接】pulsar
Apache Pulsar - distributed pub-sub messaging system
相关推荐
Apache Pulsar PIP-312 深度解析:基于 StateStoreProvider 解耦 Pulsar Functions 状态存储与 BookKeeper
Apache Pulsar PIP 312 深度解析:基于 StateStoreProvider 解耦 Pulsar Functions 状态存储与 BookK
消息队列后端Apache Pulsar函数状态管理:基于Pulsar Table的状态持久化
Apache Pulsar函数状态管理:基于Pulsar Table的状态持久化 你是否在开发流处理应用时遇到过这些痛点?函数重启后状态丢失导致数据不一致、内存
消息队列后端流处理Apache Pulsar Functions 状态存储(State Storage)开发指南:基于 BookKeeper Table Service 的有状态函数实战
Apache Pulsar Functions 状态存储(State Storage)开发指南:基于 BookKeeper Table Service 的有状态
消息队列后端流处理
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考