消息队列后端流处理【免费下载链接】pulsarApache Pulsar - distributed pub-sub messaging system项目地址https://gitcode.com/gh_mirrors/pulsar28/pulsar点击查看免费下载本文以 adaptors-storm.mdPulsar 2.3.0 版本文档为骨架结合当前仓库中 Pulsar 客户端 API 源码与官方概念文档进行纵深讲解介绍如何通过pulsar-stormAdaptor 让 Apache Pulsar 与 Apache Storm 拓扑双向互通用 Pulsar Spout 把 Topic 上的消息注入 Storm 拓扑用 Pulsar Bolt 把拓扑处理结果发回 Pulsar Topic。读完本文你将掌握 Spout/Bolt 的依赖引入、核心配置、消息映射器实现方式以及失败重放与分区键路由的底层机制。一、Pulsar Storm Adaptor 是什么Pulsar Storm 是 Apache Pulsar 为 Apache Storm 提供的集成适配层Adaptor它基于 Storm 的 Spout/Bolt 抽象为在 Pulsar 与 Storm 之间收发数据提供了核心实现Pulsar Spout以通用 spout 的形式把 Pulsar Topic 上发布的数据注入 Storm 拓扑拓扑的数据入口Pulsar Bolt以通用 bolt 的形式把 Storm 拓扑产生的数据发布到 Pulsar Topic拓扑的数据出口。两者与 Storm 的IRichSpout/IRichBolt生命周期无缝衔接应用只需实现各自的消息映射器MessageToValuesMapper/TupleToMessageMapper即可完成 Pulsar 消息与 Storm 元组Tuple之间的双向转换无需关心底层连接、订阅与确认逻辑。二、引入依赖在 Storm 应用的pom.xml中加入如下依赖即可dependency groupIdorg.apache.pulsar/groupId artifactIdpulsar-storm/artifactId version${pulsar.version}/version /dependency注意事项version应与所用 Pulsar 服务端版本保持一致如 2.3.0避免客户端与服务端协议不兼容从仓库的文档演进可以确认pulsar-storm模块早期位于 Apache Pulsar 主仓库本文档即为其 2.3.0 版本说明后期被迁移至独立的 pulsar-adapters 项目维护当前快照仓库中未包含该模块的源码因此本文的 API 行为以本文档描述及仓库内 Pulsar 客户端公共 API 为准。三、Pulsar Spout将 Pulsar 消息注入 Storm 拓扑Pulsar Spout 允许拓扑消费指定 Topic 上发布的数据。它基于收到的Message和客户端提供的MessageToValuesMapper将消息映射为 Storm 元组Values发射给下游 Bolt。3.1 失败重放语义这是 Spout 最关键的可靠性设计未被下游 Bolt 成功处理的元组fail 的 tuple会被 Spout 以指数退避exponential backoff的方式重新注入重放过程受两个上限约束——可配置的超时时间默认 60 秒或可配置的重试次数两者谁先到谁生效。达到上限后该消息才会被 Pulsar 消费者确认ack。这一机制决定了典型的at-least-once至少一次语义下游处理失败时消息会重放而一旦达到超时/重试上限Spout 主动 ack消息不会再重投。因此业务处理逻辑需要具备幂等性。3.2 Spout 构造示例完整代码MessageToValuesMapper messageToValuesMapper new MessageToValuesMapper() { Override public Values toValues(Message msg) { return new Values(new String(msg.getData())); } Override public void declareOutputFields(OutputFieldsDeclarer declarer) { // declare the output fields declarer.declare(new Fields(string)); } }; // Configure a Pulsar Spout PulsarSpoutConfiguration spoutConf new PulsarSpoutConfiguration(); spoutConf.setServiceUrl(pulsar://broker.messaging.usw.example.com:6650); spoutConf.setTopic(persistent://my-property/usw/my-ns/my-topic1); spoutConf.setSubscriptionName(my-subscriber-name1); spoutConf.setMessageToValuesMapper(messageToValuesMapper); // Create a Pulsar Spout PulsarSpout spout new PulsarSpout(spoutConf);3.3 配置项说明配置方法含义说明setServiceUrl(String)Pulsar 服务地址示例为pulsar://broker.messaging.usw.example.com:6650本地 Standalone 开发常用pulsar://localhost:66506650 是 Pulsar broker 的默认端口setTopic(String)要消费的 Topic必须是完整的 Pulsar Topic 名见下文Topic 命名规则setSubscriptionName(String)订阅名称Spout 内部以消费者身份订阅 Topic订阅名不可省略多个 spout 使用同一订阅名可共享消费进度setMessageToValuesMapper(...)消息→元组映射器决定每条 Pulsar 消息如何转换为 StormValues3.4 映射器与消息 API 的对应关系MessageToValuesMapper.toValues(Message)中的Message即 Pulsar 客户端公共 API 中的org.apache.pulsar.client.api.Message。从 Message.java 可以看到其核心方法byte[] getData()获取消息原始负载示例中用new String(msg.getData())还原为字符串String getKey()获取消息分区键可用于按 key 处理MessageId getMessageId()获取消息 ID可用于精确 ack / 去重。declareOutputFields(OutputFieldsDeclarer)用于声明 Spout 发射的字段名示例中声明了单个string字段下游 Bolt 通过tuple.getString(0)读取。3.5 Topic 命名规则补充示例中的persistent://my-property/usw/my-ns/my-topic1遵循 Pulsar Topic 命名规范。根据 concepts-messaging.md 中的说明Topic 名是结构化的 URL{persistent|non-persistent}://tenant/namespace/topic组成部分说明persistent/non-persistentTopic 类型默认为 persistent消息持久化落盘。若省略类型前缀则默认是持久化 Topictenant租户是 Pulsar 多租户隔离的基本单位旧文档中常写作propertynamespace命名空间大多数 Topic 级配置在命名空间层面完成topic最终的具体 Topic 名注意Topic 无需预先显式创建客户端首次写入或订阅时会在对应命名空间下自动创建。若未指定租户/命名空间则使用默认的public/default例如persistent://public/default/my-topic。四、Pulsar Bolt将 Storm 拓扑数据发布到 PulsarPulsar Bolt 允许把 Storm 拓扑中的数据发布到 Pulsar Topic。它基于收到的 Storm 元组和客户端提供的TupleToMessageMapper构造并发布消息。4.1 分区 Topic 与 Key 路由当目标 Topic 是**分区 Topicpartitioned topic**时可以通过在消息中设置key来实现定向路由TupleToMessageMapper的实现中需要为消息提供 key拥有相同 key 的消息会被发送到同一个分区从而保证同一 key 的消息顺序到达、可被同一消费者处理。这一行为对应 Pulsar 生产者对分区键的路由语义。4.2 Bolt 构造示例完整代码TupleToMessageMapper tupleToMessageMapper new TupleToMessageMapper() { Override public TypedMessageBuilderbyte[] toMessage(TypedMessageBuilderbyte[] msgBuilder, Tuple tuple) { String receivedMessage tuple.getString(0); // message processing String processedMsg receivedMessage -processed; return msgBuilder.value(processedMsg.getBytes()); } Override public void declareOutputFields(OutputFieldsDeclarer declarer) { // declare the output fields } }; // Configure a Pulsar Bolt PulsarBoltConfiguration boltConf new PulsarBoltConfiguration(); boltConf.setServiceUrl(pulsar://broker.messaging.usw.example.com:6650); boltConf.setTopic(persistent://my-property/usw/my-ns/my-topic2); boltConf.setTupleToMessageMapper(tupleToMessageMapper); // Create a Pulsar Bolt PulsarBolt bolt new PulsarBolt(boltConf);4.3 配置项说明配置方法含义说明setServiceUrl(String)Pulsar 服务地址与 Spout 相同pulsar://协议、默认端口 6650setTopic(String)发布目标 Topic若为分区 Topic可通过消息 key 控制分区路由setTupleToMessageMapper(...)元组→消息映射器决定每条 Storm 元组如何构造 Pulsar 消息4.4 映射器与TypedMessageBuilder的对应关系TupleToMessageMapper.toMessage(TypedMessageBuilderbyte[] msgBuilder, Tuple tuple)中的TypedMessageBuilder即 TypedMessageBuilder.java 定义的公共接口它提供了一套链式构造消息的能力value(T)设置消息负载示例中把处理结果processedMsg.getBytes()写入key(String)/keyBytes(byte[])设置消息的分区键用于分区路由分区 Topic 下相同 key 的消息进入同一分区正是上节所述机制orderingKey(byte[])设置排序键用于 Key_Shared 订阅模式下的分发控制send()/sendAsync()同步/异步发送由 Bolt 内部在映射完成后调用应用层无需手动触发。示例中的declareOutputFields为空实现因为 Bolt 通常作为拓扑终点不再发射元组若 Bolt 同时还要继续向拓扑下游发射数据则在此声明相应字段。五、在拓扑中组合 Spout 与 Bolt完整链路示意把上面两个组件串起来即可构成一条Pulsar → Storm 处理 → Pulsar的完整数据管道。以下为基于文档 API 的组合示意拓扑装配代码需按 Storm 版本 API 微调TopologyBuilder builder new TopologyBuilder(); builder.setSpout(pulsar-spout, spout); // 消费 my-topic1 builder.setBolt(pulsar-bolt, bolt) .shuffleGrouping(pulsar-spout); // 接收 spout 的元组Spout 发射的string字段经tuple.getString(0)被 Bolt 读取处理后携带 key 写回my-topic2。当目标为分区 Topic 时为消息设置相同 key 即可让同一逻辑分组的消息始终落到同一分区。六、完整示例与注意事项本文档version-2.3.0/adaptors-storm.md以及当前文档 docs/adaptors-storm.md 均指出完整的可运行示例与PulsarSpout、PulsarBolt的单元测试维护在 Apache Pulsar 的 pulsar-adaptors 系列仓库中如pulsar-storm/src/test/java/org/apache/pulsar/storm/PulsarSpoutTest.java本快照仓库未内置pulsar-storm源码模块如需阅读实现细节请到 pulsar-adapters 项目查看。版本前提本文描述的 API 与行为以 Pulsar 2.3.0 文档为准且当前仓库同时保留了多个版本如 version-2.9.1 等的同类文档内容基本一致若使用其他 Pulsar 版本请以对应版本文档为准。可靠性语义Spout 默认 60 秒超时可配置内以指数退避重放失败元组达到上限后 ack业务需按 at-least-once 设计幂等处理。密钥与鉴权如服务端启用了鉴权/TLS需要在构造 Pulsar 客户端时配置认证信息Spout/Bolt 内部创建的生产者与消费者继承 Pulsar 客户端的鉴权配置。七、小结通过pulsar-stormAdaptorApache Storm 拓扑与 Apache Pulsar 之间实现了双向、可靠的数据互通Spout 侧MessageToValuesMapper负责Pulsar 消息 → Storm 元组失败元组在默认 60 秒窗口内指数退避重放达到上限后 ackBolt 侧TupleToMessageMapper负责Storm 元组 → Pulsar 消息借助TypedMessageBuilder.key()实现分区 Topic 的按 key 路由底层依赖 Pulsar 公共客户端 APIMessage、TypedMessageBuilder、PulsarClient等与 pulsar-client-api 中的接口一一对应便于在源码层面继续深入。上述实战要点均可在仓库 site2/website-next/versioned_docs/version-2.3.0/adaptors-storm.md 及客户端 API 源码中进一步核对。赞分享消息队列后端流处理【免费下载链接】pulsarApache Pulsar - distributed pub-sub messaging system项目地址https://gitcode.com/gh_mirrors/pulsar28/pulsar点击查看免费下载相关推荐Apache Pulsar 与 Apache Storm 集成指南Pulsar Storm Adaptor 的 Spout 与 Bolt 实战Apache Pulsar 与 Apache Storm 集成指南Pulsar Storm Adaptor 的 Spout 与 Bolt 实战 导读 本文讲解消息队列后端流处理Apache Pulsar 与 Apache Storm 集成指南Pulsar Storm Adaptor 的 Spout 与 Bolt 完整实战Apache Pulsar 与 Apache Storm 集成指南Pulsar Storm Adaptor 的 Spout 与 Bolt 完整实战 本篇技术指消息队列后端流处理Apache Pulsar 集成 Apache StormPulsar Storm Adaptor 的 Spout 与 Bolt 使用指南Apache Pulsar 集成 Apache StormPulsar Storm Adaptor 的 Spout 与 Bolt 使用指南 本文以 Apach消息队列后端流处理上一篇Parabolic后台下载功能系统托盘通知与任务管理下一篇魔百盒CM201-1刷Armbian从安卓TV到7×24在线Linux主机创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
