- 消息队列
- 后端
- 流处理
【免费下载链接】pulsar
Apache Pulsar - distributed pub-sub messaging system
导读
本文基于 Apache Pulsar 官方 2.3.2 文档《Pulsar Clients》,系统讲解 Pulsar 客户端库的三大核心主题:多语言 Client API 的定位与能力、创建 Producer/Consumer 时的两阶段连接建立流程(含断线重连与指数退避),以及区别于 Consumer 的Reader 接口(手动管理游标、按 MessageId 精确定位读取)。文中所有机制均对照当前仓库pulsar-client模块的源码实现与单元测试展开佐证,读完你可以掌握 Reader 的三种起始位置配置、底层实现原理,以及ReaderBuilder的全部关键配置项,具备直接在项目中使用 Reader 实现"从指定消息开始读取"的能力。
Pulsar 客户端库:多语言绑定与 API 设计定位
Pulsar 为应用层提供了一套客户端 API,并针对主流语言提供了官方绑定:
| 语言 | 官方文档 | 仓库中的对应模块 |
|---|---|---|
| Java | client-libraries-java.md | pulsar-client-api、pulsar-client |
| Go | client-libraries-go.md | pulsar-function-go之外独立的 Go 客户端 |
| Python | client-libraries-python.md | pulsar-client-cpp/python |
| C++ | client-libraries-cpp.md | pulsar-client-cpp |
客户端 API 的核心价值在于:它将 Pulsar 客户端与 Broker 之间通信协议的复杂性封装起来,向上对应用暴露一套简单直观的编程接口,应用只需关注"发消息 / 收消息"本身。
从仓库源码结构看,这一设计体现得非常清晰:
- 接口层定义在 pulsar-client-api/src/main/java/org/apache/pulsar/client/api,其中
Reader、ReaderBuilder、PulsarClient等均为@InterfaceAudience.Public、@InterfaceStability.Stable的公开稳定接口; - 实现层位于 pulsar-client/src/main/java/org/apache/pulsar/client/impl,例如
PulsarClientImpl、ReaderImpl、ConsumerImpl、Backoff等。
底层能力由官方客户端库自动承担,应用无需关心,主要包括:
- 透明重连 / 连接故障转移:TCP 连接断开后客户端自动重连,甚至可以在 Broker 之间切换;
- 消息排队:消息在 Broker 确认之前由客户端本地排队缓冲;
- 带退避的重试:连接重试采用退避(backoff)策略,避免对 Broker 造成重连风暴。
自定义客户端库提示:如果你想从零实现自己的客户端,官方建议先研读 Pulsar 的二进制协议文档:develop-binary-protocol.md。该文档详细定义了客户端与 Broker 之间基于自定义二进制协议的指令格式。
客户端连接建立阶段:从 Topic 归属查找到授权校验
当一个应用要创建 Producer 或 Consumer 时,Pulsar 客户端库会进入一个由两个步骤组成的setup phase(建立阶段):
第一步:Topic 归属查找(Lookup)
客户端首先向 Broker 发送一次HTTP lookup 请求,目标是确定该 Topic 的"归属 Broker"。处理这次请求的可以是任意一个活跃 Broker,它通过(缓存的)ZooKeeper 元数据判断:
- 如果该 Topic 已经有 Broker 在服务,则返回这个"owner Broker"的地址;
- 如果当前没有 Broker 在服务该 Topic,则尝试将 Topic 分配给负载最低的 Broker(least loaded broker)。
从源码看,pulsar-client中的 BinaryProtoLookupService.java 通过LookupService接口提供查找能力,并通过CommandLookupTopicResponse携带LookupType(如Success/Redirect/Failed)返回归属信息;客户端还支持maxLookupRedirects配置来限制查找重定向的次数,防止查找链路过长。
第二步:建立 TCP 连接并创建 Producer/Consumer
拿到 Broker 地址后,客户端:
- 创建一条 TCP 连接(或从连接池中复用已有连接);
- 在连接上完成认证(authentication);
- 在该连接内,客户端与 Broker 通过自定义协议的二进制指令进行交互;
- 客户端发送创建 Producer/Consumer 的命令;
- Broker 在校验授权策略(authorization policy)通过后才予以应答,完成创建。
断线后的自动恢复
无论何时 TCP 连接断开,客户端都会立即重新发起上述建立阶段,并以**指数退避(exponential backoff)**的方式持续重试,直到重新建立 Producer 或 Consumer 成功。
仓库中的 Backoff.java 就是这个策略的实现:默认初始间隔100ms(DEFAULT_INTERVAL_IN_NANOSECONDS),最大退避间隔30s(MAX_BACKOFF_INTERVAL_NANOSECONDS),每次重试间隔按next = min(next * 2, max)指数翻倍增长直至上限,从而在 Broker 恢复期间避免高频冲击。
Consumer 接口与 Reader 接口:游标管理的两种模式
在理解 Reader 之前,先明确 Pulsar 中 Consumer(消费者)的标准工作方式(详细见 concepts-messaging.md):
- 应用使用 Consumer监听 Topic(参见 reference-terminology.md 对 Topic 的定义);
- 处理到达的消息,处理完成后**确认(acknowledge)**这些消息;
- 新建订阅默认初始定位在 Topic 末尾,该订阅下的 Consumer 从之后产生的第一条消息开始读取;
- 当 Consumer 用已存在的订阅连接 Topic 时,则从该订阅中最早未被确认的消息开始读取。
一句话总结:Consumer 接口的订阅游标由 Pulsar 根据消息确认(acknowledgement,参见 concepts-messaging.md#acknowledgement)自动管理。
而Reader(读取器)接口则完全不同:它让应用手动管理游标。使用 Reader 连接 Topic 时,你必须显式指定"从哪条消息开始读"。
Reader 的典型价值:实现精确一次处理语义
Reader 接口对"用 Pulsar 为流处理系统提供精确一次(effectively-once)处理语义"这类场景尤其关键:流处理系统必须能够把 Topic"回退"(rewind)到某条具体消息并从那里重新开始读取。Reader 正是为此提供了低层抽象——让客户端**手动定位(manually position)**自己在 Topic 中的读取位置。
限制:仅支持非分区主题当前 Reader 接口不能用于分区主题(partitioned topics)(分区主题的说明见 concepts-messaging.md#partitioned-topics)。从源码看,ReaderImpl.java 中通过
TopicName.getPartitionIndex(topicName)解析分区索引,Reader 面向的是单一 Topic 的读取语义。
Reader 的三种起始位置(Start Position)
使用 Reader 连接 Topic 时,可以选择以下三种起始读取位置:
| 起始位置 | 说明 | 对应 MessageId |
|---|---|---|
| 最早可用消息(earliest) | 从 Topic 中最旧的那条消息开始读 | MessageId.earliest |
| 最新可用消息(latest) | 从 Topic 末尾开始读 | MessageId.latest |
| 最早与最新之间的任意消息 | 需要显式提供MessageId,由应用提前"知道"该 ID(例如从持久化数据存储或缓存中获取) | 自定义MessageId |
关于第三种位置,ReaderBuilder.java 的接口注释补充了重要细节:定位到指定消息后,读到的第一条消息是指定消息"之后"的那一条(first message read will be the one immediatelyafterthe specified message);如果希望包含指定消息本身,需要调用startMessageIdInclusive()。
Reader 的 Java 实战示例
下面三个示例均来自原文档,并对照 Reader.java 与 ReaderBuilder.java 的 API 进行验证,可直接运行。
示例一:从最早可用消息开始读取
import org.apache.pulsar.client.api.Message; import org.apache.pulsar.client.api.MessageId; import org.apache.pulsar.client.api.Reader; // 在 topic 上创建 reader,并从最早可用消息(及其之后)开始读取 Reader<byte[]> reader = pulsarClient.newReader() .topic("reader-api-test") .startMessageId(MessageId.earliest) .create(); while (true) { Message message = reader.readNext(); // 处理消息 }示例二:从最新可用消息开始读取
Reader<byte[]> reader = pulsarClient.newReader() .topic(topic) .startMessageId(MessageId.latest) .create();注意:
MessageId.latest表示从 Topic 末尾开始读,实际读到的第一条消息是reader 创建之后新发布的消息。
示例三:从最早与最新之间的某条消息开始读取
byte[] msgIdBytes = // 某个字节数组,例如从缓存或持久化存储中取回 MessageId id = MessageId.fromByteArray(msgIdBytes); Reader<byte[]> reader = pulsarClient.newReader() .topic(topic) .startMessageId(id) .create();这里MessageId.fromByteArray(...)负责把之前序列化保存的消息 ID 字节数组还原为MessageId对象,实现"从断点精确续读"。
Reader 的底层实现原理:ReaderImpl 与 ConsumerImpl
从源码结构看,Reader 并不是一套全新的协议实现,而是对 Consumer 的封装与降级使用。核心实现位于 ReaderImpl.java,关键事实如下:
1. Reader 本质上是"非持久、独占"的订阅
ReaderImpl构造时会把ReaderConfigurationData转换为ConsumerConfigurationData,并强制设置:
SubscriptionType.Exclusive(独占订阅);SubscriptionMode.NonDurable(非持久订阅)——这正是 Reader 与普通 Consumer 最本质的差异:Reader 不依赖订阅游标持久化,重连后总是根据指定的起始位置重新定位。
2. 每次读取后立即自动确认
在 ReaderImpl.java 中,readNext()、readNext(timeout, unit)、readNextAsync()三个读取入口在拿到消息后都会立即调用consumer.acknowledgeCumulativeAsync(msg)(累积确认)。源码注释解释得很直白:
"Reader is based on non-durable subscription. When it reconnects, it will specify the subscription position anyway."(Reader 基于非持久订阅,重连时会重新指定订阅位置。)
也就是说:Reader 的确认动作本身没有持久化语义,它存在的意义只是把消息从接收队列里"消费掉",真正的定位逻辑完全由应用指定的起始 MessageId 决定。
3. 读取 API 与额外能力
Reader接口(Reader.java)提供的能力包括:
readNext()/readNext(int timeout, TimeUnit unit)/readNextAsync():同步阻塞读取、带超时读取、异步读取;hasMessageAvailable():检查当前位置之后是否还有可读消息,可用于"扫描到当前最新消息后停止"(源码给出了while (reader.hasMessageAvailable())的用法示例);hasReachedEndOfTopic():判断是否已读到被终止(sealed)的 Topic的末尾;isConnected():检查当前是否与 Broker 保持连接;seek(MessageId)/seek(long timestamp)/seek(Function):运行时将游标重新定位到指定消息 ID 或指定发布时间(注意:该操作同样仅适用于非分区 Topic)。
4. 默认配置
ReaderConfigurationData.java 给出了 Reader 的关键默认值:
receiverQueueSize = 1000:接收队列默认 1000 条消息;cryptoFailureAction = ConsumerCryptoFailureAction.FAIL:解密失败时默认直接失败;readCompacted = false:默认不读取压缩(compacted)后的 Topic 视图。
ReaderBuilder 高级配置项
除topic与startMessageId外,ReaderBuilder.java 还提供了以下常用配置,均可链式调用:
| 配置方法 | 作用 | 要点 / 默认值 |
|---|---|---|
startMessageIdInclusive() | 起始位置包含指定的消息本身(默认是其后一条) | 同样作用于seek()重置操作 |
startMessageFromRollbackDuration(long, TimeUnit) | 按时间回退定位,例如回退 5 分钟,Broker 找到该时间点前最近发布的消息作为起点 | 适合"回到 N 分钟前"的场景 |
readerListener(ReaderListener<T>) | 设置监听器回调模式;设置后不能再调用readNext() | 消息到达时自动回调 |
receiverQueueSize(int) | 控制本地接收队列大小 | 默认 1000,调大可提吞吐但增加内存占用 |
readerName(String) | 指定 Reader 名称,用于监控统计中追踪 | 默认随机生成 |
subscriptionRolePrefix(String) | 设置订阅角色前缀 | 默认前缀reader |
subscriptionName(String) | 显式指定订阅名 | 与subscriptionRolePrefix同时设置时以它为准 |
readCompacted(boolean) | 读取压缩后 Topic 的最新值视图 | 仅持久化 Topic 可用,非持久化 Topic 开启会抛PulsarClientException |
keyHashRange(Range...) | 限定只读取消息 key 哈希落在指定范围内的消息 | 总哈希范围为 65536,range 最大 end 应 ≤ 65535 |
cryptoKeyReader(...)/cryptoFailureAction(...) | 配置消息解密与失败动作 | 失败动作默认FAIL |
poolMessages(boolean) | 启用消息及底层缓冲池复用 | 启用后应用必须调用Message.release(),否则内存泄漏 |
loadConf(Map)/clone() | 从配置 Map 加载 / 克隆 Builder | 克隆便于基于同一份配置创建多个 Reader |
其中subscriptionName/subscriptionRolePrefix在 ReaderImpl.java 中的处理逻辑是:未显式设置订阅名时,自动生成reader-<sha1(UUID).substring(0,10)>,若设置了前缀则拼接为<prefix>-reader-xxxx。
单元测试佐证:Reader 的异步取消与构建行为
仓库中的测试用例印证了上述机制:
- ReaderImplTest.java 验证了
readNextAsync()返回的 Future 可被取消:调用future.cancel(false)后,Reader 内部 Consumer 的待处理接收请求(pending receive)会被移除(hasNextPendingReceive()返回 false),避免异步接收请求积压; - BuildersTest.java 验证了
client.newReader().topic(...).startMessageId(MessageId.earliest).create()的构建链路,以及 Reader 通过try-with-resources正常关闭的行为。
小结
Pulsar 客户端 API 将复杂的 lookup、认证、二进制协议交互、断线重连与退避重试全部封装在官方客户端库中,应用只需面向简单直观的接口编程。当需要手动管理读取位置(而非依赖订阅游标自动推进)时,应选择 Reader 接口:它支持从 earliest、latest 或任意指定 MessageId 开始读取,底层以非持久独占订阅封装 Consumer 实现,配合seek()与startMessageIdInclusive()等能力,为流处理系统实现精确回退与断点续读提供了坚实的低层抽象。需要注意的是,当前 Reader 仅支持非分区 Topic,实际选型时应结合 concepts-messaging.md 中的订阅模型与分区概念综合判断。
- 消息队列
- 后端
- 流处理
【免费下载链接】pulsar
Apache Pulsar - distributed pub-sub messaging system
相关推荐
Apache Pulsar 客户端概念详解:Client API、连接建立流程与 Reader 手动游标接口
Apache Pulsar 客户端概念详解:Client API、连接建立流程与 Reader 手动游标接口 导读 本文基于 Apache Pulsar 官方概
消息队列后端流处理Apache Pulsar 客户端接口深度解析:连接建立流程与 Reader 手动游标机制
Apache Pulsar 客户端接口深度解析:连接建立流程与 Reader 手动游标机制 导读 :本文以 Apache Pulsar 2.1.1 incuba
消息队列后端流处理Apache Pulsar 客户端 API 与 Reader 接口深度解析:从连接建立到手动游标控制
Apache Pulsar 客户端 API 与 Reader 接口深度解析:从连接建立到手动游标控制 Apache Pulsar 为应用提供了面向 Java、G
消息队列后端流处理
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考