ARTICLE DETAIL

资讯详情

深耕网站视觉设计与运营推广的一线实战洞察。

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

Apache Pulsar 2.1.0 特性深度解析:Pulsar IO、分层存储、有状态函数与 Avro/Protobuf Schema 消息队列后端流处理【免费下载链接】pulsarApache 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 之外新增两种主流结构化数据格式的原生 SchemaGo 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 CassandraKafka 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 等印证了发布说明中更多连接器将在后续版本推出的规划。使用流程配置 提交按发布说明给出的使用方式接入一个外部系统只需两步准备配置文件描述要连接的外部系统如 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按名称访问对应的 StateStoreCounterStateStore.java内置的分布式计数器能力提供incrCounter(key, amount)/getCounter(key)同步方法及对应的incrCounterAsync/getCounterAsync异步方法适用于词频统计、事件计数等典型场景ByteBufferStateStore.java以字节缓冲为载体的键值状态读写接口。借助这些 API函数可以在多次调用之间保持并累积状态例如统计某个窗口内出现的单词次数而无需自行对接外部存储。SchemasAvro 与 Protobuf 原生支持2.0 的 Schema 基础Pulsar 2.0 引入了 Schema 原生支持开发者可以声明消息数据的结构由 Pulsar 强制校验——只有符合声明结构的生产者才能向对应 topic 发布合法数据。2.0 仅支持String、bytes和JSON三种 Schema2.1 在此基础上新增了Avro与Protobuf两种主流序列化格式的支持。从源码看 AvroSchema 的实现当前仓库中Avro Schema 的实现位于 AvroSchema.java。从源码可以看出几个关键设计继承自AvroBaseStructSchema通过AvroReader/AvroWriter完成序列化与反序列化of(ClassT 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-cppGo 客户端的独立代码库自 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-clientGo Client让 Go 开发者拥有了基于 C 客户端的官方接入途径。对于希望深入研究的读者推荐从 pulsar-io 的连接器实现、tiered-storage 的 offloader 模块、AvroSchema.java 与 ProtobufSchema.java 等关键文件入手结合本文的脉络逐一验证各特性的实际行为。赞分享消息队列后端流处理【免费下载链接】pulsarApache Pulsar - distributed pub-sub messaging system项目地址https://gitcode.com/gh_mirrors/pulsar28/pulsar点击查看免费下载相关推荐Apache Pulsar PIP-312 深度解析基于 StateStoreProvider 解耦 Pulsar Functions 状态存储与 BookKeeperApache 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),仅供参考
返回列表