- 消息队列
- 后端
- 流处理
【免费下载链接】pulsar
Apache Pulsar - distributed pub-sub messaging system
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 明确列出了以下官方支持的客户端库:
| 语言 | 官方文档 |
|---|---|
| Java | client-libraries-java.md |
| Go | client-libraries-go.md |
| Python | client-libraries-python.md |
| C++ | client-libraries-cpp.md |
| Node.js | client-libraries-node.md |
| WebSocket | client-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.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 变为
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.api | Producer/Consumer API | org.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();| 参数 | 类型 | 说明 | 默认值 |
|---|---|---|---|
serviceUrl | String | Pulsar 服务地址 | 无 |
operationTimeoutMs | long | 操作超时时间 | 30000 |
statsIntervalSeconds | long | 统计信息输出间隔(>=1 秒) | 60 |
numIoThreads | int | 处理 Broker 连接的 IO 线程数 | 1 |
numListenerThreads | int | 处理消息监听器的线程数 | 1 |
useTcpNoDelay | boolean | 是否禁用 Nagle 算法 | true |
useTls | boolean | 连接是否启用 TLS | false |
tlsTrustCertsFilePath | String | 受信 TLS 证书路径 | 无 |
tlsAllowInsecureConnection | boolean | 是否接受不受信证书 | false |
concurrentLookupRequest | int | 每条 Broker 连接上并发的 lookup 请求数 | 5000 |
maxLookupRequest | int | 每条 Broker 连接上允许的最大 lookup 请求数 | 50000 |
keepAliveIntervalSeconds | int | 客户端-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.so与libpulsar.a,perf目录下则有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.0 | 2.3.0 或更高 |
| 1.1.0 | 2.4.0 或更高 |
| 1.2.0 | 2.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 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 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/...?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 接口让应用手动管理游标位置,连接主题时可选择三种起始位置:
- 主题中最早可用消息(
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; // 从最早消息开始读 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 的实现原理(源码结构可佐证):它本质上是使用一个随机的、独占的非持久化订阅连接到主题。由此带来两个重要注意事项(官方文档明确强调):
- 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.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,并在客户端配置useTls、tlsTrustCertsFilePath、tlsAllowInsecureConnection、tlsHostnameVerificationEnable等参数(配置项及默认值见上文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):
| 语言 | 项目 | 维护者 | 许可证 | 说明 |
|---|---|---|---|---|
| Go | pulsar-client-go | Comcast | Apache 2.0 | 原生 Go 客户端 |
| Go | go-pulsar | t2y | Apache 2.0 | Go 客户端实现 |
| Haskell | supernova | Chatroulette | Apache 2.0 | Haskell 原生 Pulsar 客户端 |
| Scala | neutron | Chatroulette | Apache 2.0 | 基于 Fs2 的纯函数式 Scala 客户端 |
| Scala | pulsar4s | sksamuel | Apache 2.0 | 类型安全、响应式的 Scala 客户端 |
| Rust | pulsar-rs | Wyyerd Group | Apache 2.0 | 基于 Future 的 Rust 绑定 |
| .NET | pulsar-client-dotnet | Lanayx | MIT | C#/F#/VB 原生 .NET 客户端 |
| Node.js | pulsar-flex | ayeo-flex-org | MIT | 原生 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 地址列表与pulsar+ssl://;配置operationTimeoutMs、keepAliveIntervalSeconds等参数控制连接行为。 - API 选择:需要自动游标管理、消息确认交给 Pulsar 时用 Consumer;需要手动指定起始位置(重放到指定消息)时用 Reader,且务必配套配置数据保留策略。
- 可靠性:客户端自带指数退避重连与连接池复用,应用侧只需关注消息确认与负确认策略;使用 Shared 订阅时可借助 Key_Shared 实现按 key 保序。
- 深入排查:遇到连接问题可对照 PulsarClientImpl.java 理解 lookup 与连接池路径,或查阅二进制协议文档自研或调试客户端。
延伸阅读(仓库内文档):
- 客户端概念与连接流程
- 消息概念:订阅、确认与保留
- Java 客户端完整参数表
- WebSocket API 端点与示例
- 数据保留与过期策略
- 消息队列
- 后端
- 流处理
【免费下载链接】pulsar
Apache Pulsar - distributed pub-sub messaging system
相关推荐
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),仅供参考