消息队列后端流处理【免费下载链接】pulsarApache Pulsar - distributed pub-sub messaging system项目地址https://gitcode.com/gh_mirrors/pulsar28/pulsar点击查看免费下载本指南以 Apache Pulsar 官方 C 客户端pulsar-client-cpp为对象完整讲解其在 Linux 与 macOS 上的依赖安装、从源码编译的全过程并基于仓库真实 API 与示例代码深入演示如何使用 C 客户端创建 Producer、Consumer 完成消息的发布与订阅以及如何通过 TLS 与认证插件安全地连接 Broker。读完本文你将能够独立编译出libpulsar库并写出可直接运行的 C 生产消费程序与带认证的安全客户端。概览Pulsar C 客户端Apache Pulsar 官方 C 客户端是与 Java、Python、Go 客户端并列的第一方语言客户端代码位于仓库根目录下的 pulsar-client-cpp 目录。整个目录结构清晰地区分了对外 API 与内部实现pulsar-client-cpp/include/pulsar对外公开的头文件即 Doxygen 文档的主体也是开发者编写应用时唯一需要 include 的接口层pulsar-client-cpp/lib内部实现ClientImpl、ProducerImpl、ConsumerImpl、Commands等pulsar-client-cpp/examples可直接运行的示例程序SampleProducer.cc、SampleConsumer.cc等pulsar-client-cpp/perf性能压测工具perfProducer与perfConsumer的源码pulsar-client-cpp/tests基于 Google Test 的单元测试覆盖消息、Topic 名解析、压缩、认证等多个维度。支持的平台与系统依赖平台支持Pulsar C 客户端已在macOS与Linux上经过完整测试这也是本指南编译章节覆盖的两个平台。从 CMakeLists.txt 可以看到项目要求cmake_minimum_required(VERSION 3.4)默认按 C11 标准编译set(CMAKE_CXX_STANDARD 11)。系统依赖清单使用 C 客户端前需要预先安装以下依赖依赖作用CMake构建系统版本需 ≥ 3.4Boost通用 C 库系统组件、regex、program_options 等Protocol Buffers 2.6消息协议序列化Log4CXX日志组件可选默认关闭见下文 CMake 选项libcurlHTTP 服务发现与认证请求Google Test单元测试编译测试时需要OpenSSL / ZLib / Zstd / SnappyTLS 加密与压缩编解码Zstd、Snappy 为可选需要说明的是CMakeLists.txt 中开放了多个编译选项实际构建所需的依赖会随选项变化BUILD_DYNAMIC_LIB默认 ON构建动态库BUILD_STATIC_LIB默认 ON构建静态库BUILD_TESTS默认 ON构建单元测试需要 Google TestBUILD_PYTHON_WRAPPER默认 ON构建 Python 封装需要 Boost.Python 与 Python 头文件BUILD_PERF_TOOLS默认 OFF构建perfProducer/perfConsumer压测工具USE_LOG4CXX默认 OFF启用 Log4CXX 日志后端LINK_STATIC默认 OFF静态链接全部依赖。如果你只需要客户端库而不需要测试与 Python 封装可在cmake阶段显式关闭对应选项下文会给出示例。从源码编译无论 Linux 还是 macOS第一步都是获取 Pulsar 源码仓库$ git clone https://gitcode.com/gh_mirrors/pulsar28/pulsar如果你已有 Pulsar 源码树直接进入pulsar-client-cpp子目录即可。Linux 编译步骤第一步安装全部依赖。在 Debian/Ubuntu 系列发行版上一条命令即可装齐$ apt-get install cmake libssl-dev libcurl4-openssl-dev liblog4cxx-dev \ libprotobuf-dev libboost-all-dev libgtest-dev libjsoncpp-dev第二步编译并安装 Google Test仅在需要运行单元测试时需要$ git clone https://github.com/google/googletest.git cd googletest $ sudo cmake . $ sudo make $ sudo cp *.a /usr/lib第三步编译 Pulsar C 客户端库。进入客户端源码目录执行 CMake 与 make$ cd pulsar-client-cpp $ cmake . $ make编译完成后产物会自动落到仓库对应子目录动态库libpulsar.so与静态库libpulsar.a位于pulsar-client-cpp/lib目录压测工具perfProducer与perfConsumer位于pulsar-client-cpp/perf目录构建时需开启BUILD_PERF_TOOLS。若只想快速产出客户端库、跳过测试与 Python 封装可这样配置$ cd pulsar-client-cpp $ cmake -DBUILD_TESTSOFF -DBUILD_PYTHON_WRAPPEROFF -DBUILD_PERF_TOOLSOFF . $ makemacOS 编译步骤第一步安装依赖。推荐使用 Homebrew# OpenSSL 安装 $ brew install openssl $ export OPENSSL_INCLUDE_DIR/usr/local/opt/openssl/include/ $ export OPENSSL_ROOT_DIR/usr/local/opt/openssl/ # Protocol Buffers 安装 $ brew tap homebrew/versions $ brew install protobuf260 $ brew install boost $ brew install log4cxx # Google Test 安装 $ git clone https://github.com/google/googletest.git $ cd googletest $ cmake . $ make install需要留意的是现代 macOS特别是 Apple Silicon上 OpenSSL 的 Homebrew 路径可能不同。源码在 CMakeLists.txt 中已经兼容了/usr/local/opt/openssl/与/opt/homebrew/opt/openssl/两种前缀因此export的两条变量可以按实际安装路径调整。第二步编译客户端库$ cd pulsar-client-cpp $ cmake . $ make核心 API 速览在动手写代码前先建立对客户端公共 API 的整体认知。公开头文件全部位于 pulsar-client-cpp/include/pulsarClient.h客户端入口负责创建 Producer / Consumer / Reader、执行subscribe、close等操作所有方法都提供同步与异步*Async两个版本ClientConfiguration.h客户端级配置如 TLS、认证、IO 线程数、操作超时等Producer.h 与 ProducerConfiguration.h消息发布端Consumer.h 与 ConsumerConfiguration.h消息消费端Message.h 与 MessageBuilder.h消息对象与消息构建Authentication.h认证插件体系TLS、Token、Athenz、OAuth2Result.h所有 API 统一返回的Result枚举如ResultOk、ResultTimeout、ResultProducerQueueIsFull等。客户端连接方式客户端通过serviceUrl连接集群。仓库示例SampleProducer.cc默认使用pulsar://localhost:6650连接本地单机。两种构造方式// 默认配置 Client client(pulsar://localhost:6650); // 自定义配置TLS、认证、线程数等 Client client(pulsarssl://my-broker.com:6651, config);常见的 serviceUrl 协议包括pulsar://明文 TCPpulsarssl://TLS 加密 TCPhttp:///https://HTTP 服务发现配合认证中的 HTTP 头部使用。客户端级关键配置ClientConfiguration.h 中定义的常用配置项及其默认值如下配置方法默认值说明setOperationTimeoutSeconds30 秒客户端操作subscribe、createProducer、close、unsubscribe的超时时间setIOThreads1IO 线程数用于网络读写setMessageListenerThreads1消息 listener 投递线程数超过 1 时不同 listener 可并行投递但同一 listener 始终固定在同一线程setConcurrentLookupRequest50000每条 Broker 连接上允许的并发 lookup 请求数防止 Broker 过载在单客户端需生产/订阅上千个 Topic 时可调大setStatsIntervalInSeconds600统计信息打印间隔秒设为 0 表示关闭统计setPartititionsUpdateInterval60分区数更新间隔秒为 0 时不主动探测分区数变化setConnectionTimeout10000建立 Broker 连接的等待时间毫秒setMemoryLimit0不限制客户端实例可分配的内存上限字节可用于控制积压缓冲占用setUseTls/setTlsTrustCertsFilePath/setTlsAllowInsecureConnection/setValidateHostName见下节TLS 相关配置setAuthDisabled设置认证插件消费者Consumer实战原文档给出的 Consumer 核心示例在仓库的 SampleConsumer.cc 中有完全对应的可运行版本。完整代码Client client(pulsar://localhost:6650); Consumer consumer; Result result client.subscribe(persistent://sample/standalone/ns1/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)使用默认的 ConsumerConfiguration 创建消费者。Topic 使用完整的persistent://tenant/namespace/topic三段式域名。订阅名subscription是逻辑隔离单位同一订阅下多个消费者可分摊消息Shared 模式或主备切换Failover 模式。接收消息consumer.receive(msg)是阻塞调用无消息时一直等待也支持带超时的重载consumer.receive(msg, timeoutMs)超时返回ResultTimeout。同步接收与消息 listener 不能混用若在配置中设置了 listener调用receive会返回ResultInvalidConfiguration见 Consumer.h。确认消息consumer.acknowledge(msg)通知 Broker 该消息已成功处理避免重新投递。Consumer还提供acknowledgeCumulative累积确认Shared 模式下不可用与negativeAcknowledge否定确认触发延迟重投延迟可通过ConsumerConfiguration::setNegativeAckRedeliveryDelayMs配置。收尾client.close()会等待所有在途写请求持久化后有序关闭全部 Producer、Consumer 与 ReaderClient.h而shutdown()则立即释放资源、不等待在途操作。从源码结构看subscribe底层由ClientImpl创建ConsumerImpl对应 pulsar-client-cpp/lib/ConsumerImpl.cc并支持subscribeWithRegex以正则一次订阅同一 namespace 下的多个 Topic对应行为在 pulsar-client-cpp/tests/ConsumerTest.cc 等测试中均有覆盖。生产者Producer实战仓库中 SampleProducer.cc 展示了与原文档一致的生产端用法。原文档中的批量发布版本Client client(pulsar://localhost:6650); Producer producer; Result result client.createProducer(persistent://sample/standalone/ns1/my-topic, producer); if (result ! ResultOk) { LOG_ERROR(Error creating producer: result); return -1; } // Publish 10 messages to the topic for(int i0;i10;i){ Message msg MessageBuilder().setContent(my-message).build(); Result res producer.send(msg); LOG_INFO(Message sent: res); } client.close();代码走查与要点创建生产者client.createProducer(topic, producer)返回ResultOk即表示与 Broker 建立连接成功可用producer.getProducerName()获取 Broker 分配的 producer 名称。构建消息MessageBuilder采用链式调用。除setContent外它还支持见 MessageBuilder.hsetProperty(name, value)/setProperties(map)附加键值对属性setPartitionKey(key)设置分区键分区 Topic 会根据该键的哈希路由setOrderingKey(key)为 Key_Shared 订阅设置排序键setDeliverAfter(milliseconds)/setDeliverAt(timestamp)延迟投递 / 定时投递setEventTimestamp(timestamp)设置事件时间戳setSequenceId(id)设置自定义序号要求 0且严格递增可用于去重setReplicationClusters(clusters)/disableReplication(flag)覆盖 namespace 级跨集群复制策略setAllocatedContent(ptr, size)零拷贝接入调用方已分配的内存缓冲区。同步发送producer.send(msg)会阻塞直到 Broker 接受并持久化消息发送失败时客户端库会自动重连并切换到其他 Broker 重试。send(msg, messageId)重载还能返回 Broker 分配的MessageId。可能返回的错误码包括ResultTimeout超过ProducerConfiguration::getSendTimeout、ResultProducerQueueIsFull队列满且未开启队列满时阻塞、ResultMessageTooBig、ResultAlreadyClosed、ResultCryptoError、ResultInvalidMessage见 Producer.h。异步发送与冲刷高吞吐场景建议使用sendAsync(msg, callback)批量/缓冲模式下可调用flush()/flushAsync()强制将缓冲消息全部持久化。getLastSequenceId()可获取已确认的最后一条消息序号。关闭producer.close()等待在途写请求持久化后释放资源。认证与 TLS 加密连接原文档给出了启用 TLS 并加载libauthtls.so认证插件的完整示例ClientConfiguration config ClientConfiguration(); config.setUseTls(true); std::string certfile /path/to/cacert.pem; ParamMap params; params[tlsCertFile] /path/to/client-cert.pem; params[tlsKeyFile] /path/to/client-key.pem; config.setTlsTrustCertsFilePath(certfile); config.setTlsAllowInsecureConnection(false); AuthenticationPtr auth pulsar::AuthFactory::create(/path/to/libauthtls.so, params); config.setAuth(auth); Client client(pulsarssl://my-broker.com:6651,config);结合 Authentication.h 与 ClientConfiguration.h 的源码可深入理解这套机制TLS 连接配置配置方法默认值说明setUseTls(true)false对连接启用 TLS 加密同时 serviceUrl 需使用pulsarssl://协议setTlsTrustCertsFilePath(path)空信任的 CA 证书cacert.pem路径用于验证 Broker 证书链setTlsAllowInsecureConnection(false)false是否允许接受来自 Broker 的未受信任证书生产环境应保持falsesetValidateHostName(bool)false是否按 RFC 2818 校验服务端证书 CN/SAN 与连接主机名一致建议开启认证插件体系AuthFactory::create(pluginNameOrDynamicLibPath, params)支持两种形式传入内置插件名可直接创建内置认证tls→AuthTlstoken或org.apache.pulsar.client.impl.auth.AuthenticationToken→AuthTokenathenz→AuthAthenzoauth2token→AuthOauth2。传入动态库路径如示例中的/path/to/libauthtls.so运行时加载外部认证插件params作为键值参数传入插件的create方法。这正是示例代码采用的动态插件路线。除动态加载插件外仓库还提供了纯静态的创建方式无需.so文件// TLS 双向认证直接指定客户端证书与私钥 AuthenticationPtr auth AuthTls::create(/path/to/client-cert.pem, /path/to/client-key.pem); // Token 认证直接传入 token 字符串 AuthenticationPtr auth AuthToken::createWithToken(my-auth-token); // Token 认证从环境变量或文件读取通过参数串 AuthenticationPtr auth AuthToken::create(file:///path/to/token-file);其中AuthToken::create(ParamMap)支持三种 token 来源参数键为token时值为 token 本身键为file时值为存放 token 的文件路径键为env时值为保存 token 的环境变量名Authentication.h。此外仓库还内置AuthAthenz需tenantDomain、tenantService、providerDomain、privateKey、ztsUrl五个参数与AuthOauth2client_credentials 流程核心参数为issuer_url、private_key、audience、client_id、client_secret对应的 C 实现位于 pulsar-client-cpp/lib/auth 目录。值得说明的是TLS 传输加密与客户端身份认证是两个独立维度。setUseTls系列配置解决链路是否加密、是否信任对端而setAuth解决我是谁、Broker 是否放行。两者配合使用才能构成生产环境的完整安全链路。代码格式化修改过客户端源码后可用项目自带的格式化目标统一代码风格make format对应的实现位于 CMakeLists.txt它调用python build-support/run_clang_format.py对lib、perf、examples、tests、include、python/src、wireshark目录下的代码执行自动格式化排除规则见build-support/clang_format_exclusions.txt。CI 环境则使用make check-format校验格式是否合规。使用自动格式化需要安装clang-format 5.0。进一步探索可运行示例完整的 Producer / Consumer / Reader 示例位于 pulsar-client-cpp/examples包含SampleProducer.cc、SampleConsumer.cc、SampleReaderCApi.c、SampleConsumerListener.cc基于 listener 的异步消费等C API面向纯 C 调用方的封装头文件在 pulsar-client-cpp/include/pulsar/c性能压测开启BUILD_PERF_TOOLSON后构建的perfProducer/perfConsumer源码在 pulsar-client-cpp/perf单元测试测试用例位于 pulsar-client-cpp/tests覆盖 Producer/Consumer 全流程、批量消息、压缩编解码、认证插件等是理解 API 语义的最佳辅助资料Python 封装基于同一 C 核心构建的 Python 客户端源码位于 pulsar-client-cpp/python。常见问题Q1编译报错找不到 Google Test如果不需要运行单元测试用-DBUILD_TESTSOFF关闭测试构建即可若需要则按上文 Linux 步骤预先编译安装 googletest 到系统路径。Q2make format提示找不到 clang-format格式化目标依赖 clang-format 5.0请先安装对应版本并确保其在PATH中。Q3连接远程集群一直超时先确认 serviceUrl 协议pulsar:///pulsarssl://与 Broker 监听端口默认 6650 明文、6651 TLS是否匹配其次检查 ClientConfiguration.h 中的setConnectionTimeout毫秒与setOperationTimeoutSeconds秒是否过小。Q4开启 TLS 后握手失败依次核对三件事setUseTls(true)是否设置、serviceUrl 是否为pulsarssl://、setTlsTrustCertsFilePath指向的 CA 证书是否能验证 Broker 证书链若为测试环境可临时将setTlsAllowInsecureConnection(true)生产环境必须保持false并建议开启setValidateHostName(true)。赞分享消息队列后端流处理【免费下载链接】pulsarApache Pulsar - distributed pub-sub messaging system项目地址https://gitcode.com/gh_mirrors/pulsar28/pulsar点击查看免费下载相关推荐Apache Pulsar Go 客户端pulsar-client-go使用指南生产、消费、读取与安全认证实战Apache Pulsar Go 客户端pulsar client go使用指南生产、消费、读取与安全认证实战 Apache Pulsar 官方为 Go消息队列后端流处理Apache Pulsar C 客户端开发指南构建安装、生产者与消费者实战Apache Pulsar C 客户端开发指南构建安装、生产者与消费者实战 本指南基于当前仓库的官方文档 client libraries cpp.md消息队列后端流处理Apache Pulsar C 客户端实战指南三平台编译安装、生产者/消费者开发与 Schema 配置Apache Pulsar C 客户端实战指南三平台编译安装、生产者/消费者开发与 Schema 配置 Apache Pulsar 提供了原生 C 客消息队列后端流处理创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考