ARTICLE DETAIL

资讯详情

深耕郑州网站建设与运营推广的一线实战洞察。

Apache Pulsar 客户端库全景指南:从语言生态选型到源码级原理(Getting Started with Pulsar Clients)

Apache Pulsar 客户端库全景指南:从语言生态选型到源码级原理(Getting Started with Pulsar Clients) 消息队列后端流处理【免费下载链接】pulsarApache Pulsar - distributed pub-sub messaging system项目地址https://gitcode.com/gh_mirrors/pulsar28/pulsar点击查看免费下载Apache Pulsar 官方为开发者提供了 Java、Go、Python、C、Node.js、WebSocket、C# 等完整客户端生态本文以site2/docs/getting-started-clients.md为主线系统梳理各语言客户端的安装接入、连接 URL 规范、Producer/Consumer/Reader 三大核心 API并深入客户端连接建立、查找lookup与重连的底层源码实现。读完本文你将能够根据项目语言栈正确选型客户端快速完成从“连接 Pulsar”到“生产/消费消息”的端到端落地并理解客户端在断线重连、认证授权等场景下的内部工作机制。本文所有代码与配置均取自当前仓库的真实文档与源码可直接对照仓库文件路径进一步查阅。一、Pulsar 客户端库生态概览Pulsar 是一个分布式发布订阅pub-sub消息系统其客户端 API 将 Pulsar 自定义的客户端- Broker 通信协议封装起来对外暴露简单直观的编程接口。当前仓库 site2/docs/getting-started-clients.md 明确列出了以下官方支持的客户端库语言官方文档Javaclient-libraries-java.mdGoclient-libraries-go.mdPythonclient-libraries-python.mdCclient-libraries-cpp.mdNode.jsclient-libraries-node.mdWebSocketclient-libraries-websocket.mdC#client-libraries-dotnet.md从客户端概念文档可以了解到官方客户端的底层行为是一致的支持透明的重连与 Broker 故障切换、在消息被 Broker 确认前进行排队并内置带退避backoff的连接重试启发式策略。这意味着无论你使用哪种语言客户端都会替你处理连接恢复等琐碎工作。各语言客户端的实现形态从仓库源码结构可以观察到各客户端的实现差异Java 客户端完全原生实现核心代码位于 pulsar-client 与 pulsar-client-api 模块org.apache.pulsar.client.api包提供 Producer/Consumer/Reader 编程接口C 客户端原生实现于 pulsar-client-cpp构建产物包括libpulsar.so、libpulsar.a等并提供perfProducer/perfConsumer性能测试工具Python 客户端是 C 客户端之上的封装wrapper源码位于 pulsar-client-cpp/python因此继承 C 客户端全部能力Node.js 客户端同样基于 C 客户端通过node-addon-api模块包装因此要求先安装 C 客户端库C# 客户端DotPulsar与Go 客户端pulsar-client-go则是独立于本仓库的官方项目。二、Pulsar 协议连接 URL 规范无论使用哪个语言客户端接入 Pulsar 的第一步都是构造正确的Pulsar protocol URL。官方文档Java、Go、Node.js 等给出了统一规范协议 scheme 为pulsar默认端口6650本地开发standalone 模式默认地址pulsar://localhost:6650多 Broker 时可用逗号分隔多个地址启用 TLS 时 scheme 变为pulsarssl默认端口 6651。各场景 URL 示例# 单机本地 pulsar://localhost:6650 # 多 Broker逗号分隔 pulsar://localhost:6650,localhost:6651,localhost:6652 # 生产集群 pulsar://pulsar.us-west.example.com:6650 # 启用 TLS pulsarssl://pulsar.us-west.example.com:6651在 getting-started-standalone.md 描述的 standalone 模式下Broker 默认就监听pulsar://localhost:6650本地跑通示例代码前无需修改任何配置。三、各语言客户端快速上手3.1 Java 客户端Java 客户端支持 Producer、Consumer、Reader 以及 TableView且所有方法线程安全。API 按包分为两个域包说明Maven 构件org.apache.pulsar.client.apiProducer/Consumer APIorg.apache.pulsar:pulsar-clientorg.apache.pulsar.client.admin管理 API见 admin-api-overview.mdorg.apache.pulsar:pulsar-client-adminorg.apache.pulsar.client.all同时包含以上两者且统一打 shade避免重复类org.apache.pulsar:pulsar-client-allMaven 依赖properties pulsar.version2.8.0/pulsar.version /properties dependencies dependency groupIdorg.apache.pulsar/groupId artifactIdpulsar-client/artifactId version${pulsar.version}/version /dependency /dependenciesGradle 依赖def pulsarVersion 2.8.0 dependencies { compile group: org.apache.pulsar, name: pulsar-client, version: pulsarVersion }创建客户端并指定loadConf常用参数完整参数表见 client-libraries-java.mdPulsarClient client PulsarClient.builder() .serviceUrl(pulsar://localhost:6650) .build();参数类型说明默认值serviceUrlStringPulsar 服务地址无operationTimeoutMslong操作超时时间30000statsIntervalSecondslong统计信息输出间隔1 秒60numIoThreadsint处理 Broker 连接的 IO 线程数1numListenerThreadsint处理消息监听器的线程数1useTcpNoDelayboolean是否禁用 Nagle 算法trueuseTlsboolean连接是否启用 TLSfalsetlsTrustCertsFilePathString受信 TLS 证书路径无tlsAllowInsecureConnectionboolean是否接受不受信证书falseconcurrentLookupRequestint每条 Broker 连接上并发的 lookup 请求数5000maxLookupRequestint每条 Broker 连接上允许的最大 lookup 请求数50000keepAliveIntervalSecondsint客户端-Broker 连接保活间隔303.2 Go 客户端官方推荐使用纯 Go 实现的pulsar-client-goCGo 客户端已标记为弃用见 client-libraries-cgo.md。安装与导入$ go get -u github.com/apache/pulsar-client-go/pulsarimport github.com/apache/pulsar-client-go/pulsar创建客户端import ( log time github.com/apache/pulsar-client-go/pulsar ) func main() { client, err : pulsar.NewClient(pulsar.ClientOptions{ URL: pulsar://localhost:6650, OperationTimeout: 30 * time.Second, ConnectionTimeout: 30 * time.Second, }) if err ! nil { log.Fatalf(Could not instantiate Pulsar client: %v, err) } defer client.Close() }多 Broker 时只需将URL改为逗号分隔列表即可。3.3 Python 客户端Python 客户端是 C 客户端的封装因此依赖安装比较简单# 基础安装 $ pip install pulsar-client2.8.0 # 附带 Avro 序列化 $ pip install pulsar-client[avro]2.8.0 # 附带 Functions 运行时 $ pip install pulsar-client[functions]2.8.0 # 全部可选组件 $ pip install pulsar-client[all]2.8.0Producer 示例向my-topic发送 10 条消息import pulsar client pulsar.Client(pulsar://localhost:6650) producer client.create_producer(my-topic) for i in range(10): producer.send((Hello-%d % i).encode(utf-8)) client.close()Consumer 示例订阅my-topic并逐条确认import pulsar client pulsar.Client(pulsar://localhost:6650) consumer client.subscribe(my-topic, my-subscription) while True: msg consumer.receive() try: print(Received message {} id{}.format(msg.data(), msg.message_id())) consumer.acknowledge(msg) except Exception: consumer.negative_acknowledge(msg) client.close()3.4 C 客户端C 客户端支持 Linux、MacOS 与 Windows。Linux 下既可以通过源码编译也可以直接安装预构建的 RPM/DEB 包。源码编译方式在仓库根目录执行# 安装依赖Debian/Ubuntu $ apt-get install cmake libssl-dev libcurl4-openssl-dev liblog4cxx-dev \ libprotobuf-dev protobuf-compiler libboost-all-dev google-mock libgtest-dev libjsoncpp-dev $ cd pulsar-client-cpp $ cmake . $ make编译完成后仓库lib目录下会出现libpulsar.so与libpulsar.aperf目录下则有perfProducer和perfConsumer性能测试工具。若安装 RPM/DEB 包则/usr/lib下会包含libpulsar.so、libpulsarnossl.so、libpulsar.a、libpulsarwithdeps.a四个库文件。从 2.1.0 版本起 Pulsar 即随发布附带预构建的 RPM 与 Debian 包依赖的详细版本清单可查看 pulsar-client-cpp/pkg/rpm 与 pulsar-client-cpp/pkg/deb 下的构建文件。3.5 Node.js 客户端Node.js 客户端基于 C 客户端安装前需先按 client-libraries-cpp.md#compilation 编译安装 C 客户端库且仅支持 Node.js 10.x 及以上依赖node-addon-api。版本兼容矩阵如下Node.js 客户端C 客户端1.0.02.3.0 或更高1.1.02.4.0 或更高1.2.02.5.0 或更高安装与创建客户端$ npm install pulsar-clientconst Pulsar require(pulsar-client); (async () { const client new Pulsar.Client({ serviceUrl: pulsar://localhost:6650, }); // ... 创建 producer / consumer await client.close(); })();3.6 C# 客户端DotPulsar通过 dotnet CLI 安装$ dotnet new console $ dotnet add package DotPulsarcsproj中会生成如下引用ItemGroup PackageReference IncludeDotPulsar Version2.0.1 / /ItemGroup创建客户端using DotPulsar; var client PulsarClient.Builder().Build();常用构建选项ServiceUrl默认pulsar://localhost:6650、RetryInterval操作或重连前的等待时间默认 3 秒。3.7 WebSocket 客户端WebSocket API 面向没有官方客户端语言的场景通过浏览器或任意 WebSocket 库即可发布与消费消息所有交换数据均为 JSON 格式。standalone 模式下 WebSocket 服务默认已启用非 standalone 模式有两种部署方式方式一内嵌在 Broker 中在 conf/broker.conf 设置webSocketServiceEnabledtrue方式二独立组件运行在 conf/websocket.conf 至少配置三个参数configurationMetadataStoreUrlzk1:2181,zk2:2181,zk3:2181 webServicePort8080 clusterNamemy-cluster随后用pulsar-daemon启动$ bin/pulsar-daemon start websocketWebSocket 提供三类端点路径中的:tenant/:namespace/:topic为占位符实际按持久化主题完整路径替换ws://broker-service-url:8080/ws/v2/producer/persistent/:tenant/:namespace/:topic ws://broker-service-url:8080/ws/v2/consumer/persistent/:tenant/:namespace/:topic/:subscription ws://broker-service-url:8080/ws/v2/reader/persistent/:tenant/:namespace/:topic启用 TLS 时 scheme 变为wss认证令牌通过查询参数token传递例如ws://broker-service-url:8080/ws/v2/producer/...?tokentoken。完整示例可查阅 client-libraries-websocket.md其中包含 Python 与 Node.js 的 WebSocket 客户端示例代码。四、客户端核心 APIProducer、Consumer 与 ReaderPulsar 客户端围绕三个核心抽象展开详见消息概念文档与客户端概念文档。4.1 Producer生产者以 Java 为例创建 Producer 并发送消息Producerbyte[] producer client.newProducer() .topic(my-topic) .create(); // 同步发送 producer.send(Hello Pulsar.getBytes()); // 异步发送 producer.sendAsync(Hello Pulsar.getBytes()) .thenAccept(msgId - System.out.println(Message ID: msgId));关键配置能力详见 client-libraries-java.md#producer消息路由分区主题下可配置MessageRouter决定消息发往哪个分区消息属性与顺序键通过newMessage()构造Message可设置 key、properties、eventTime 等消息分块chunking大消息自动切块发送、消费端重组通过enableChunking(true)开启压缩支持 LZ4、ZLib、Zstd、Snappy 等压缩算法仓库对应实现位于 pulsar-client-cpp/lib 的CompressionCodec*系列文件。4.2 Consumer消费者与订阅模式消费者通过**订阅subscription**绑定主题不同订阅类型决定消息在多个消费者之间的分发方式订阅类型行为适用场景Exclusive一个订阅只允许一个消费者严格有序、单消费者Shared消息在多个消费者间轮流分发无顺序保证高吞吐并行消费Failover主消费者失败后由备用消费者接管保序有序且需要故障切换Key_Shared按消息 key 将同一 key 的消息固定分发到同一消费者按 key 保序的并行消费Java 消费者示例Shared 订阅异步接收Consumerbyte[] consumer client.newConsumer() .topic(my-topic) .subscriptionName(my-subscription) .subscriptionType(SubscriptionType.Shared) .messageListener((consumer, msg) - { System.out.println(Received: new String(msg.getData())); consumer.acknowledgeAsync(msg); }) .subscribe();其他常用消费者能力批量接收batch receive一次接收多条消息减少 RTT负确认重投退避negativeAckRedeliveryBackoff控制消息重投延迟策略确认超时重投退避ackTimeoutRedeliveryBackoff配合ackTimeout使用多主题订阅subscribe时传入多个 topic或使用PatternMultiTopicsConsumer按正则匹配主题。4.3 Reader读取器接口手动管理游标与 Consumer 不同Reader 接口让应用手动管理游标位置连接主题时可选择三种起始位置主题中最早可用消息MessageId.earliest主题中最新可用消息MessageId.latest最早与最新之间的任意消息显式传入MessageId。Java 示例import org.apache.pulsar.client.api.Message; import org.apache.pulsar.client.api.MessageId; import org.apache.pulsar.client.api.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();从指定消息 ID 开始读byte[] msgIdBytes // 从持久化存储等处取回 MessageId id MessageId.fromByteArray(msgIdBytes); Readerbyte[] reader pulsarClient.newReader() .topic(topic) .startMessageId(id) .create();Reader 的实现原理源码结构可佐证它本质上是使用一个随机的、独占的非持久化订阅连接到主题。由此带来两个重要注意事项官方文档明确强调Reader 非持久化不会阻止主题中的数据被删除因此强烈建议配置数据保留策略cookbooks-retention-expiry.md否则未读消息可能被删除导致 Reader 跳消息Reader 的 backlog 指标仅用于观测落后程度不参与任何 backlog 配额计算。Reader 接口适合“精确回放”场景例如流处理系统基于 Pulsar 实现 exactly-once 语义时需要把主题回退到指定消息重新读取。五、客户端连接建立的源码级原理理解客户端如何连上 Broker有助于排查连接超时、鉴权失败等问题。官方客户端概念文档将客户端初始化划分为两个阶段主题归属查找lookup客户端向任一活跃 Broker 发送 HTTP lookup 请求Broker 通过缓存的ZooKeeper 元数据判断主题当前由哪个 Broker 服务若主题尚无人服务则尝试将其指派给负载最低的 Broker建立二进制协议连接拿到 Broker 地址后客户端创建或复用连接池中的TCP 连接并进行认证随后在连接上通过自定义二进制协议交换命令客户端发出创建 Producer/Consumer 的命令Broker 校验授权策略后予以响应。从 Java 客户端源码 pulsar-client/src/main/java/org/apache/pulsar/client/impl/PulsarClientImpl.java 可以看到PulsarClientImpl构造时会根据配置实例化两类查找服务if (conf.isUseTls()) { lookup new HttpLookupService(conf, this.eventLoopGroup); } else { lookup new BinaryProtoLookupService(this, conf.getServiceUrl(), ...); }即非 TLS 场景默认走二进制协议查找BinaryProtoLookupService对应 pulsar-client-cpp/lib/BinaryProtoLookupService.ccTLS 场景走 HTTP 查找HttpLookupService。同一文件中new ConnectionPool(conf, this.eventLoopGroup)表明客户端维护一个连接池可复用已有 TCP 连接避免每个 Producer/Consumer 都新建连接。断线重连机制无论何时 TCP 连接断开客户端都会立即重新发起上述两个阶段并以**指数退避exponential backoff**持续重试直到 Producer/Consumer 重建成功。仓库中的 Backoff.cc 与 Java 侧对应实现均体现了这一策略——这与开头所述“透明重连、带退避的重试”一脉相承。六、客户端安全能力各语言客户端均支持以下安全能力详见 security-overview.mdTLS 传输加密连接 URL 使用pulsarssl://scheme并在客户端配置useTls、tlsTrustCertsFilePath、tlsAllowInsecureConnection、tlsHostnameVerificationEnable等参数配置项及默认值见上文loadConf参数表认证插件Java 客户端通过authPluginClassNameauthParams配置支持 TLS、Athenz、OAuth2 等认证方式见 client-libraries-java.md#authenticationWebSocket 令牌通过 URL 查询参数token传递认证令牌。仓库对应的安全认证文档包括 security-tls-authentication.md、security-athenz.md、security-oauth2.md 等配置生产环境前建议逐一查阅。七、第三方客户端生态除了官方客户端社区还提供了多种语言的第三方客户端getting-started-clients.md#third-party-clients语言项目维护者许可证说明Gopulsar-client-goComcastApache 2.0原生 Go 客户端Gogo-pulsart2yApache 2.0Go 客户端实现HaskellsupernovaChatrouletteApache 2.0Haskell 原生 Pulsar 客户端ScalaneutronChatrouletteApache 2.0基于 Fs2 的纯函数式 Scala 客户端Scalapulsar4ssksamuelApache 2.0类型安全、响应式的 Scala 客户端Rustpulsar-rsWyyerd GroupApache 2.0基于 Future 的 Rust 绑定.NETpulsar-client-dotnetLanayxMITC#/F#/VB 原生 .NET 客户端Node.jspulsar-flexayeo-flex-orgMIT原生 Node.js 客户端如果你开发了新的 Pulsar 客户端官方欢迎通过提交 Pull Request 将项目补充到上述列表。注意第三方客户端的特性支持可能与官方客户端存在差异生产选型时建议核对官方 Client Features Matrix文档中引用自外部分享表格仓库内另有developing-binary-protocol.md可帮助你基于自定义二进制协议自研客户端。八、总结与最佳实践结合官方文档与仓库源码使用 Pulsar 客户端时可遵循以下要点选型官方推荐优先使用 Java、Go、Python、C、Node.js、C# 官方客户端无官方客户端的语言可通过 WebSocket API 接入对延迟与特性完整度要求高时核对特性矩阵。连接本地开发用pulsar://localhost:6650生产环境按需使用多 Broker 地址列表与pulsarssl://配置operationTimeoutMs、keepAliveIntervalSeconds等参数控制连接行为。API 选择需要自动游标管理、消息确认交给 Pulsar 时用 Consumer需要手动指定起始位置重放到指定消息时用 Reader且务必配套配置数据保留策略。可靠性客户端自带指数退避重连与连接池复用应用侧只需关注消息确认与负确认策略使用 Shared 订阅时可借助 Key_Shared 实现按 key 保序。深入排查遇到连接问题可对照 PulsarClientImpl.java 理解 lookup 与连接池路径或查阅二进制协议文档自研或调试客户端。延伸阅读仓库内文档客户端概念与连接流程消息概念订阅、确认与保留Java 客户端完整参数表WebSocket API 端点与示例数据保留与过期策略赞分享消息队列后端流处理【免费下载链接】pulsarApache Pulsar - distributed pub-sub messaging system项目地址https://gitcode.com/gh_mirrors/pulsar28/pulsar点击查看免费下载相关推荐Apache Pulsar Go客户端库指南Apache Pulsar Go客户端库指南 项目介绍 Apache Pulsar Go客户端库是专为Apache Pulsar设计的纯Go语言实现旨在提供对消息队列Apache Pulsar 客户端库完全指南官方多语言客户端与 WebSocket 接入方案Apache Pulsar 客户端库完全指南官方多语言客户端与 WebSocket 接入方案 Pulsar 提供 Java、Go、Python、C、Nod消息队列后端流处理Apache Pulsar Java 客户端 OAuth 2.0 认证插件实战指南从配置到源码原理Apache Pulsar Java 客户端 OAuth 2.0 认证插件实战指南从配置到源码原理 导读 本文围绕 Apache Pulsar Java 客户消息队列后端流处理创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
返回列表