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属性可查看存储大小、消息条数等状态。 - 列出所有 Stream:
ListStreamsAsync()返回异步可枚举集合,配合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");这正是NatsJSStream(src/NATS.Client.JetStream/NatsJSStream.cs)中提供的便捷能力,同时它还支持RefreshAsync刷新状态、PurgeAsync清空数据等操作。
六、最佳实践与注意事项
- 先建 Stream,再谈 Consumer:Consumer 必须依附于已存在的 Stream。
- 善用 CreateOrUpdate 语义:开发阶段用
CreateOrUpdateStreamAsync/CreateOrUpdateConsumerAsync可避免"已存在"报错。 - 删除操作不可逆:
DeleteStreamAsync会连数据一起删除,生产环境务必先备份或使用 Purge 替代。 - 临时 Consumer 会自清理:不设置
DurableName的消费者在无订阅一段时间后会被自动删除。 - 异步 API 统一:所有管理方法都返回
ValueTask,全程异步、支持CancellationToken取消。 - 消息发布用 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),仅供参考