首页/新闻资讯/正文详情

Apache Pulsar 客户端深入解析:Client API、连接建立流程与 Reader 手动游标接口

发布时间:2026/9/27 21:57:14 来源:云帆数科 栏目:资讯中心
Apache Pulsar 客户端深入解析:Client API、连接建立流程与 Reader 手动游标接口
消息队列后端流处理【免费下载链接】pulsarApache Pulsar - distributed pub-sub messaging system项目地址https://gitcode.com/gh_mirrors/pulsar28/pulsar点击查看免费下载导读本文基于 Apache Pulsar 官方 2.3.2 文档《Pulsar Clients》系统讲解 Pulsar 客户端库的三大核心主题多语言 Client API 的定位与能力、创建 Producer/Consumer 时的两阶段连接建立流程含断线重连与指数退避以及区别于 Consumer 的Reader 接口手动管理游标、按 MessageId 精确定位读取。文中所有机制均对照当前仓库pulsar-client模块的源码实现与单元测试展开佐证读完你可以掌握 Reader 的三种起始位置配置、底层实现原理以及ReaderBuilder的全部关键配置项具备直接在项目中使用 Reader 实现从指定消息开始读取的能力。Pulsar 客户端库多语言绑定与 API 设计定位Pulsar 为应用层提供了一套客户端 API并针对主流语言提供了官方绑定语言官方文档仓库中的对应模块Javaclient-libraries-java.mdpulsar-client-api、pulsar-clientGoclient-libraries-go.mdpulsar-function-go之外独立的 Go 客户端Pythonclient-libraries-python.mdpulsar-client-cpp/pythonCclient-libraries-cpp.mdpulsar-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 分配给负载最低的 Brokerleast 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 就是这个策略的实现默认初始间隔100msDEFAULT_INTERVAL_IN_NANOSECONDS最大退避间隔30sMAX_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并从最早可用消息及其之后开始读取 Readerbyte[] reader pulsarClient.newReader() .topic(reader-api-test) .startMessageId(MessageId.earliest) .create(); while (true) { Message message reader.readNext(); // 处理消息 }示例二从最新可用消息开始读取Readerbyte[] reader pulsarClient.newReader() .topic(topic) .startMessageId(MessageId.latest) .create();注意MessageId.latest表示从 Topic 末尾开始读实际读到的第一条消息是reader 创建之后新发布的消息。示例三从最早与最新之间的某条消息开始读取byte[] msgIdBytes // 某个字节数组例如从缓存或持久化存储中取回 MessageId id MessageId.fromByteArray(msgIdBytes); Readerbyte[] 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(ReaderListenerT)设置监听器回调模式设置后不能再调用readNext()消息到达时自动回调receiverQueueSize(int)控制本地接收队列大小默认 1000调大可提吞吐但增加内存占用readerName(String)指定 Reader 名称用于监控统计中追踪默认随机生成subscriptionRolePrefix(String)设置订阅角色前缀默认前缀readersubscriptionName(String)显式指定订阅名与subscriptionRolePrefix同时设置时以它为准readCompacted(boolean)读取压缩后 Topic 的最新值视图仅持久化 Topic 可用非持久化 Topic 开启会抛PulsarClientExceptionkeyHashRange(Range...)限定只读取消息 key 哈希落在指定范围内的消息总哈希范围为 65536range 最大 end 应 ≤ 65535cryptoKeyReader(...)/cryptoFailureAction(...)配置消息解密与失败动作失败动作默认FAILpoolMessages(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 中的订阅模型与分区概念综合判断。赞分享消息队列后端流处理【免费下载链接】pulsarApache Pulsar - distributed pub-sub messaging system项目地址https://gitcode.com/gh_mirrors/pulsar28/pulsar点击查看免费下载相关推荐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消息队列后端流处理上一篇Apache Beam GCP 安全日志分析器从 Log Sink 配置到每周 IAM 安全告警的自动化实践下一篇wgpu Mesh Shader 实战指南基于任务着色器、Payload 与逐图元数据渲染三角形创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考

