完整使用指南:安装、生产者/消费者/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),仅供参考