nats jetstream 不是独立“排队网关”,而是需显式启用服务端(-js)、客户端初始化 js 后创建 stream 并配合 queuesubscribe 实现负载均衡与持久化;缺一环消息即丢。

直接说结论:NATS JetStream 不是“部署成队列网关”,而是用 nats-server 启动一个带 JetStream 模块的实例,再让 Go 微服务通过 nats.go 连接它——所谓“排队网关”本质是 JetStream Stream + Queue Group 的组合用法,不是独立组件。
JetStream 必须显式启用并建 Stream,否则消息真会丢
很多人以为连上 nats://localhost:4222 就能用 JetStream,结果 js.Publish() panic 或消息一断线就蒸发。根本原因是:JetStream 是 NATS Server 的可选模块,默认不开启;即使开了,也必须手动创建 Stream,否则所有发布都落空。
- 启动服务端时加
-js参数:nats-server -js -c nats.conf(或 Docker 中设NATS_JETSTREAM=1) - Go 客户端初始化后立刻调
js, err := nc.JetStream(),且err必须检查——别等到第一次 publish 才发现nil pointer - Stream 创建不可省:
js.AddStream(&jetstream.StreamConfig{Subjects: []string{"gateway.>"}, RetentionPolicy: jetstream.InterestPolicy}),不建 Stream,gateway.queue主题发的消息就只是内存广播,重启或掉线即丢 -
InterestPolicy适合网关场景:只保留当前有消费者订阅的主题消息,比LimitsPolicy更省空间,但要注意——若所有消费者全下线,消息会被自动清理
Queue Group 是实现“排队网关”的核心,不是靠单个 Subscribe
想让多个微服务实例像线程池一样分担任务?别用 Subscribe(),那是广播模式。真正做负载分摊、避免重复消费,必须用 QueueSubscribe() 并指定相同 queue name。
- 发布端仍用
js.Publish("gateway.task", data),主题名要和 Stream 的Subjects匹配(比如gateway.>) - 消费端写:
js.QueueSubscribe("gateway.task", "gateway-workers", handler, nats.ManualAck(), nats.AckWait(30*time.Second)),其中"gateway-workers"是队列组名,同名的多个实例自动负载均衡 - 漏掉
nats.ManualAck()就等于没开持久化:消息投递后立即标记为“已处理”,断线重连后收不到积压 -
AckWait要大于业务处理耗时,否则超时自动重发,导致重复执行——比如扣库存操作耗时 5 秒,AckWait至少设 10 秒以上
连接参数不设超时和重连,上线就卡死
本地跑通不代表生产可用。nats.Connect("nats://localhost:4222") 在测试环境大概率 hang 住,因为默认无限重试、无单次超时、无 jitter 防抖。
- 必须设
nats.Timeout(5 * time.Second):单次连接最多等 5 秒,DNS 慢或防火墙拦截时快速失败 -
nats.MaxReconnects(-1)表示永久重试(别用 0 或正数),配合nats.ReconnectWait(2 * time.Second)和nats.ReconnectJitter(100*time.Millisecond, time.Second)避免雪崩式重连 - 集群地址写成逗号分隔:
"nats://n1:4222,nats://n2:4222,nats://n3:4222",客户端自动剔除故障节点并轮询 - 生产环境必须加认证:
nats.UserCredentials("nats.creds")或nats.Token("xxx"),否则日志里全是Authorization Error
事件结构体没 Type/Version 字段,网关就变成黑盒
网关转发事件给下游微服务,如果 payload 是裸 json.RawMessage 或 map[string]interface{},下游根本没法路由、没法升级、一加字段就 panic。
- 强制定义顶层结构:
type Event struct { Type string `json:"type"` Version string `json:"version"` Payload json.RawMessage `json:"payload"` } - 下游按
event.Type == "order.created" && event.Version == "v1"分发到对应 handler,版本升级时旧服务继续跑 v1,新服务接 v2 - 别把时间戳存字符串,用
Timestamp time.Time `json:"timestamp,string"`,防解析错且便于日志追踪 - 发布前校验:
if event.Type == "" { return errors.New("missing event type") },避免无效消息污染流
最常被忽略的是:JetStream Stream 的 RetentionPolicy 和 Queue Group 的 ManualAck 必须同时生效,缺一不可。前者决定消息存多久,后者决定是否重投——两者不配对,网关就既不能保序,也无法容错。
golang免费学习笔记(深入):立即使用
在学习笔记中,你将探索golang的核心概念和高级技巧!











