NATS.Net JetStream管理API指南:Stream与Consumer的增删改查全解析
2026/8/21 13:04:19 网站建设 项目流程

NATS.Net JetStream管理API指南:Stream与Consumer的增删改查全解析

【免费下载链接】nats.netThe official C# Client for NATS项目地址: https://gitcode.com/gh_mirrors/na/nats.net

NATS.Net 是 NATS 官方的 C# 客户端库,而 JetStream 是 NATS 内置的持久化消息流引擎。本文是一份面向初学者的NATS.Net JetStream 管理 API实战指南,带你完整掌握 **Stream(消息流)**与 **Consumer(消费者)**的创建、查询、更新与删除全流程,也就是"增删改查"的每一个细节,并附上可以直接运行的 C# 代码示例。

如果你正在用 C# 开发消息队列、事件驱动或数据流应用,这份 NATS.Net 管理 API 教程将是你的最佳起点。🚀

一、NATS.Net JetStream 管理 API 是什么?

在 NATS.Net 中,管理 Stream 和 Consumer 的统一入口是NatsJSContext(JetStream 上下文)。它封装了与 JetStream 服务器交互的所有管理指令,核心实现在src/NATS.Client.JetStream/NatsJSContext.cs,而 Stream 与 Consumer 的具体管理方法分别位于:

  • src/NATS.Client.JetStream/NatsJSContext.Streams.cs—— Stream 增删改查
  • src/NATS.Client.JetStream/NatsJSContext.Consumers.cs—— Consumer 增删改查

它的典型用法是:先用NatsConnection建立连接,再创建上下文,之后所有管理操作都通过上下文发起。

二、快速上手:如何创建 JetStream 上下文

最简单的方式是使用NatsClient扩展方法CreateJetStreamContext()

await using var nc = new NatsClient(); // 创建 JetStream 管理上下文 INatsJSContext js = nc.CreateJetStreamContext();

也可以基于已有的连接手动创建:

await using var nats = new NatsConnection(); // 传入连接对象创建 JetStream 上下文 var js = new NatsJSContext(nats);

拿到js之后,下面的所有增删改查操作就都能执行了。参考官方文档示例:tests/NATS.Net.DocsExamples/JetStream/ManagingPage.cs

三、JetStream Stream 增删改查全解析

Stream 是消息的持久化存储,相当于"消息仓库"。我们按增删改查四步走。

1. 创建 Stream(增)

使用CreateStreamAsync传入StreamConfig配置,指定名称和订阅的 subject 即可:

await js.CreateStreamAsync(new StreamConfig( name: "ORDERS", subjects: new[] { "orders.>" }));

这样所有发布到orders.>主题的消息都会被 ORDERS 这个 Stream 持久化保存。💾

2. 查询 Stream(查)

  • 获取单个 Stream 信息GetStreamAsync("ORDERS")返回INatsJSStream对象,通过其Info属性可查看存储大小、消息条数等状态。
  • 列出所有 StreamListStreamsAsync()返回异步可枚举集合,配合await foreach遍历;还支持按 subject 过滤,例如ListStreamsAsync("orders.>")
  • 只列出 Stream 名称ListStreamNamesAsync()轻量级列出名字列表。
// 列出全部 Stream 名称 await foreach (var name in js.ListStreamNamesAsync()) { Console.WriteLine($"Stream: {name}"); }

3. 更新 Stream(改)

修改 Stream 配置(如调整最大存储、保留策略)有两个选择:

  • UpdateStreamAsync(config):要求 Stream 必须已存在,否则报错。
  • CreateOrUpdateStreamAsync(config):存在则更新,不存在则创建,更省心。👍
