
缓存KV存储消息队列流处理后端【免费下载链接】hazelcastHazelcast is a unified real-time data platform combining stream processing with a fast data store, allowing customers to act instantly on>项目地址https://gitcode.com/gh_mirrors/ha/hazelcast点击查看免费下载Apache Pulsar 连接器Contrib 模块让 Hazelcast Jet 既能把 Pulsar 主题中的消息摄入 Jet 管道又能把处理结果发布回 Pulsar 主题。本文以 009-pulsar-connector.md 设计文档为核心结合extensions/pulsar模块的源码实现深入讲解 Consumer API 源与 Reader API 源在订阅模型、接收策略、确认机制、容错语义上的取舍以及 Sink 的异步发布与重试设计并给出可直接运行的完整代码示例。读完你将掌握如何在 Jet Pipeline 中正确选择 Pulsar 源、理解 exactly-once 语义的前提条件以及如何通过 Builder 配置自定义行为。概述连接器做什么Pulsar 连接器实现两件事Source源从 Pulsar 主题读取消息进入 Jet Pipeline用于数据摄取ingestionSink汇将 Jet 管道处理后的结果发布到 Pulsar 主题。它完全基于 Pulsar 官方客户端库。Pulsar 客户端为“从主题读消息”提供了两套截然不同的 APIConsumer API基于订阅subscription的高层抽象Reader API基于MessageId的低层抽象。连接器的巧妙之处在于两种 API 的源都被实现了——Consumer 源带来分布式能力Reader 源带来 exactly-once 容错能力Sink 则基于 Producer API。三个入口类分别是PulsarSources、PulsarSinks位于 extensions/pulsar/src/main/java/com/hazelcast/jet/pulsar/自 6.0 起随 Hazelcast 发布需引入extensions/pulsar模块作为依赖。双源策略为什么同时用 Consumer 与 Reader两种 API 各有优势连接器按需各取所长维度Consumer API 源Reader API 源分区映射由 Pulsar 订阅抽象自动处理无需关注 partition 与 processor 的对应关系不支持分布式/分区主题需自定义分区映射机制当前版本缺失消息游标不向应用暴露最后消费消息的 cursor无法做精确回滚可从指定MessageId开始读取天然支持快照恢复容错语义无法保证详见下文可达到 exactly-once前提是管道其余部分同样保证消息顺序Shared 订阅轮询分发顺序无法保持单一 Reader 顺序读取下文分别剖析两套源的实现细节。Consumer API 源Consumer API 建立在主题订阅之上应用订阅主题后从该订阅的第一条未确认消息开始消费。其核心操作只有 4 个订阅主题、从未确认消息开始消费、消费消息、向 broker 发送确认ack。订阅模式为什么选 SharedPulsar 提供三种订阅模式exclusive独占、shared共享、failover故障切换。连接器的 Consumer 源选择sharedround-robin模式源码中通过subscriptionType(SubscriptionType.Shared)明确指定见 PulsarConsumerBuilder.java。Shared 模式允许多个 consumer 挂到同一订阅上broker 以轮询方式把消息分发给各个 consumer且每条消息只投递给一个 consumer当某个 consumer 断开时已投递给它但未确认的消息会被重新调度给其余 consumer。这带来两点分布式消费这正是 Jet 需要的——不同 Jet 处理器可并行消费同一主题broker 在分区与消费者之间做负载均衡无需维护 partition→processor 的一一映射顺序性丧失轮询分发破坏消息顺序因此文档明确提示对顺序敏感的场景需考虑其他订阅选项。此外Consumer 源被标记为distributed(2)PulsarConsumerBuilder.java即本地并行度偏好为 2。消息接收策略同步批量接收Consumer API 提供 4 种接收策略Blocking Single Receive同步单条Async Single Receive异步单条Blocking Batch Receive同步批量连接器所选Async Batch Receive异步批量设计文档解释了放弃异步策略的原因SourceBuilder.SourceBuffer与SourceBuilder.TimestampedSourceBuffer不是线程安全的异步回调若另起线程处理结果会产生并发问题而这与异步编程的惯常做法冲突因此选择了同步批量接收。连接器认为批量接收性能更优。Pulsar 的批量接收满足“以下任一条件即完成一批”消息数量达到上限消息字节数/条数达到设定阈值等待超时。默认批量接收策略为超时 1000ms、最多 512 条。由于该策略日后可能变化Builder 提供了独立的 setterbatchReceivePolicySupplier(...)供自定义。源码中的默认值PulsarConsumerBuilder.javaprivate static SupplierExBatchReceivePolicy getDefaultBatchReceivePolicySupplier() { final int maxNumMessages 512; final int timeoutInMs 1000; return () - BatchReceivePolicy.builder() .maxNumMessages(maxNumMessages) .timeout(timeoutInMs, TimeUnit.MILLISECONDS) .build(); }在fillBuffer中一批消息被批量取出、逐条投影后写入TimestampedSourceBuffer随后立即调用consumer.acknowledgeAsync(messages)异步确认若确认失败仅记录 warn 级别日志PulsarConsumerBuilder.java。Consumer 源的默认配置还包含consumerName hazelcast-jet-consumer与subscriptionName hazelcast-jet-subscriptionPulsarConsumerBuilder.java并可通过consumerConfig(Map)覆盖。确认机制与存储语义默认情况下Pulsar broker 会在该主题的所有订阅都确认某条消息后将其从主题中删除也可以通过调整 broker 配置让这些“全量确认”的消息继续留存。Consumer 的消费恢复逻辑是“从未确认的第一条开始”看起来类似游标机制但存在隐蔽的丢失风险为支持故障恢复回滚机制必须能在失败时正确工作而Consumer API 的确认机制没有回滚能力。为什么 Consumer 源不支持容错设计文档记录了一次真实的尝试作者曾在两次快照之间把已消费的MessageId存入列表打算在快照时刻统一确认——假设任务不会在快照期间失败。但如果任务恰好在快照过程中失败这些确认无法回滚任务重启后就会永久丢失这些已被确认的消息。若 Consumer API 有提交commit机制本可做到 exactly-once。另一种极端做法是干脆不做确认但这会导致所有消息被永久保存除非配置了淘汰机制对持续运行的任务造成存储压力。结论当前 Consumer 源在消费到一批消息后就立即确认这种实现连 at-least-once 都无法保证更遑论 exactly-once。文档认为“确认的存在性与时机”仍需进一步讨论。Consumer 源设计决策小结使用Shared 订阅模式同步批量接收消息消费完一批消息后立即向 broker 发送确认。关于 broker 不可用的行为值得注意当任务运行期间 Pulsar broker 关闭时Consumer 源不会快速失败——Pulsar 客户端会以指数退避方式持续重连期间只输出 warn 级别日志而不抛异常任务因此继续“挂起”而非报错。Reader API 源Reader API 是 Pulsar 更底层的接口没有订阅概念用户只需告诉它“从哪条消息开始读”。Reader 提供最早earliest与最新latest消息的MessageId可作为读取起点用户也可给出介于两者之间的具体MessageId精确指定起点。注意当前 Reader 源的实现会覆盖用户对起始点的任何偏好设置一律从最早消息开始读取。PulsarReaderBuilder中MessageId offset MessageId.earliest与startMessageId(offset)印证了这一点PulsarReaderBuilder.java。文档将此列为“可能需要增强”的点。Reader 源的容错快照即游标Reader 源之所以能支持 exactly-once关键在于它可以把“最后读到哪条消息”序列化进 Jet 快照。源码实现非常直接PulsarReaderBuilder.javabyte[] createSnapshot() { return offset.toByteArray(); } void restoreSnapshot(Listbyte[] snapshots) throws IOException { offset MessageId.fromByteArray(snapshots.get(0)); }工作流程为每次读取消息后更新offset message.getMessageId()createSnapshot把最新MessageId的字节数组写入快照失败恢复时restoreSnapshot从快照取回该MessageIdReader 从该位置继续读。因此只要管道其余部分同样提供相应保证Reader 源即可达成exactly-once。内部实现上Reader 使用容量 1024 的ArrayBlockingQueue缓冲通过readerListener回调入队fillBuffer每次最多处理 128 条QUEUE_CAP 1024、MAX_FILL_MESSAGES 128见 PulsarReaderBuilder.java默认readerName hazelcast-jet-reader。与 Consumer 相同的 broker 断连行为Reader 源与 Consumer 源一样broker 关闭时不快速失败仅输出 warn 日志因为 Pulsar 客户端会指数退避重连且不抛异常。Sink基于 Producer APIProducer API 提供 4 种发布选项Sync Single Publish同步单条Async Single Publish异步单条连接器所选Sync Batch Publish同步批量Async Batch Publish异步批量连接器选择Async Single Publish每条消息通过messageBuilder.sendAsync()异步发送发送失败时把异常记录到AtomicReferenceThrowable error在下一次add/flush时以JetException(Error during message send, ...)重新抛出PulsarSinkBuilder.java。设计文档还说明该异步发送带有重试机制最大重试次数设为 10该参数可通过producerConfig覆盖。Sink Builder 与可选字段Sink 同样采用 Builder 模式PulsarSinkBuilder构造函数接收必填字段其余字段经 setter 追加。除topic、connectionSupplier、schemaSupplier、extractValueFn必填外可选 setter 包括Setter作用extractKeyFn从流元素提取消息 keyextractPropertiesFn从流元素提取消息属性 MapextractTimestampFn从流元素提取事件时间戳写入消息 eventTimeproducerConfig覆盖 Pulsar Producer 配置preferredLocalParallelismSink 本地并行度默认 2add方法依次应用 value/key/properties/timestamp 提取函数后异步发送PulsarSinkBuilder.javaflush时调用producer.flush()并重抛累积错误。Sink 的容错为什么用不了 deduplicationPulsar 作为消息系统刻意追求低延迟没有提供提交commit机制它用去重机制deduplication来支持 exactly-once开启后 broker 依据消息发布时写入的SequenceId字段识别并剔除重复消息。问题在于Jet Pipeline 中处理元素的顺序无法保证在任务重启后保持不变因此无法在重启后为同一条消息分配相同的SequenceId而错误的SequenceId分配会导致消息丢失。于是连接器无法利用 Pulsar 的去重机制。源与 Sink 的公共属性Builder 模式与默认值源的PulsarSources与 Sink 的PulsarSinks均以 Builder 模式创建对象必填参数进构造函数可选参数走 setter全部配置都有默认值。可验证的默认值包括ConsumerconsumerName、subscriptionName、batch receive 策略512 条 / 1000msReaderreaderNameSinkpreferredLocalParallelism 2、空producerConfig。build()时会做参数校验topics非空、connectionSupplier与dataConnectionRef二选一等非法配置抛IllegalArgumentException对应的校验用例见 PulsarSourceTest.java 与 PulsarSinkTest.java。Projection 函数数据格式转换无论源还是 Sink都需要用户提供投影函数作为必填参数负责按数据传输方向转换数据格式源方向把收到的 Pulsar 消息转换为 Jet 可序列化的发射项FunctionExMessageM, TSink 方向把处理后的流元素转换为 Pulsar 消息形式extractValueFn。PulsarSources提供了多个重载可只传 schema 而省略投影函数此时默认FunctionEx.identity()直接发射MessageM。时间戳语义Jet 源的发射项必须携带时间戳。连接器的规则是若 Pulsar 消息存在Event Time则以其作为发射项时间戳否则使用Publish Time总是存在。这一逻辑在两个源的fillBuffer中完全一致PulsarConsumerBuilder.java、PulsarReaderBuilder.java。对应地Sink 可通过extractTimestampFn把自定义时间写入消息的 eventTime。Schema强制类型安全Pulsar 提供名为 Schema 的类型安全机制。虽然 Pulsar 原生允许无 Schema 发消息但连接器在创建源与 Sink 时强制要求 SchemaSchemaSupplier为必填参数。文档特别提醒不使用任何 Schema 在实践中几乎等同于使用Schema.BYTE。连接两种方式connectionSupplier 与 DataConnection除传统的connectionSupplier(() - PulsarClient.builder().serviceUrl(...).build())外三个 Builder 均支持.dataConnectionRef(...)方式复用共享连接。Utils.getClient负责二选一Utils.java。共享连接由PulsarDataConnection管理PulsarDataConnection.java其属性包括属性说明brokerUrlPulsar Broker URL必填httpServiceUrlPulsar HTTP 管理端点用于列出主题资源tenant/namespace可选过滤资源列表配置示例与测试一致见 PulsarSourceTest.javaConfig conf smallInstanceConfig(); conf.addDataConnectionConfig(new DataConnectionConfig(myPulsar) .setShared(true) .setType(Pulsar) .setProperty(PulsarDataConnection.Properties.BROKER_URL, pulsar://localhost:6650));测试细节如何验证连接器测试支撑类PulsarTestSupport extends JetTestSupportPulsarTestSupport.java负责提供TestcontainersPulsarContainerPulsar 4 与 Pulsar 5 两个镜像版本参数化测试供测试使用的 Pulsar 客户端、Producer/Consumer源与 Sink 的 setup 函数setupConsumerSource、setupReaderSource、setupSink。源测试流程先启动 Pulsar 容器用测试专用客户端向预定义主题发布一批消息然后用 Pulsar 源运行 Jet 任务读取该主题最后收集发射项断言全部消息都被读到例如 PulsarSourceTest.java 在双成员集群上发布并校验 1000 条消息。Sink 测试则完全反向执行见 PulsarSinkTest.java。Reader 源容错测试模拟分布式节点故障恢复——创建两个 Jet 实例提交任务确认任务已产生至少一个快照后通过终止其中一个 Jet 实例强制任务重启最后校验消息既不丢失也不重复。测试同时覆盖ProcessingGuarantee.NONE与EXACTLY_ONCE两种保证级别并显式断言快照确实产生PulsarSourceTest.java。此外PulsarSinkTest.java 还用 Toxiproxy 注入带宽故障验证 sink 在 broker 短暂不可用时的恢复与最终失败行为。未来改进方向设计文档明确列出两项未完成工作让 Reader 源支持分布式当前 Reader 源不支持分布式处理可通过为其添加自定义分区映射机制启用Sink 的 exactly-once若 Jet Pipeline 引入处理元素的有序机制则可在 Pulsar Sink 上实现 exactly-once 处理语义。代码示例说明以下两个示例完整保留自设计文档。使用前需在 Maven/Gradle 中引入extensions/pulsar模块依赖并补充com.hazelcast.jet.pulsar.PulsarSources/PulsarSinks的 import当前仓库中对应工厂类源码位于 PulsarSources.java 与 PulsarSinks.java。Pulsar Reader 源的完整用法可参考设计文档示例的模式与 Consumer 源的区别仅在于工厂方法换成pulsarReaderBuilder(...)Java 代码所需 import 类似不再重复列举。Pulsar Consumer 源示例下面的程序创建一个任务连接位于localhost:6650默认地址的本地 Pulsar 集群并从主题hazelcast-demo-topic消费消息package com.hazelcast.jet.contrib.pulsar; import com.hazelcast.jet.Jet; import com.hazelcast.jet.JetInstance; import com.hazelcast.jet.pipeline.Pipeline; import com.hazelcast.jet.pipeline.Sinks; import com.hazelcast.jet.pipeline.StreamSource; import org.apache.pulsar.client.api.BatchReceivePolicy; import org.apache.pulsar.client.api.Message; import org.apache.pulsar.client.api.PulsarClient; import org.apache.pulsar.client.api.Schema; import java.util.Collections; import java.util.HashMap; import java.util.Map; import java.util.concurrent.TimeUnit; public class PulsarConsumerDemo { public static void main(String[] args) { JetInstance jet Jet.bootstrappedInstance(); final StreamSourceInteger pulsarConsumerSource PulsarSources.pulsarConsumerBuilder( hazelcast-demo-topic, () - PulsarClient.builder().serviceUrl(pulsar://localhost:6650).build(), () - Schema.INT32, Message::getValue).build(); Pipeline pipeline Pipeline.create(); pipeline.readFrom(pulsarConsumerSource) .withoutTimestamps() .writeTo(Sinks.logger()); jet.newJob(pipeline).join(); } }上面创建的源使用默认客户端配置可通过 Builder 方法修改例如用batchReceivePolicySupplier(() - BatchReceivePolicy.builder().maxNumMessages(256).timeout(500, TimeUnit.MILLISECONDS).build())调整批量策略或用consumerConfig(map)覆盖 consumer/subscription 名称。Pulsar Sink 示例下面的程序创建一个任务连接本地 Pulsar 集群localhost:6650并把消息发布到主题hazelcast-demo-topicpackage com.hazelcast.jet.contrib.pulsar; import com.hazelcast.function.FunctionEx; import com.hazelcast.jet.Jet; import com.hazelcast.jet.JetInstance; import com.hazelcast.jet.config.JobConfig; import com.hazelcast.jet.pipeline.Pipeline; import com.hazelcast.jet.pipeline.Sink; import com.hazelcast.jet.pipeline.test.TestSources; import org.apache.pulsar.client.api.PulsarClient; import org.apache.pulsar.client.api.Schema; import java.util.HashMap; import java.util.Map; public class PulsarProducerDemo { public static void main(String[] args) { JetInstance jet Jet.bootstrappedInstance(); Pipeline p Pipeline.create(); SinkInteger pulsarSink PulsarSinks.builder( hazelcast-demo-topic, () - PulsarClient.builder() .serviceUrl(pulsar://localhost:6650) .build(), () - Schema.INT32, FunctionEx.identity()).build(); p.readFrom(TestSources.itemStream(15)) .withoutTimestamps() .map(x - (int) x.sequence()) .writeTo(pulsarSink); JobConfig jobConfig new JobConfig(); jobConfig.setName(hazelcast-pulsar-producer); jet.newJob(p, jobConfig).join(); } }结语与源码索引Pulsar 连接器是“以取舍换场景适配”的典型设计Consumer 源以分布式消费换取容错缺失适合吞吐优先、允许消息重复或丢失的摄取场景Reader 源以单机读取换取exactly-once 快照恢复适合对数据准确性敏感的流处理任务Sink 则用异步发布 重试换低延迟但因 Jet 无法保证重启后元素顺序暂时无法利用 Pulsar 去重实现 exactly-once。想深入源码的读者可按以下路径继续探索工厂入口PulsarSources.java、PulsarSinks.java三套 BuilderPulsarConsumerBuilder.java、PulsarReaderBuilder.java、PulsarSinkBuilder.java共享连接PulsarDataConnection.java、Utils.java集成测试PulsarTestSupport.java、PulsarSourceTest.java、PulsarSinkTest.java赞分享缓存KV存储消息队列流处理后端【免费下载链接】hazelcastHazelcast is a unified real-time data platform combining stream processing with a fast data store, allowing customers to act instantly on>项目地址https://gitcode.com/gh_mirrors/ha/hazelcast点击查看免费下载相关推荐Apache Pulsar Pulsar IO 连接器全解析Source 与 Sink 架构、内置连接器与管理实战Apache Pulsar Pulsar IO 连接器全解析Source 与 Sink 架构、内置连接器与管理实战 Pulsar IO 是 Apache Pu消息队列后端流处理Apache Pulsar Go 客户端实战指南安装、连接配置与 Producer/Consumer/Reader 开发Apache Pulsar Go 客户端实战指南安装、连接配置与 Producer/Consumer/Reader 开发 本篇指南以 Apache Pulsa消息队列后端流处理Apache Pulsar Java 客户端完全指南Producer、Consumer、Reader 与 Schema 实战Apache Pulsar Java 客户端完全指南Producer、Consumer、Reader 与 Schema 实战 Apache Pulsar 的消息队列后端流处理上一篇Shader Park Core实战案例10个惊艳效果展示与代码解析下一篇livego的安全审计直播服务的操作日志与合规性检查创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考