Apache Pulsar 客户端库全景指南:从语言生态选型到源码级原理(Getting Started with Pulsar Clients)
2026/9/23 2:37:45 网站建设 项目流程
  • 消息队列
  • 后端
  • 流处理

【免费下载链接】pulsar

Apache 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.md
Goclient-libraries-go.md
Pythonclient-libraries-python.md
C++client-libraries-cpp.md
Node.jsclient-libraries-node.md
WebSocketclient-libraries-websocket.md
C#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.solibpulsar.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 变为pulsar+ssl,默认端口 6651。

各场景 URL 示例:

# 单机本地 pulsar://localhost:6650 # 多 Broker(逗号分隔) pulsar://localhost:6650,localhost:6651,localhost:6652 # 生产集群 pulsar://pulsar.us-west.example.com:6650 # 启用 TLS pulsar+ssl://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-client
org.apache.pulsar.client.admin管理 API(见 admin-api-overview.md)org.apache.pulsar:pulsar-client-admin
org.apache.pulsar.client.all同时包含以上两者且统一打 shade,避免重复类org.apache.pulsar:pulsar-client-all

Maven 依赖:

<properties> <pulsar.version>2.8.0</pulsar.version> </properties> <dependencies> <dependency> <groupId>org.apache.pulsar</groupId> <artifactId>pulsar-client</artifactId> <version>${pulsar.version}</version> </dependency> </dependencies>

Gradle 依赖:

def pulsarVersion = '2.8.0' dependencies { compile group: 'org.apache.pulsar', name: 'pulsar-client', version: pulsarVersion }

创建客户端并指定loadConf常用参数(完整参数表见 client-libraries-java.md):

PulsarClient client = PulsarClient.builder() .serviceUrl("pulsar://localhost:6650") .build();
参数类型说明默认值
serviceUrlStringPulsar 服务地址
operationTimeoutMslong操作超时时间30000
statsIntervalSecondslong统计信息输出间隔(>=1 秒)60
numIoThreadsint处理 Broker 连接的 IO 线程数1
numListenerThreadsint处理消息监听器的线程数1
useTcpNoDelayboolean是否禁用 Nagle 算法true
useTlsboolean连接是否启用 TLSfalse
tlsTrustCertsFilePathString受信 TLS 证书路径
tlsAllowInsecureConnectionboolean是否接受不受信证书false
concurrentLookupRequestint每条 Broker 连接上并发的 lookup 请求数5000
maxLookupRequestint每条 Broker 连接上允许的最大 lookup 请求数50000
keepAliveIntervalSecondsint客户端-Broker 连接保活间隔30

3.2 Go 客户端

官方推荐使用纯 Go 实现的pulsar-client-go(CGo 客户端已标记为弃用,见 client-libraries-cgo.md)。安装与导入:

$ go get -u "github.com/apache/pulsar-client-go/pulsar"
import "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-client==2.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.0'

