NATS.Net JetStream入门:5步创建Stream与Consumer实现消息持久化
NATS.Net JetStream入门:5步创建Stream与Consumer实现消息持久化
【免费下载链接】nats.netThe official C# Client for NATS项目地址: https://gitcode.com/gh_mirrors/na/nats.net
NATS.Net 是 NATS 官方推出的 C# 客户端库,支持 .NET Standard 2.0 及以上版本,而 JetStream 是 NATS 内置的分布式消息持久化系统。这篇 NATS.Net JetStream入门教程面向零基础新手,用 5 个简单步骤带你创建 Stream 与 Consumer,实现可靠的消息持久化——即使发布者下线、消费者重启,消息也不会丢失。
JetStream 是什么?为什么需要消息持久化?
普通 NATS 发布/订阅模式有一个限制:订阅者不在线时发布的消息会直接丢失。JetStream 正是为解决这个问题而生——它内置于 nats-server,无需额外部署,只需一个参数即可开启。
JetStream 的两个核心概念:
- Stream(消息流):相当于一个消息仓库,把发布到指定主题(Subject)的消息持久化存储起来;
- Consumer(消费者):Stream 上的"视图",负责跟踪每条消息的投递和确认(ACK)状态。
🚀 换句话说:Stream 负责存,Consumer 负责读,两者配合就构成了完整的消息持久化链路。
第1步:启动支持 JetStream 的 NATS 服务器
JetStream 是服务器端能力,首先需要启动一个开启 JetStream 的 NATS 服务。两种方式任选其一:
方式一:直接运行(需先安装 nats-server)
nats-server -js方式二:使用 Docker(最省事)
docker run -p 4222:4222 nats -js看到JetStream is enabled的日志即表示持久化功能已就绪。生产环境建议使用 3 或 5 节点集群实现容错,开发调试单节点完全够用。
第2步:创建项目并安装 NATS.Net 包
打开终端创建控制台项目,并安装聚合包:
dotnet new console -n JetStreamDemo cd JetStreamDemo dotnet add package NATS.NetNATS.Net 是一个元包(Meta Package),一次安装即可获得 Core、JetStream、KeyValueStore、ObjectStore 等全部子库能力,无需分别引用。
第3步:连接服务器并创建 JetStream 上下文
JetStream 上下文(INatsJSContext)是所有流管理操作的入口。使用简化的NatsClient只需要两行代码:
await using var nc = new NatsClient(); // 默认连接 nats://localhost:4222 var js = nc.CreateJetStreamContext(); // 创建 JetStream 上下文这一模式在官方示例中随处可见,例如 tests/NATS.Net.DocsExamples/JetStream/IntroPage.cs 就是从这里开始的。
第4步:创建 Stream 实现消息持久化
创建名为ORDERS的 Stream,并让它监听orders.>通配主题(即以orders.开头的所有主题,如orders.new.1):
await js.CreateStreamAsync(new StreamConfig( name: "ORDERS", subjects: ["orders.>"]));Stream 创建成功后,向对应主题发布的消息就会被持久化存储。一定要使用 JetStream 上下文发布,这样能收到服务器返回的存储确认:
var ack = await js.PublishAsync("orders.new.1", order); ack.EnsureSuccess(); // 确认消息已成功落盘StreamConfig还支持配置保留策略、存储方式(内存/文件)、消息上限等,完整字段可参考 src/NATS.Client.JetStream/Models/StreamConfig.cs。
第5步:创建 Consumer 并消费持久化消息
创建名为order_processor的 Consumer(重复调用会自动更新或复用),然后循环消费:
var consumer = await js.CreateOrUpdateConsumerAsync( stream: "ORDERS", new ConsumerConfig("order_processor")); await foreach (var msg in consumer.ConsumeAsync<Order>()) { Console.WriteLine($"处理订单: {msg.Data.Id}"); await msg.AckAsync(); // 手动确认,处理成功后才算消费完成 }✨ 关键点:只有调用AckAsync()确认后,消息才会从待处理列表移除;若消费者崩溃,未确认消息会被重新投递,这正是 JetStream 消息持久化可靠性的体现。
Consumer 还支持三种消费模式,可按场景选用:
- Consume:持续批量推送消费(适合长时间运行的任务);
- Fetch:按批次拉取指定数量消息(如
MaxMsgs = 1000); - Next:单条拉取下一条消息。
更多消费写法可参考 tests/NATS.Net.DocsExamples/JetStream/ConsumePage.cs,完整的可运行示例见 examples/Example.JetStream.PullConsumer/Program.cs。
进阶技巧:Durable 与 Ephemeral Consumer
创建 Consumer 时设置了DurableName即为持久消费者,断开后状态保留、可随时恢复;不设置则创建临时消费者(Ephemeral),空闲一段时间后会被服务器自动清理。
var durable = new ConsumerConfig { Name = "durable_processor", DurableName = "durable_processor", };日常开发建议先跑通入门示例再深入调优。你也可以通过git clone https://gitcode.com/gh_mirrors/na/nats.net获取完整仓库,其中examples/目录下有覆盖拉取消费、JetStream、KV 存储等场景的现成代码可直接运行。
总结
通过以上 5 步,你已经完成了 NATS.Net JetStream 的入门闭环:开启服务器 → 安装客户端 → 建立上下文 → 创建 Stream 持久化消息 → 创建 Consumer 可靠消费。借助 JetStream,你的应用可以轻松获得消息持久化、重放、限流等生产级能力,而这一切都基于官方 C# 客户端 NATS.Net,上手成本极低。
下一步建议阅读仓库中的 JetStream 文档目录tools/site_src/documentation/jetstream/,深入了解发布重试、流管理、有序消费等进阶特性。
【免费下载链接】nats.netThe official C# Client for NATS项目地址: https://gitcode.com/gh_mirrors/na/nats.net
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