相关推荐

CodeBuddy-CN 介绍与使用指南:从安装到 TaoToken 统一 Key 配置实战
CodeBuddy-CN 介绍与使用指南:从安装到 TaoToken 统一 Key 配置实战

/* MD / 富文本中的 .toc(含博客园搬家等嵌套结构);.toc-box 在侧栏,不受影响 */#content_views .toc,/* 编辑器常在目录前后插入空 p(:empty 仍占 20px),一并去掉避免顶空隙 */#content_views.markdown_views > p:empty:has(+ .toc),#content_views.markdown_views … · 2026/9/27 21:57:08

VSCode 插件分享:6 个 Vue3 开发必备插件,附 TaoToken 统一 Key 配置骨架
VSCode 插件分享:6 个 Vue3 开发必备插件,附 TaoToken 统一 Key 配置骨架

/* MD / 富文本中的 .toc(含博客园搬家等嵌套结构);.toc-box 在侧栏,不受影响 */#content_views .toc,/* 编辑器常在目录前后插入空 p(:empty 仍占 20px),一并去掉避免顶空隙 */#content_views.markdown_views > p:empty:has(+ .toc),#content_views.markdown_views … · 2026/9/27 21:56:49

RL-赵-(七)-不基于模型2-计算V/StateValue-TD算法:狭义TD算法04【TD对比MC:在线/离线、持续/一次性任务、BS/非BS、低/高方差、开始时偏差大/小】
RL-赵-(七)-不基于模型2-计算V/StateValue-TD算法:狭义TD算法04【TD对比MC:在线/离线、持续/一次性任务、BS/非BS、低/高方差、开始时偏差大/小】

5、Algorithm properties:TD算法与蒙特卡洛算法比较 尽管TD learning和MC learning都是model-free方法,与MC learning相比,TD learning的优点和缺点分别是什么呢?TD/Sarsa learningMC learning在线学习:时差学习是在线… · 2026/9/27 21:56:49

互功率谱计算实战:jiufang V1.4工具与Welch平均参数详解
互功率谱计算实战:jiufang V1.4工具与Welch平均参数详解

简介:面向信号处理研究人员与MATLAB学习者的时延估计实用方案,以互功率谱(互功率密度函数)为理论基础,引入压缩传感技术,突破传统奈奎斯特采样速率限制,在稀疏信号条件下仍能保持较高估计精度&a… · 2026/9/27 23:04:53

安徽公有云软件有限公司专业吗
安徽公有云软件有限公司专业吗

从千禧年初企业信息化刚刚起步,到如今数字经济浪潮席卷千行百业,安徽企业数字化服务赛道已经走过了近二十载的变迁。不少服务商来了又走,也有扎根本土的团队始终坚守在客户身边,陪伴一代安徽企业从中小微成长为行业中坚。安徽公有… · 2026/9/27 23:04:53

农行BRIDGE商户直连Java V1.4 Demo接入与避坑指南
农行BRIDGE商户直连Java V1.4 Demo接入与避坑指南

简介:这份资源是中国农业银行缴费中心BRIDGE新版商户直连的Java版DEMO(V1.4),面向需要接入农行缴费支付接口的商户研发与系统集成人员。DEMO覆盖订单处理、支付确认、退款及回调通知等核心交易环节,并附有对应版本的接… · 2026/9/27 23:04:53

临沂太阳能一体化光源电路设计与工程选型标准解析
临沂太阳能一体化光源电路设计与工程选型标准解析

临沂太阳能一体化光源电路设计与工程选型标准解析在太阳能路灯工程应用中,一体化光源因其集成度高、安装便捷、免布线等优势,正逐步成为道路照明、园区亮化、新农村建设等场景的主流选择。本文从电路设计底层逻辑出发,结合工程选型中的常见误… · 2026/9/27 23:04:47

用Python打造销售数据可视化看板:Pandas+Flask+ECharts实战
用Python打造销售数据可视化看板:Pandas+Flask+ECharts实战