Producer 示例(向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.solibpulsar.aperf目录下则有perfProducerperfConsumer性能测试工具。若安装 RPM/DEB 包,则/usr/lib下会包含libpulsar.solibpulsarnossl.solibpulsar.alibpulsarwithdeps.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-client
const 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 DotPulsar

csproj中会生成如下引用:

<ItemGroup> <PackageReference Include="DotPulsar" Version="2.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 设置:

webSocketServiceEnabled=true

方式二:独立组件运行,在 conf/websocket.conf 至少配置三个参数:

configurationMetadataStoreUrl=zk1:2181,zk2:2181,zk3:2181 webServicePort=8080 clusterName=my-cluster

随后用pulsar-daemon启动:

$ bin/pulsar-daemon start websocket

WebSocket 提供三类端点(路径中的: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/...?token=<token>。完整示例可查阅 client-libraries-websocket.md,其中包含 Python 与 Node.js 的 WebSocket 客户端示例代码。

四、客户端核心 API:Producer、Consumer 与 Reader

Pulsar 客户端围绕三个核心抽象展开,详见消息概念文档与客户端概念文档。

4.1 Producer(生产者)

以 Java 为例,创建 Producer 并发送消息:

Producer<byte[]> 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 订阅,异步接收):

Consumer<byte[]> 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 接口让应用手动管理游标位置,连接主题时可选择三种起始位置:

  1. 主题中最早可用消息(MessageId.earliest);
  2. 主题中最新可用消息(MessageId.latest);
  3. 最早与最新之间的任意消息(显式传入MessageId)。

Java 示例:

import org.apache.pulsar.client.api.Message; import org.apache.pulsar.client.api.MessageId; import org.apache.pulsar.client.api.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();

从指定消息 ID 开始读:

byte[] msgIdBytes = // 从持久化存储等处取回 MessageId id = MessageId.fromByteArray(msgIdBytes); Reader<byte[]> reader = pulsarClient.newReader() .topic(topic) .startMessageId(id) .create();

Reader 的实现原理(源码结构可佐证):它本质上是使用一个随机的、独占的非持久化订阅连接到主题。由此带来两个重要注意事项(官方文档明确强调):

  1. Reader 非持久化,不会阻止主题中的数据被删除,因此强烈建议配置数据保留策略(cookbooks-retention-expiry.md),否则未读消息可能被删除导致 Reader 跳消息;
  2. Reader 的 backlog 指标仅用于观测落后程度,不参与任何 backlog 配额计算

Reader 接口适合“精确回放”场景,例如流处理系统基于 Pulsar 实现 exactly-once 语义时,需要把主题回退到指定消息重新读取。

五、客户端连接建立的源码级原理

理解客户端如何连上 Broker,有助于排查连接超时、鉴权失败等问题。官方客户端概念文档将客户端初始化划分为两个阶段:

  1. 主题归属查找(lookup):客户端向任一活跃 Broker 发送 HTTP lookup 请求,Broker 通过(缓存的)ZooKeeper 元数据判断主题当前由哪个 Broker 服务;若主题尚无人服务,则尝试将其指派给负载最低的 Broker;
  2. 建立二进制协议连接:拿到 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.cc),TLS 场景走 HTTP 查找(HttpLookupService)。同一文件中new ConnectionPool(conf, this.eventLoopGroup)表明客户端维护一个连接池,可复用已有 TCP 连接,避免每个 Producer/Consumer 都新建连接。

断线重连机制:无论何时 TCP 连接断开,客户端都会立即重新发起上述两个阶段,并以**指数退避(exponential backoff)**持续重试,直到 Producer/Consumer 重建成功。仓库中的 Backoff.cc 与 Java 侧对应实现均体现了这一策略——这与开头所述“透明重连、带退避的重试”一脉相承。

六、客户端安全能力

各语言客户端均支持以下安全能力(详见 security-overview.md):

  • TLS 传输加密:连接 URL 使用pulsar+ssl://scheme,并在客户端配置useTlstlsTrustCertsFilePathtlsAllowInsecureConnectiontlsHostnameVerificationEnable等参数(配置项及默认值见上文loadConf参数表);
  • 认证插件:Java 客户端通过authPluginClassName+authParams配置,支持 TLS、Athenz、OAuth2 等认证方式,见 client-libraries-java.md#authentication;
  • WebSocket 令牌:通过 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 客户端时可遵循以下要点:

  1. 选型:官方推荐优先使用 Java、Go、Python、C++、Node.js、C# 官方客户端;无官方客户端的语言可通过 WebSocket API 接入;对延迟与特性完整度要求高时核对特性矩阵。
  2. 连接:本地开发用pulsar://localhost:6650,生产环境按需使用多 Broker 地址列表与pulsar+ssl://;配置operationTimeoutMskeepAliveIntervalSeconds等参数控制连接行为。
  3. API 选择:需要自动游标管理、消息确认交给 Pulsar 时用 Consumer;需要手动指定起始位置(重放到指定消息)时用 Reader,且务必配套配置数据保留策略。
  4. 可靠性:客户端自带指数退避重连与连接池复用,应用侧只需关注消息确认与负确认策略;使用 Shared 订阅时可借助 Key_Shared 实现按 key 保序。
  5. 深入排查:遇到连接问题可对照 PulsarClientImpl.java 理解 lookup 与连接池路径,或查阅二进制协议文档自研或调试客户端。

延伸阅读(仓库内文档):

  • 客户端概念与连接流程
  • 消息概念:订阅、确认与保留
  • Java 客户端完整参数表
  • WebSocket API 端点与示例
  • 数据保留与过期策略
  • 消息队列
  • 后端
  • 流处理

【免费下载链接】pulsar

Apache Pulsar - distributed pub-sub messaging system

项目地址:https://gitcode.com/gh_mirrors/pulsar28/pulsar
点击查看免费下载

创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考

需要专业的网站建设服务?

联系我们获取免费的网站建设咨询和方案报价,让我们帮助您实现业务目标

立即咨询