nats 客户端必须在 kratos 配置加载完成后初始化,使用 app.new() 回调或 di 注入;关键路径用 publish() 并检查 error,异步场景配 flushtimeout;消费端用 queuesubscribe() 和 maxinflight 实现并发;jetstream 必须启用并正确配置 stream、消费者 ack 及重试策略。

别在 main.go 里直接 nats.Connect() —— 配置没加载完就连,conf.Nats.URL 是空字符串,连接地址变成 nats://,报错 nats: no servers specified,服务直接 panic。
Kratos 等框架里 NATS 客户端初始化时机不对
典型错误是在 main() 函数开头、甚至 init() 里就调 nats.Connect(),此时 Kratos 的配置还没解析,conf.Nats.URL 还是零值。必须等 config.ToType(&conf.Nats) 执行完毕再初始化客户端。
- 正确做法:把 NATS 初始化放在 Kratos
app.New()的回调里,或封装成一个依赖项,在 DI 容器中按需注入 - 别在
initNATS()函数里写defer nc.Close()—— 连接立刻被关,后续所有Publish()都 panic - 连接对象
*nats.Conn必须存进 app 生命周期管理,由app.Run()启动时建立、app.Stop()时关闭
Publish() 和 PublishAsync() 选错导致消息静默丢失
Publish() 是同步阻塞,等服务器 ACK 才返回;PublishAsync() 是异步非阻塞,发出去就继续执行,但若 NATS 不可达或网络抖动,消息会丢且不报错。
在 Golang 中使用 samber/hot 进行内存缓存,支持 LRU、LFU、TinyLFU、W‑TinyLFU、S3FIFO、ARC、TwoQueue、SIEVE、FIFO 等淘汰算法,提供 TTL、缓存加载器及分片功能。
- 订单创建、支付回调等关键路径,必须用
Publish(),并检查返回的error - 日志、埋点、通知类弱一致性消息可用
PublishAsync(),但得配nc.FlushTimeout(5*time.Second),否则缓冲区满后新消息直接被丢 - 别漏掉
nc.Flush()或nc.FlushTimeout()—— 异步模式下不 flush,消息可能卡在 client buffer 里出不去
高吞吐消费端用 Subscribe() 就是单点瓶颈
Web 服务里常见错误:用 Subscribe() 处理事件流,所有消息串行进同一个 goroutine,CPU 利用率低、延迟飙升,QPS 上不去。
- 微服务场景必须用
QueueSubscribe(),配合相同 queue name 的多个实例,实现负载分摊 - 记得设
nats.MaxInflight(256)(或根据业务调整),不然默认 inflight=1,相当于又串行了 - 主题名是纯字符串匹配:
"order.created"和"order.created.v1"完全无关,发布和订阅两端必须完全一致
JetStream 没开、Stream 没配、消费者没 ACK = 白连
裸连 NATS 默认不持久、不重试、不幂等。消费者掉线期间发的消息直接蒸发,不是你代码错了,是它压根没存。
- 要“不丢”,必须启用 JetStream:
js, _ := jetstream.New(nc),且连接后立即调用,否则首次js.Publish()可能 panic - 创建 Stream 时必须指定
RetentionPolicy,比如jetstream.InterestPolicy(按兴趣保留)或jetstream.WorkQueuePolicy(仅保留未确认) - 消费者必须显式调
msg.Ack(),否则 AckWait 超时后重投;失败要用msg.NakWithDelay()+nats.MaxDeliver(3)控制重试次数
真正难的不是连上 NATS,而是让每条消息在断网、重启、扩容缩容时都可追溯、可重放、不重复——这需要 JetStream 配置、消费者幂等、发布端 MsgID、去重窗口、Ack 机制全部对齐,漏一个环节,线上就出数据问题。
golang免费学习笔记(深入):立即使用
在学习笔记中,你将探索golang的核心概念和高级技巧!










