首页/新闻资讯/正文详情

Apache Pulsar C 客户端(DotPulsar)完整使用指南:安装、生产者/消费者/Reader 开发与状态监控

发布时间:2026/9/23 16:55:25 来源:云帆数科 栏目:资讯中心
Apache Pulsar C 客户端(DotPulsar)完整使用指南:安装、生产者/消费者/Reader 开发与状态监控
Apache Pulsar C# 客户端DotPulsar完整使用指南安装、生产者/消费者/Reader 开发与状态监控【免费下载链接】pulsarApache Pulsar - distributed pub-sub messaging system项目地址: https://gitcode.com/gh_mirrors/pulsar28/pulsarApache Pulsar 官方为 .NET / C# 开发者提供了基于 DotPulsar 的 C# 客户端库本文以 site2/website-next/docs/client-libraries-dotnet.md 为核心系统讲解如何在 .NET Core 项目中安装并创建 PulsarClient、Producer、Consumer 与 Reader覆盖消息发送/接收/确认、加密策略、TLS 与 JWT 认证以及基于状态机的事件驱动监控方案。读完本文你将能够用 C# 完整接入 Apache Pulsar 集群编写可运行的发布订阅程序并像官方文档示例一样对客户端生命周期进行健壮的状态监控。背景说明C# 客户端由官方社区贡献的 DotPulsar 中“Contribute DotPulsar to Apache Pulsar”的条目。当前仓库各语言客户端 API 语义一致concepts-clients.md 概述了所有官方客户端共享的“查找主题 → 建立 TCP 连接 → 认证 → 创建生产者/消费者”的建立流程。安装与项目准备前置条件使用 C# 客户端前需要先安装 .NET Core SDK它提供了dotnet命令行工具。从 Visual Studio 2017 开始dotnet CLI 会随任何 .NET Core 相关的工作负载自动安装因此也可以直接在 VS 环境中操作。安装步骤创建项目文件夹并打开终端切换到该目录。初始化控制台项目dotnet new console使用dotnet run运行一次验证应用已正确创建。添加 DotPulsar NuGet 包dotnet add package DotPulsar命令执行完成后打开.csproj文件即可看到自动加入的包引用官方文档示例中的版本为 0.11.0ItemGroup PackageReference IncludeDotPulsar Version0.11.0 / /ItemGroup后续若需升级版本只需修改该Version或重新执行dotnet add package DotPulsar即可。客户端PulsarClient配置PulsarClient 是 C# 应用与 Pulsar 集群通信的入口负责管理底层连接、自动重连与资源生命周期。其所有方法都是线程安全的因此可以在多线程场景中共享同一个 client 实例。创建客户端连接本地集群默认地址pulsar://localhost:6650的最简写法var client PulsarClient.Builder().Build();使用 Builder 时可以指定以下核心选项Option说明默认值ServiceUrl设置 Pulsar 集群的服务地址pulsar://localhost:6650RetryInterval设置操作或重连前的等待时间3s结合 concepts-clients.md 的客户端建立流程可以更准确地理解 ServiceUrl 的作用应用创建 producer/consumer 前客户端会先通过 HTTP 查找请求确定 topic 归属的 broker再建立 TCP 连接并完成认证最后在连接上创建生产者/消费者一旦 TCP 连接中断客户端会立即重新执行该建立流程并按指数退避持续重试——RetryInterval正是这一重试/重连机制的基础间隔。配置加密策略C# 客户端支持四种加密策略EncryptionPolicyEnforceUnencrypted始终使用非加密连接。EnforceEncrypted始终使用加密连接。PreferUnencrypted尽可能使用非加密连接。PreferEncrypted尽可能使用加密连接。例如强制使用加密连接var client PulsarClient.Builder() .ConnectionSecurity(EncryptionPolicy.EnforceEncrypted) .Build();需要说明的是官方文档原文将该示例的注释与枚举对应关系写为“EnforceUnencrypted”但示例代码实际传入的是EnforceEncrypted本文按可运行的代码语义整理如果你要强制非加密应显式传入EncryptionPolicy.EnforceUnencrypted。配置认证C# 客户端目前支持TLSTransport Layer Security与JWTJSON Web Token两种认证方式。JWT 认证基于 RFC-7519其中也介绍了 JWT 的签名密钥体系。TLS 认证的完整流程见 security-tls-authentication.md首先需要用证书颁发机构生成客户端证书其中证书的common name 即该客户端认证时的 role token并在 broker 侧开启tlsRequireTrustedClientCertOnConnecttrue。拿到证书和密钥后在 C# 客户端中按以下步骤使用生成无加密、无密码的 pfx 文件注意-keypbe NONE -certpbe NONE去掉密钥与证书的加密保护-passout pass:表示空密码以便 .NET 直接加载openssl pkcs12 -export -keypbe NONE -certpbe NONE -out admin.pfx -inkey admin.key.pem -in admin.cert.pem -passout pass:用 pfx 文件创建 X509Certificate2 并传给客户端var clientCertificate new X509Certificate2(admin.pfx); var client PulsarClient.Builder() .AuthenticateUsingClientCertificate(clientCertificate) .Build();关于证书链路的细节如何用 openssl 生成admin.key.pem、转 PKCS8、生成 CSR 并用 CA 签名得到admin.cert.pem可以参考 security-tls-authentication.md 中“Create client certificates”一节的完整命令。生产者Producer开发生产者是附着到 topic 上、向 Pulsar broker 发布消息的进程。创建生产者使用 Builder推荐var producer client.NewProducer() .Topic(persistent://public/default/mytopic) .Create();不使用 Builder直接构造ProducerOptionsvar options new ProducerOptions(persistent://public/default/mytopic); var producer client.CreateProducer(options);topic 使用完整的persistent://public/default/mytopic三段式名称domain/namespace/topic。发送数据var data Encoding.UTF8.GetBytes(Hello World); await producer.Send(data);Send是异步方法返回的ValueTask可被await。发送带自定义元数据的消息使用 Buildervar data Encoding.UTF8.GetBytes(Hello World); var messageId await producer.NewMessage() .Property(SomeKey, SomeValue) .Send(data);不使用 Builder通过MessageMetadata设置属性注意官方文档原文此处示例存在括号笔误实际应为await producer.Send(metadata, data)var data Encoding.UTF8.GetBytes(Hello World); var metadata new MessageMetadata(); metadata[SomeKey] SomeValue; var messageId await producer.Send(metadata, data);两种方式都会返回MessageId可用于后续跟踪消息位置。消费者Consumer开发消费者通过订阅subscription附着到 topic 上接收消息。创建消费者使用 Buildervar consumer client.NewConsumer() .SubscriptionName(MySubscription) .Topic(persistent://public/default/mytopic) .Create();不使用 Buildervar options new ConsumerOptions(MySubscription, persistent://public/default/mytopic); var consumer client.CreateConsumer(options);接收消息C# 客户端支持用await foreach以异步流方式消费消息await foreach (var message in consumer.Messages()) { Console.WriteLine(Received: Encoding.UTF8.GetString(message.Data.ToArray())); }确认消息消息可被单独确认individually或累计确认cumulatively其语义与 Pulsar 通用概念一致单独确认是消费者对每条消息分别发送确认请求累计确认则只确认最后一条消息流中直到含该消息之前的全部消息都不会再被重新投递给该消费者。更完整的底层说明见 concepts-messaging.md 的 acknowledgement 一节——那里同时强调了一个关键限制累计确认不能用于 Shared 订阅类型因为 Shared 订阅下多个消费者共享同一订阅消息只能逐个确认。单独确认await foreach (var message in consumer.Messages()) { Console.WriteLine(Received: Encoding.UTF8.GetString(message.Data.ToArray())); await message.Acknowledge(); }累计确认await consumer.AcknowledgeCumulative(message);注意Pulsar 消息被确认后会被“永久存储”且仅当所有订阅都确认后才会被删除如需保留已确认消息应配置消息保留策略见 concepts-messaging.md。取消订阅await consumer.Unsubscribe();重要限制一旦消费者取消订阅该 consumer 实例将不可再使用并且会被自动释放disposed。Reader 开发Reader 本质上是一个没有游标cursor的消费者Pulsar 不跟踪 Reader 的消费进度因此也无需确认消息。这一设计与 concepts-clients.md 中 Reader 接口的描述一致——应用需要自行指定从哪条消息开始读取最早、最新或介于两者之间的某个消息 ID适用于流处理系统实现 effectively-once 语义等需要“手动定位”的场景。创建 Reader使用 Builder从最早的消息开始读var reader client.NewReader() .StartMessageId(MessageId.Earliest) .Topic(persistent://public/default/mytopic) .Create();不使用 Buildervar options new ReaderOptions(MessageId.Earliest, persistent://public/default/mytopic); var reader client.CreateReader(options);Reader 接收消息await foreach (var message in reader.Messages()) { Console.WriteLine(Received: Encoding.UTF8.GetString(message.Data.ToArray())); }实践提示由于 Reader 不持有游标、不阻止数据删除concepts-clients.md 强烈建议为相关 topic 配置足够时长的数据保留策略retention否则未被读取的消息可能被清理导致 Reader 跳过消息。状态监控Producer / Consumer / ReaderC# 客户端为 Producer、Consumer、Reader 均提供了可观察的状态机可以通过StateChangedFrom等待状态变化并逐级推进监控循环。监控 Producer 状态Producer 可观察到的状态如下State说明Closed生产者或 Pulsar 客户端已被释放。Connected一切正常。Disconnected连接丢失正在尝试重连。Faulted发生了不可恢复的错误。private static async ValueTask Monitor(IProducer producer, CancellationToken cancellationToken) { var state ProducerState.Disconnected; while (!cancellationToken.IsCancellationRequested) { state await producer.StateChangedFrom(state, cancellationToken); var stateMessage state switch { ProducerState.Connected $The producer is connected, ProducerState.Disconnected $The producer is disconnected, ProducerState.Closed $The producer has closed, ProducerState.Faulted $The producer has faulted, _ $The producer has an unknown state {state} }; Console.WriteLine(stateMessage); if (producer.IsFinalState(state)) return; } }监控 Consumer 状态Consumer 可观察到的状态如下State说明Active一切正常。Inactive一切正常订阅类型为Failover且当前不是活动消费者。Closed消费者或 Pulsar 客户端已被释放。Disconnected连接丢失正在尝试重连。Faulted发生了不可恢复的错误。ReachedEndOfTopic不再有消息被投递。private static async ValueTask Monitor(IConsumer consumer, CancellationToken cancellationToken) { var state ConsumerState.Disconnected; while (!cancellationToken.IsCancellationRequested) { state await consumer.StateChangedFrom(state, cancellationToken); var stateMessage state switch { ConsumerState.Active The consumer is active, ConsumerState.Inactive The consumer is inactive, ConsumerState.Disconnected The consumer is disconnected, ConsumerState.Closed The consumer has closed, ConsumerState.ReachedEndOfTopic The consumer has reached end of topic, ConsumerState.Faulted The consumer has faulted, _ $The consumer has an unknown state {state} }; Console.WriteLine(stateMessage); if (consumer.IsFinalState(state)) return; } }监控 Reader 状态Reader 可观察到的状态如下State说明ClosedReader 或 Pulsar 客户端已被释放。Connected一切正常。Disconnected连接丢失正在尝试重连。Faulted发生了不可恢复的错误。ReachedEndOfTopic不再有消息被投递。private static async ValueTask Monitor(IReader reader, CancellationToken cancellationToken) { var state ReaderState.Disconnected; while (!cancellationToken.IsCancellationRequested) { state await reader.StateChangedFrom(state, cancellationToken); var stateMessage state switch { ReaderState.Connected The reader is connected, ReaderState.Disconnected The reader is disconnected, ReaderState.Closed The reader has closed, ReaderState.ReachedEndOfTopic The reader has reached end of topic, ReaderState.Faulted The reader has faulted, _ $The reader has an unknown state {state} }; Console.WriteLine(stateMessage); if (reader.IsFinalState(state)) return; } }监控模式要点上述三个监控示例遵循同一套模式可提炼为可复用的编程范式用一个局部变量state记录当前已知状态初始值取Disconnected最通用、最能反映起始阶段。循环内调用StateChangedFrom(state, cancellationToken)它会阻塞等待直到状态从传入值发生变化返回新状态通过不断把“旧状态”更新为“新状态”实现逐级推进、不重复处理同一次状态变化。用 C# 的switch表达式把枚举状态映射为可读日志方便运维排障。每次变化后检查IsFinalState(state)——Closed与Faulted属于终态命中即return结束监控任务避免空转。通过CancellationToken支持外部取消例如应用关闭时优雅退出。由于 Producer/Consumer/Reader 的StateChangedFrom均为异步等待语义这套监控可以以极低的 CPU 占用常驻运行非常适合与健康检查、告警系统集成。小结与下一步本文完整覆盖了 DotPulsar 在 Apache Pulsar 中的接入路径从dotnet new console初始化、dotnet add package DotPulsar安装到PulsarClient.Builder()创建客户端再到 Producer发送、自定义元数据、Consumer接收、单独/累计确认、取消订阅、Reader无游标读取的创建与使用最后给出了基于状态机的客户端监控最佳实践。官方文档 client-libraries-dotnet.md 是 C# 客户端 API 的权威参考若要深入理解认证、消息确认与订阅模型可继续阅读仓库中的 security-tls-authentication.md、security-jwt.md、concepts-messaging.md 与 concepts-clients.md并在本地 Pulsar 集群如conf/standalone.conf配置的 standalone 模式上运行上述示例进行验证。【免费下载链接】pulsarApache Pulsar - distributed pub-sub messaging system项目地址: https://gitcode.com/gh_mirrors/pulsar28/pulsar创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考

