go中nats订阅需按场景选型:异步订阅适合事件通知但不保序不持久;同步订阅需配deliverall和durable才能续读;jetstream订阅必须显式设deliverpolicy(如deliverall)且与流的retentionpolicy协同,否则收不到历史消息或重复消费。

Go 里用 nats 订阅消息,不是调个 Subscribe() 就完事——多数“收不到消息”“重启后丢数据”“重复消费”问题,都卡在订阅方式选错、参数没配对、或压根没启用 JetStream。
异步订阅(nc.Subscribe())适合什么场景
这是最常用、也最容易误用的方式。它启动一个 goroutine 监听 subject,收到消息立刻回调,不阻塞主流程。
- 适用于“发了就不管”的事件通知,比如日志上报、监控打点
- 不保证顺序:多个
Subscribe()并发处理时,同一 subject 的消息到达顺序和发布顺序可能不一致 - 不持久:纯 NATS 模式下,订阅者断开期间的消息直接丢弃,不会补推
- 别在循环里反复
Subscribe()同一个 subject —— 会创建多个独立订阅,消息被重复投递多次
同步订阅(nc.ChannelSubscribe() 或 sub.NextMsg())怎么避免阻塞
同步订阅本质是把消息塞进 Go channel,由你手动从 channel 取;或者用 sub.NextMsg(timeout) 主动拉取。它不自动启 goroutine,可控性更强。
- 用
nc.ChannelSubscribe()时,记得配nats.DeliverAll和nats.Durable("my-processor"),否则重启后无法续读历史消息 -
sub.NextMsg(5 * time.Second)适合短时轮询任务,但 timeout 太小会频繁返回nats.ErrTimeout,太大又拖慢响应 - channel 容量必须显式设:比如
ch := make(chan *nats.Msg, 1024),否则默认 0 容量,一发消息就卡死
JetStream 订阅必须带 nats.DeliverPolicy() 才能读历史
连上 JetStream 后,js.Subscribe() 默认行为仍是跳过已发布消息——这不是 bug,是设计。想从头消费,必须显式声明策略。
-
nats.DeliverAll:从流最早一条开始读(适合初始化、重放) -
nats.DeliverLastPerSubject:只取每个 subject 最新一条(适合状态同步) -
nats.DeliverNew:只收订阅后新发布的消息(默认值,等同于裸 NATS 行为) - 漏掉
nats.Durable("xxx")?那每次都是全新 offset,重启等于重头开始
队列组订阅(nats.QueueSubscribe())为什么消息只被一个实例处理
这是实现负载均衡的关键机制:相同 queue name 的多个订阅者构成一个“队列组”,NATS 保证每条消息只投给其中一人。
- 适合 worker 类服务,比如订单处理集群,避免多实例重复扣款
- 必须搭配 JetStream 才能持久化未确认消息(否则 worker 崩溃,消息就丢了)
- 注意:
nats.DeliverAll在队列组里仍有效,但只对当前被选中的那个 worker 生效 - 别把 queue name 写成随机值(如
uuid.New().String()),否则失去负载均衡意义
真正容易被忽略的是:JetStream 流(Stream)的 RetentionPolicy 和订阅端的 DeliverPolicy 必须协同——比如设了 InterestPolicy 却用 DeliverAll,可能读到空流;而 WorkQueuePolicy 下用 DeliverAll 则根本拿不到旧消息。配置不是写完就跑,得看它和流定义是否咬合。
golang免费学习笔记(深入):立即使用
在学习笔记中,你将探索golang的核心概念和高级技巧!