var config = new StreamConfig( name: "ORDERS", subjects: new[] { "orders.>", "archived.orders.>" }) { MaxAge = TimeSpan.FromDays(30), // 消息最多保留 30 天 }; await js.CreateOrUpdateStreamAsync(config);

4. 删除 Stream(删)

DeleteStreamAsync("ORDERS")会连同其中所有消息一起删除,操作不可逆,请谨慎使用 ⚠️:

bool ok = await js.DeleteStreamAsync("ORDERS"); Console.WriteLine($"删除结果: {ok}");

5. 附加操作:清空与删除单条消息

  • PurgeStreamAsync:清空 Stream 中的数据但保留 Stream 本身,适合"重置仓库"。
  • DeleteMessageAsync:按序列号删除 Stream 中的某一条消息。

四、JetStream Consumer 增删改查全解析

Consumer 是 Stream 的"消费视图",负责跟踪哪些消息已投递、已确认。

1. 创建 Consumer(增)

使用CreateOrUpdateConsumerAsync指定所属 Stream 和ConsumerConfig

// 创建(或获取)一个名为 order_processor 的消费者 INatsJSConsumer consumer = await js.CreateOrUpdateConsumerAsync( stream: "ORDERS", new ConsumerConfig("order_processor"));

Consumer 分为两种类型:

类型配置方式生命周期
持久型(Durable)设置DurableName长期存在,直到手动删除
临时型(Ephemeral)不设置名称无订阅后自动清理

2. 查询 Consumer(查)

  • GetConsumerAsync(stream, consumer):获取单个 Consumer 的详细信息。
  • ListConsumersAsync(stream):列出指定 Stream 下的所有 Consumer 对象。
  • ListConsumerNamesAsync(stream):列出 Consumer 名称列表。
// 查看 ORDERS 下有哪些消费者 await foreach (var name in js.ListConsumerNamesAsync("ORDERS")) { Console.WriteLine($"Consumer: {name}"); }

3. 更新 Consumer(改)

UpdateConsumerAsync(stream, config)可修改 Consumer 的配置,例如改变确认策略或过滤主题;CreateOrUpdateConsumerAsync同样适用,语义为"有则改、无则建"。

4. 删除 Consumer(删)

bool ok = await js.DeleteConsumerAsync("ORDERS", "order_processor");

删除后该 Consumer 对象不可再使用。🧹

5. 进阶操作:暂停、恢复与有序消费者

  • PauseConsumerAsync(stream, consumer, pauseUntil):让消费者暂停到指定时间。
  • ResumeConsumerAsync(stream, consumer):恢复被暂停的消费者。
  • CreateOrderedConsumerAsync(stream):创建自动处理重排的有序消费者,适合对消息顺序敏感的场景。

五、通过 Stream 对象管理 Consumer

拿到INatsJSStream后,也可以直接通过它管理下属消费者,代码更简洁:

var stream = await js.GetStreamAsync("ORDERS"); // 直接在该 Stream 下创建、查询、删除 Consumer var consumer = await stream.CreateOrUpdateConsumerAsync(new ConsumerConfig("worker")); await stream.DeleteConsumerAsync("worker");

这正是NatsJSStreamsrc/NATS.Client.JetStream/NatsJSStream.cs)中提供的便捷能力,同时它还支持RefreshAsync刷新状态、PurgeAsync清空数据等操作。

六、最佳实践与注意事项

  1. 先建 Stream,再谈 Consumer:Consumer 必须依附于已存在的 Stream。
  2. 善用 CreateOrUpdate 语义:开发阶段用CreateOrUpdateStreamAsync/CreateOrUpdateConsumerAsync可避免"已存在"报错。
  3. 删除操作不可逆DeleteStreamAsync会连数据一起删除,生产环境务必先备份或使用 Purge 替代。
  4. 临时 Consumer 会自清理:不设置DurableName的消费者在无订阅一段时间后会被自动删除。
  5. 异步 API 统一:所有管理方法都返回ValueTask,全程异步、支持CancellationToken取消。
  6. 消息发布用 JetStream API:向 Stream 发布消息建议使用js.PublishAsync,可获得服务器 ACK 确认写入成功。

完整的示例工程可参考examples/Example.JetStream.PullConsumer/Program.cs,本地运行前先启动支持 JetStream 的 NATS 服务器。

七、总结

通过本文,你已经掌握了 NATS.Net JetStream 管理 API 的核心:用NatsJSContext一个入口完成Stream 与 Consumer 的增删改查。这套 API 设计统一、全程异步,配合CreateOrUpdate语义和有序消费者等高级特性,足以支撑从开发调试到生产运维的各种场景。

现在就可以动手创建你的第一个 Stream 和 Consumer,开启 C# 与 JetStream 的持久化消息之旅吧!🎉

【免费下载链接】nats.netThe official C# Client for NATS项目地址: https://gitcode.com/gh_mirrors/na/nats.net

创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考

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

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

立即咨询