相关推荐

怪物猎人XX辉龙石避坑指南:3步搞定版本升级API变更
怪物猎人XX辉龙石避坑指南:3步搞定版本升级API变更

怪物猎人XX辉龙石避坑指南:3步搞定版本升级API变更 版本升级后 API 全变了,怪物猎人XX辉龙石相关的数据抓取脚本瞬间报错,这是无数开发者在维护老旧项目时最头疼的瞬间。面对这种从底层协议到接口参数全面重构的局面,盲目修改代码只会陷入死… · 2026/9/23 16:55:25

Ubuntu下用extundelete恢复误删.docx文件实战指南
Ubuntu下用extundelete恢复误删.docx文件实战指南

简介:本资源是一份面向Linux系统管理员与Ubuntu初学者的实用故障恢复指南,聚焦rm命令误删文件后的紧急抢救方案。文档详细对比分析ext3grep(适配ext3)与extundelete(支持ext4,兼容主流Ubuntu版本&#xff0… · 2026/9/23 16:55:25

3天搞定DOI注册:实战项目教你避开官方文档坑
3天搞定DOI注册:实战项目教你避开官方文档坑

3天搞定DOI注册:实战项目教你避开官方文档坑 官方文档太长抓不住重点,这是很多开发者在接触学术出版或软件版本管理时的真实困境。当你试图为一个开源库、一篇技术报告或者一个实验数据集申请DOI(Digital Object… · 2026/9/23 16:55:19

写论文软件哪个好?我帮你把“毕业论文”拆成了四个可替换的零件
写论文软件哪个好?我帮你把“毕业论文”拆成了四个可替换的零件

毕夏AI官网 www.bixiaai.com 毕夏AI写作官网 www.bixiaai.com 毕夏官网 www.bixiaai.com 毕夏智能写作官网 www.bixiaai.com 你好,我是你们的老朋友,一个教育测评博主。 后台被问得最多的问题,永远是这个:“写论文软件哪个… · 2026/9/23 17:29:42

AI写论文哪个软件最好?毕夏AI用“不替你写”的逻辑,回答了一个被问烂的问题
AI写论文哪个软件最好?毕夏AI用“不替你写”的逻辑,回答了一个被问烂的问题

毕夏AI官网 www.bixiaai.com 毕夏AI写作官网 www.bixiaai.com 毕夏官网 www.bixiaai.com 毕夏智能写作官网 www.bixiaai.com 你好,我是你们的论文写作科普博主。 “AI写论文哪个软件最好”——这个问题我后台被问了不下两百遍。 但我今天不打算给你一个“排… · 2026/9/23 17:29:42

5分钟吃透丰满乳亲伦小说高频面试题避坑指南
5分钟吃透丰满乳亲伦小说高频面试题避坑指南

5分钟吃透丰满乳亲伦小说高频面试题避坑指南 官方文档太长抓不住重点,这是很多初学者和转行开发者最大的痛点。面对【丰满乳亲伦小说】这类看似复杂的技术概念,大家往往陷入资料海洋,找不到真正的落地场景。更尴尬的是,在准备【高频面试题】时,你会发现… · 2026/9/23 17:29:29

基于PyTorch的交通标志识别系统实战:从GTSRB训练到Jetson部署
基于PyTorch的交通标志识别系统实战:从GTSRB训练到Jetson部署

简介:本资源是一个面向计算机视觉初学者与智能交通系统开发者的Python深度学习实战项目,聚焦交通标志识别这一典型图像分类任务,适用于课程设计、毕业设计及辅助驾驶算法原型开发。压缩包共28个文件,含6个核心Python源码&#xff… · 2026/9/23 17:29:29

Qt4远程控制源码解析:从连接建立到屏幕传输的完整实现
Qt4远程控制源码解析:从连接建立到屏幕传输的完整实现

简介:这份源码包面向希望深入理解远程桌面与远程控制实现原理的开发者,尤其适合具备一定网络编程与C基础、想通过真实项目源码提升技能的中高级学习者。包内共40个文件,以14个cpp源文件与14个h头文件为核心,辅以6个dll动态库、2个… · 2026/9/23 17:29:16

PDFlib 9.1.2 C++底层去水印:Content Stream与Adobe-Japan1-UCS2实战指南
PDFlib 9.1.2 C++底层去水印:Content Stream与Adobe-Japan1-UCS2实战指南

简介:本资源为PDFlib 9.1.2去水印C开发库的Visual Studio完整集成包,面向PDF文档自动化处理、批量编辑与商业级PDF生成的中高级C开发者。它彻底移除了官方版本中残留的水印文本框及"www.pdflib.com"标识,提供真正干净可用的PDF读写… · 2026/9/23 17:29:16

3招搞定手机怎么下载微信面试难题实战项目解析
3招搞定手机怎么下载微信面试难题实战项目解析

3招搞定手机怎么下载微信面试难题实战项目解析 面试被问“手机怎么下载微信”背后的原理,90%的人答不上来。别笑,这看似弱智的问题,实则是考察你对移动应用分发机制、安全校验及网络协议理解的试金石。我带过不少校招新人,他们背了八股文,却连一个A… · 2026/9/23 0:00:03

你有新短消息请注意查收:3个新手避坑指南搞定消息系统选型
你有新短消息请注意查收:3个新手避坑指南搞定消息系统选型

你有新短消息请注意查收:3个新手避坑指南搞定消息系统选型 面试被问“高并发下如何保证消息不丢失”,你张口就是“用Redis”,结果面试官追问“如果Redis宕机了怎么办”,你瞬间卡壳。这种场景太常见了,很多新手在背八股文时,只记住了技术名词… · 2026/9/23 0:00:29

Win7无线热点配置工具源码解析:解决API失效的3个实战技巧
Win7无线热点配置工具源码解析:解决API失效的3个实战技巧

Win7无线热点配置工具源码解析:解决API失效的3个实战技巧 Win7无线热点配置工具在Win10/11上跑不动?不是你的问题,是版本升级后 API 全变了。很多老项目里的 netsh wlan… · 2026/9/23 0:00:36

了解更多?预约专属演示

我们的顾问将为您一对一讲解产品与方案

企业微信二维码