简介:这是一套完整的Python销售数据可视化看板项目,面向希望提升数据可视化技能的数据分析学习者与业务人员,解决如何用Python高效构建交互式销售数据看板的问题。压缩包共3个文件,涵盖Python源码、Excel销售数据集及依赖说明文本… · 2026/9/27 23:04:40

卫星链路计算中信号带宽的顶层设计:从符号速率到链路余量
卫星链路计算中信号带宽的顶层设计:从符号速率到链路余量

简介:这是一套基于MATLAB的卫星链路计算工具包,面向卫星通信工程师与相关研究者,可用于星地链路中的信号带宽、空间损耗、天线增益及链路预算等关键参数计算。压缩包为RAR格式,大小约4KB,包含多个.m脚本文件&#xff0… · 2026/9/27 23:04:40

MATLAB雷达信号脉冲压缩仿真:LFM线性调频、匹配滤波与距离分辨率实现
MATLAB雷达信号脉冲压缩仿真:LFM线性调频、匹配滤波与距离分辨率实现

简介:这套Matlab仿真工具完整呈现雷达信号脉冲压缩过程,从线性调频(LFM)信号生成、目标回波仿真到匹配滤波压缩处理均有可运行代码支撑,面向电子信息工程、计算机、数学等专业学生,适用于课程设计、期末大作… · 2026/9/27 0:00:01

汕头网站建设制作厂家避坑指南:5大注意事项救急
汕头网站建设制作厂家避坑指南:5大注意事项救急

汕头网站建设制作厂家避坑指南:5大注意事项救急 改个需求建站公司拖一周,这种憋屈事我见得太多了。 很多汕头老板找本地建站团队,签合同前看着方案挺美,一上线就变脸。 今天不聊虚的,直接拆解找 汕头网站建设制作厂家 时的5个核心 注意事项… · 2026/9/27 0:00:01

多模态虚假新闻检测实战:BERT+ResNet双塔与对比学习
多模态虚假新闻检测实战:BERT+ResNet双塔与对比学习

简介:基于PyTorch的多模态虚假新闻检测项目完整代码包,面向自然语言处理与计算机视觉交叉方向的开发者、科研人员及毕业设计选题者,解决社交媒体中文本与图像联合识别虚假新闻的问题。系统以BERT预训练模型提取文本语义特征,以Res… · 2026/9/27 0:00:01

MATLAB雷达信号脉冲压缩仿真:LFM线性调频、匹配滤波与距离分辨率实现
MATLAB雷达信号脉冲压缩仿真:LFM线性调频、匹配滤波与距离分辨率实现

简介:这套Matlab仿真工具完整呈现雷达信号脉冲压缩过程,从线性调频(LFM)信号生成、目标回波仿真到匹配滤波压缩处理均有可运行代码支撑,面向电子信息工程、计算机、数学等专业学生,适用于课程设计、期末大作… · 2026/9/27 0:00:01

汕头网站建设制作厂家避坑指南:5大注意事项救急
汕头网站建设制作厂家避坑指南:5大注意事项救急

汕头网站建设制作厂家避坑指南:5大注意事项救急 改个需求建站公司拖一周,这种憋屈事我见得太多了。 很多汕头老板找本地建站团队,签合同前看着方案挺美,一上线就变脸。 今天不聊虚的,直接拆解找 汕头网站建设制作厂家 时的5个核心 注意事项… · 2026/9/27 0:00:01

多模态虚假新闻检测实战:BERT+ResNet双塔与对比学习
多模态虚假新闻检测实战:BERT+ResNet双塔与对比学习

简介:基于PyTorch的多模态虚假新闻检测项目完整代码包,面向自然语言处理与计算机视觉交叉方向的开发者、科研人员及毕业设计选题者,解决社交媒体中文本与图像联合识别虚假新闻的问题。系统以BERT预训练模型提取文本语义特征,以Res… · 2026/9/27 0:00:01

了解更多?预约专属演示

我们的顾问将为您一对一讲解产品与方案

企业微信二维码