- 消息队列
- 后端
- 流处理
【免费下载链接】pulsar
Apache Pulsar - distributed pub-sub messaging system
本文以 adaptors-storm.md(Pulsar 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> <groupId>org.apache.pulsar</groupId> <artifactId>pulsar-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:6650(6650 是 Pulsar broker 的默认端口) |
setTopic(String) | 要消费的 Topic | 必须是完整的 Pulsar Topic 名,见下文"Topic 命名规则" |
setSubscriptionName(String) | 订阅名称 | Spout 内部以消费者身份订阅 Topic,订阅名不可省略;多个 spout 使用同一订阅名可共享消费进度 |
setMessageToValuesMapper(...) | 消息→元组映射器 | 决定每条 Pulsar 消息如何转换为 StormValues |
3.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-persistent | Topic 类型;默认为 persistent(消息持久化落盘)。若省略类型前缀,则默认是持久化 Topic |
tenant | 租户,是 Pulsar 多租户隔离的基本单位(旧文档中常写作property) |
namespace | 命名空间,大多数 Topic 级配置在命名空间层面完成 |
topic | 最终的具体 Topic 名 |
注意:Topic 无需预先显式创建,客户端首次写入或订阅时会在对应命名空间下自动创建。若未指定租户/命名空间,则使用默认的public/default,例如persistent://public/default/my-topic。
四、Pulsar Bolt:将 Storm 拓扑数据发布到 Pulsar
Pulsar Bolt 允许把 Storm 拓扑中的数据发布到 Pulsar Topic。它基于收到的 Storm 元组和客户端提供的TupleToMessageMapper构造并发布消息。
4.1 分区 Topic 与 Key 路由
当目标 Topic 是**分区 Topic(partitioned topic)**时,可以通过在消息中设置key来实现定向路由:TupleToMessageMapper的实现中需要为消息提供 key,拥有相同 key 的消息会被发送到同一个分区,从而保证同一 key 的消息顺序到达、可被同一消费者处理。这一行为对应 Pulsar 生产者对分区键的路由语义。
4.2 Bolt 构造示例(完整代码)
TupleToMessageMapper tupleToMessageMapper = new TupleToMessageMapper() { @Override public TypedMessageBuilder<byte[]> toMessage(TypedMessageBuilder<byte[]> 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://协议、默认端口 6650 |
setTopic(String) | 发布目标 Topic | 若为分区 Topic,可通过消息 key 控制分区路由 |
setTupleToMessageMapper(...) | 元组→消息映射器 | 决定每条 Storm 元组如何构造 Pulsar 消息 |
4.4 映射器与TypedMessageBuilder的对应关系
TupleToMessageMapper.toMessage(TypedMessageBuilder<byte[]> 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-stormAdaptor,Apache Storm 拓扑与 Apache Pulsar 之间实现了双向、可靠的数据互通:
- Spout 侧:
MessageToValuesMapper负责"Pulsar 消息 → Storm 元组",失败元组在默认 60 秒窗口内指数退避重放,达到上限后 ack; - Bolt 侧:
TupleToMessageMapper负责"Storm 元组 → Pulsar 消息",借助TypedMessageBuilder.key()实现分区 Topic 的按 key 路由; - 底层依赖 Pulsar 公共客户端 API(
Message、TypedMessageBuilder、PulsarClient等),与 pulsar-client-api 中的接口一一对应,便于在源码层面继续深入。
上述实战要点均可在仓库 site2/website-next/versioned_docs/version-2.3.0/adaptors-storm.md 及客户端 API 源码中进一步核对。
- 消息队列
- 后端
- 流处理
【免费下载链接】pulsar
Apache Pulsar - distributed pub-sub messaging system
相关推荐
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 Storm:Pulsar Storm Adaptor 的 Spout 与 Bolt 使用指南
Apache Pulsar 集成 Apache Storm:Pulsar Storm Adaptor 的 Spout 与 Bolt 使用指南 本文以 Apach
消息队列后端流处理
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考