Apache Pulsar C 客户端(DotPulsar)完整使用指南:安装、生产者/消费者/Reader 开发与状态监控
2026/9/23 16:55:21 网站建设 项目流程

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

【免费下载链接】pulsarApache Pulsar - distributed pub-sub messaging system项目地址: https://gitcode.com/gh_mirrors/pulsar28/pulsar

Apache 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 环境中操作。

安装步骤

  1. 创建项目文件夹,并打开终端切换到该目录。

  2. 初始化控制台项目

    dotnet new console
  3. 使用dotnet run运行一次,验证应用已正确创建。

  4. 添加 DotPulsar NuGet 包

    dotnet add package DotPulsar
  5. 命令执行完成后,打开.csproj文件即可看到自动加入的包引用(官方文档示例中的版本为 0.11.0):

    <ItemGroup> <PackageReference Include="DotPulsar" Version="0.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:6650
RetryInterval设置操作或重连前的等待时间3s

结合 concepts-clients.md 的客户端建立流程,可以更准确地理解 ServiceUrl 的作用:应用创建 producer/consumer 前,客户端会先通过 HTTP 查找请求确定 topic 归属的 broker,再建立 TCP 连接并完成认证,最后在连接上创建生产者/消费者;一旦 TCP 连接中断,客户端会立即重新执行该建立流程,并按指数退避持续重试——RetryInterval正是这一重试/重连机制的基础间隔。

配置加密策略

C# 客户端支持四种加密策略(EncryptionPolicy):

  • EnforceUnencrypted:始终使用非加密连接。
  • EnforceEncrypted:始终使用加密连接。
  • PreferUnencrypted:尽可能使用非加密连接。
  • PreferEncrypted:尽可能使用加密连接。

例如强制使用加密连接:

var client = PulsarClient.Builder() .ConnectionSecurity(EncryptionPolicy.EnforceEncrypted) .Build();

需要说明的是,官方文档原文将该示例的注释与枚举对应关系写为“EnforceUnencrypted”,但示例代码实际传入的是EnforceEncrypted,本文按可运行的代码语义整理:如果你要强制非加密,应显式传入EncryptionPolicy.EnforceUnencrypted

配置认证

C# 客户端目前支持TLS(Transport Layer Security)JWT(JSON Web Token)两种认证方式。JWT 认证基于 RFC-7519(其中也介绍了 JWT 的签名密钥体系)。

TLS 认证的完整流程见 security-tls-authentication.md:首先需要用证书颁发机构生成客户端证书,其中证书的common name 即该客户端认证时的 role token,并在 broker 侧开启tlsRequireTrustedClientCertOnConnect=true。拿到证书和密钥后,在 C# 客户端中按以下步骤使用:

  1. 生成无加密、无密码的 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:
  2. 用 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,直接构造ProducerOptions

    var 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

发送带自定义元数据的消息

  • 使用 Builder

    var 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 上接收消息。

创建消费者

  • 使用 Builder

    var consumer = client.NewConsumer() .SubscriptionName("MySubscription") .Topic("persistent://public/default/mytopic") .Create();
  • 不使用 Builder

    var 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();
  • 不使用 Builder

    var 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 / Reader

C# 客户端为 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; } }

监控模式要点

上述三个监控示例遵循同一套模式,可提炼为可复用的编程范式:

  1. 用一个局部变量state记录当前已知状态,初始值取Disconnected(最通用、最能反映起始阶段)。
  2. 循环内调用StateChangedFrom(state, cancellationToken),它会阻塞等待直到状态从传入值发生变化,返回新状态;通过不断把“旧状态”更新为“新状态”,实现逐级推进、不重复处理同一次状态变化。
  3. 用 C# 的switch表达式把枚举状态映射为可读日志,方便运维排障。
  4. 每次变化后检查IsFinalState(state)——ClosedFaulted属于终态,命中即return结束监控任务,避免空转。
  5. 通过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),仅供参考

需要专业的网站建设服务?

联系我们获取免费的网站建设咨询和方案报价,让我们帮助您实现业务目标

立即咨询