- 消息队列
- 后端
- 流处理
【免费下载链接】pulsar
Apache Pulsar - distributed pub-sub messaging system
导读
本文聚焦 Apache Pulsar 官方 C++ 客户端库(libpulsar),完整覆盖在 Linux 与 macOS 上通过预编译包或源码构建的方式安装客户端、理解pulsar://与pulsar+ssl://协议 URL 的连接规则,以及基于Client、Consumer、Producer三大核心类编写消息生产、消费与 TLS 认证程序的实战方法。文中所有命令、配置与 API 说明均以当前仓库 pulsar-client-cpp 目录下的源码、打包脚本与头文件为准,读者可据此在真实环境中直接落地使用。
支持的平台
Pulsar C++ 客户端已在macOS与Linux两大平台上完成测试验证。客户端库本体位于仓库的 pulsar-client-cpp 目录,公共头文件集中在 pulsar-client-cpp/include/pulsar,既提供 C++ API,也通过 pulsar-client-cpp/include/pulsar/c 提供 C API,可被 C/C++ 程序直接引用。
Linux 安装
自 Pulsar2.1.0版本起,官方随发布版本提供预编译的 RPM 与 Debian 包,用户可以直接下载安装,无需自行编译。
RPM 包安装
RPM 发布包含三个包:
| 包 | 内容 |
|---|---|
| pulsar-client | 动态库libpulsar.so |
| pulsar-client-devel | 静态库libpulsar.a以及 C++ 与 C 头文件 |
| pulsar-client-debuginfo | libpulsar.so的调试符号 |
下载对应架构的 RPM 包后,在包所在目录执行:
$ rpm -ivh apache-pulsar-client*.rpmDEB 包安装
Debian/Ubuntu 系列发行版提供两个包:
| 包 | 内容 |
|---|---|
| pulsar-client | 动态库libpulsar.so |
| pulsar-client-dev | 静态库libpulsar.a以及 C++ 与 C 头文件 |
下载 DEB 包后安装:
$ dpkg -i apache-pulsar-client*.deb(原版文档中的apt-install命令并非标准工具,实际安装 DEB 包请使用dpkg -i;如需自动解析依赖,可先apt-get update后再配合apt-get install -f修复依赖关系。)
包内容与目录约定
以 RPM 为例,仓库中的打包规格文件 pulsar-client-cpp/pkg/rpm/SPECS/pulsar-client.spec 明确了安装路径与文件清单:
- 动态库安装为
/usr/lib/libpulsar.so.<version>,并创建符号链接/usr/lib/libpulsar.so,同时提供不依赖 OpenSSL 的libpulsarnossl.so变体; - 开发包安装
/usr/lib/libpulsar.a、/usr/lib/libpulsarwithdeps.a(将第三方依赖一并静态链入的版本)以及/usr/include/pulsar头文件目录; - 文档与 LICENSE 文件安装至
/usr/share/doc/下。
也就是说,安装 devel/dev 包后即可在代码中直接#include <pulsar/Client.h>并链接-lpulsar。
从源码构建 RPM / Debian 包
如果你希望基于最新 master 分支自行产出安装包,仓库提供了现成的 Docker 构建脚本,所有命令均在 Pulsar 仓库根目录执行。
先构建 Java 模块
C++ 库的打包过程会引用版本信息与部分生成资源,因此需要先构建 Java 模块:
mvn install -DskipTests构建 RPM
pulsar-client-cpp/pkg/rpm/docker-build-rpm.sh该脚本会拉取apachepulsar/pulsar-build:centos-7构建镜像,将仓库根目录挂载进容器后执行容器内的构建(详见 docker-build-rpm.sh)。构建产物 RPM 文件位于:
pulsar-client-cpp/pkg/rpm/RPMS/x86_64/包内产物与前面表格一致:pulsar-client(共享库)、pulsar-client-devel(静态库与头文件)、pulsar-client-debuginfo(调试符号)。
构建 Debian 包
pulsar-client-cpp/pkg/deb/docker-build-deb.sh对应地,该脚本使用apachepulsar/pulsar-build:debian-9镜像(见 docker-build-deb.sh),Debian 包产出目录为:
pulsar-client-cpp/pkg/deb/BUILD/DEB/构建细节:静态链接与多目标产物
从 RPM 的 spec 文件(pulsar-client.spec)可以看到实际构建命令:
cmake . -DBUILD_TESTS=OFF -DLINK_STATIC=ON -DBUILD_PYTHON_WRAPPER=OFF make pulsarShared pulsarSharedNossl pulsarStatic pulsarStaticWithDeps -j 3这段构建逻辑说明了几点关键事实:
LINK_STATIC=ON使能静态链接依赖,因此libpulsar.a/libpulsarwithdeps.a内已包含 Boost、Protobuf、cURL 等依赖的代码,用户链接时无需再单独指定;- 同时产出
pulsarShared(带 SSL 的动态库)、pulsarSharedNossl(不带 SSL 的动态库)、pulsarStatic(静态库)与pulsarStaticWithDeps(含依赖的静态库)四种目标; - 打包时排除测试(
BUILD_TESTS=OFF)与 Python wrapper(BUILD_PYTHON_WRAPPER=OFF)。
如果不想走 Docker,也可以直接在本地运行 build-rpm.sh / build-deb.sh,但需要自行保证本机具备完整构建依赖。
macOS 安装
在 macOS 上,Pulsar 发布版本通过 Homebrew 提供。可以直接安装:
brew install libpulsar安装完成后,库文件(libpulsar.dylib、libpulsar.a)与include/pulsar头文件会一并就位。仓库中保留了对应的 Homebrew 配方 homebrew/libpulsar.rb,从中可以看出它依赖cmake、openssl、boost、jsoncpp、protobuf@2.6等构建组件,并通过cmake . -DBUILD_TESTS=OFF -DLINK_STATIC=ON与make pulsarShared pulsarStatic完成编译安装,与 Linux 上的构建目标一致。
Connection URLs(连接地址规则)
使用客户端连接 Pulsar 时,必须指定一个 Pulsar 协议 URL:
- Pulsar 协议 URL 与具体集群绑定;
- 使用
pulsarURI 协议方案; - 默认端口为6650。
本地连接的示例:
pulsar://localhost:6650生产集群的典型地址:
pulsar://pulsar.us-west.example.com:6650启用 TLS 加密后,协议方案变为pulsar+ssl,端口相应改为6651:
pulsar+ssl://pulsar.us-west.example.com:6651从源码结构看,客户端通过Client构造函数的serviceUrl参数解析上述地址,例如 Client.h 中定义的两个构造重载(默认配置与自定义配置)都接收该 URL 字符串;底层由ClientImpl负责与 broker 建立连接。
Consumer:订阅并消费消息
下面的完整示例创建一个连接本地 broker 的客户端,订阅my-topic并以订阅名my-subscribtion-name持续接收、打印并确认消息:
#include <pulsar/Client.h> #include <iostream> using namespace pulsar; int main() { Client client("pulsar://localhost:6650"); Consumer consumer; Result result = client.subscribe("my-topic", "my-subscribtion-name", consumer); if (result != ResultOk) { LOG_ERROR("Failed to subscribe: " << result); return -1; } Message msg; while (true) { consumer.receive(msg); LOG_INFO("Received: " << msg << " with payload '" << msg.getDataAsString() << "'"); consumer.acknowledge(msg); } client.close(); }要点:
client.subscribe(topic, subscriptionName, consumer)同步订阅成功后,consumer即被填充为可用实例;订阅失败时返回Result错误码,可用result != ResultOk判断;consumer.receive(msg)阻塞式接收消息,msg.getDataAsString()取出 payload 文本;- 每条消息处理完成后务必调用
consumer.acknowledge(msg)进行确认,broker 据此推进消费位点、避免消息重复投递; - 除了单 topic 订阅,Client.h 还提供了多 topic 订阅(
subscribe(const std::vector<std::string>&, ...))与正则订阅(subscribeWithRegex),以及对应的subscribeAsync异步变体,适合需要一次性消费多个 topic 的场景。
Producer:向 topic 发布消息
生产者示例创建一个指向my-topic的生产者,循环发送 10 条消息:
#include <pulsar/Client.h> #include <iostream> using namespace pulsar; int main() { Client client("pulsar://localhost:6650"); Producer producer; Result result = client.createProducer("my-topic", producer); if (result != ResultOk) { LOG_ERROR("Error creating producer: " << result); return -1; } // Publish 10 messages to the topic for (int i = 0; i < 10; i++) { Message msg = MessageBuilder().setContent("my-message").build(); Result res = producer.send(msg); LOG_INFO("Message sent: " << res); } client.close(); }要点:
MessageBuilder().setContent(...).build()是构造消息的标准方式,也可通过setProperties、setPartitionKey等进一步定制(对应头文件 MessageBuilder.h);producer.send(msg)为同步发送,返回Result表示本次发送成功与否;- 需要更高吞吐时,
Client提供createProducerAsync异步创建生产者的回调接口(见 Client.h),生产端还有sendAsync可避免逐个等待 broker 确认; client.close()会等待所有挂起的写请求持久化完成后再释放资源;若需立即释放,可使用shutdown()。
Authentication:TLS 双向认证
当 broker 开启 TLS 与双向认证(mTLS)时,通过ClientConfiguration配置证书与认证插件:
#include <pulsar/Client.h> using namespace pulsar; int main() { ClientConfiguration config = ClientConfiguration(); config.setUseTls(true); config.setTlsTrustCertsFilePath("/path/to/cacert.pem"); config.setTlsAllowInsecureConnection(false); config.setAuth(pulsar::AuthTls::create( "/path/to/client-cert.pem", "/path/to/client-key.pem")); Client client("pulsar+ssl://my-broker.com:6651", config); }配置项说明:
setUseTls(true):启用 TLS 加密,此时连接地址必须使用pulsar+ssl://方案;setTlsTrustCertsFilePath:指定用于校验 broker 证书的 CA 证书(cacert.pem)路径;setTlsAllowInsecureConnection(false):禁止接受 broker 端未受信任的证书;setAuth(pulsar::AuthTls::create(clientCert, clientKey)):使用 TLS 客户端证书/私钥完成客户端身份认证,这是 Pulsar TLS 双向认证的标准配置方式;- 相关接口定义见 ClientConfiguration.h,同头文件还提供了
setValidateHostName(bool)用于开启基于 RFC 2818 的主机名(CN/SAN)校验。
客户端核心 API 与配置深入
Client 对象
Client.h 是 C++ 客户端的门面,核心能力包括:
| 能力 | 方法 |
|---|---|
| 创建生产者 | createProducer/createProducerAsync |
| 订阅消费 | subscribe/subscribeAsync(单 topic、多 topic、正则) |
| 创建 Reader | createReader/createReaderAsync(按消息 ID 定位,支持MessageId::earliest、MessageId::latest或指定位置) |
| 分区发现 | getPartitionsForTopic/getPartitionsForTopicAsync(自 2.3.0 起,返回 topic 的全部分区列表) |
| 关闭释放 | close/closeAsync(有序关闭并等待写请求持久化)、shutdown(立即释放) |
| 状态查询 | getNumberOfProducers/getNumberOfConsumers |
此外,Reader 提供不依赖订阅的低层读取能力,适合需要手动定位消息位置的场景(如重放指定消息 ID 之后的数据),但只能作用于非分区 topic。
ClientConfiguration 常用参数
结合 ClientConfiguration.h,客户端级配置的核心参数如下:
| 配置项 | 默认值 | 说明 |
|---|---|---|
setOperationTimeoutSeconds | 30 秒 | 订阅、创建 producer、关闭、取消订阅等客户端操作的超时 |
setIOThreads | 1 | 客户端使用的 IO 线程数 |
setMessageListenerThreads | 1 | 消息 listener 投递线程数;多线程时不同 listener 分派到不同线程,但单个 listener 始终固定同一线程 |
setConcurrentLookupRequest | 50000 | 每条 broker 连接上允许的并发 lookup 请求数,避免 broker 过载;在单客户端需要创建/订阅数千 topic 时再调高 |
setConnectionTimeout | 10000 ms | 建立 broker 连接的等待超时,超时后放弃该次连接尝试 |
setStatsIntervalInSeconds | 600 秒 | 统计信息打印与重置间隔,设为 0 表示关闭统计 |
setPartititionsUpdateInterval | 60 秒 | 分区 topic 元数据(分区数)的刷新间隔;分区扩容后客户端自动为新增分区创建 producer/consumer |
setMemoryLimit | 0(不限制) | 客户端实例允许分配的内存上限(字节),用于控制内存占用 |
setListenerName | - | 指定 broker 返回的advertisedListener对应的 listener 名称 |
小结
Apache Pulsar C++ 客户端提供了一条从安装到上线的完整路径:Linux 上可直接使用 2.1.0 起发布的 RPM/DEB 预编译包(rpm -ivh/dpkg -i),也可借助仓库自带的 Docker 脚本在 pulsar-client-cpp/pkg/rpm 与 pulsar-client-cpp/pkg/deb 产出包含静态链接依赖的安装包;macOS 用户通过 Homebrew 的libpulsar公式一键安装。连接层牢记pulsar://(6650)与pulsar+ssl://(6651)两种 URL 方案,业务侧则围绕Client、Producer、Consumer三大对象编写同步或异步的生产消费逻辑,并通过ClientConfiguration精确控制超时、线程数、内存上限与 TLS 认证参数。如需在项目里进一步研究源码细节,可重点阅读 Client.h、ClientConfiguration.h 与客户端测试目录 pulsar-client-cpp/tests 中的端到端用例。
- 消息队列
- 后端
- 流处理
【免费下载链接】pulsar
Apache Pulsar - distributed pub-sub messaging system
相关推荐
Apache Pulsar C++ 客户端开发指南:构建安装、生产者与消费者实战
Apache Pulsar C++ 客户端开发指南:构建安装、生产者与消费者实战 本指南基于当前仓库的官方文档( client libraries cpp.md
消息队列后端流处理Apache Pulsar C++ 客户端完全指南:安装、构建、连接与收发消息实战
Apache Pulsar C++ 客户端完全指南:安装、构建、连接与收发消息实战 本指南系统讲解 Apache Pulsar C++ 客户端( libpuls
消息队列后端流处理Apache Pulsar C++ 客户端完整指南:安装、构建与消息收发实战
Apache Pulsar C++ 客户端完整指南:安装、构建与消息收发实战 Apache Pulsar 官方提供基于 C++ 编写的客户端库( libpuls
消息队列后端流处理
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考