jetstream不能直接当消息队列用,因其本质是流式存储引擎,不自动删除已消费消息,也不保证“每条消息只被消费一次”,必须显式配置consumer的ackpolicy、deliverpolicy等参数并手动ack才能模拟可靠队列语义。

JetStream 为什么不能直接当“消息队列”用
JetStream 本质是流式存储引擎,不是传统意义上的队列。它不自动删除已消费消息,也不保证“每条消息只被消费一次”——除非你显式配置消费者(Consumer)为 deliver_policy = "by_start_time" 或启用 ack 机制。很多 Golang 开发者一上来就调 js.Publish() 然后等 js.Subscribe() 自动收,结果发现消息堆积、重复消费、重启后丢失 offset,根源就在这儿。
关键判断:JetStream 必须配合 Consumer + Ack 才能模拟可靠事件队列语义。
创建 Stream 和 Consumer 的最小可行配置
用 nats.go v1.29+ 创建 JetStream 实例后,Stream 要设 RetentionPolicy: nats.InterestPolicy(否则消息永不清除),Consumer 必须设 AckPolicy: nats.AckExplicit 并指定 DeliverPolicy: nats.DeliverAll 或 DeliverByStartSequence。
-
Stream名字必须全小写、无下划线(如"orders"),否则nats.go会静默失败 -
Consumer名字建议带服务标识(如"orders-processor-v1"),避免多个服务共用同一 consumer 导致 ack 冲突 - 若需严格顺序,
MaxAckPending设为 1;若要吞吐,可设为 100~1000,但需确保业务逻辑能处理乱序
示例片段:
Go 配置库,使用 spf13/viper — 分层优先级(flag > env >file > KV > default),提供 BindPFlag/BindPFlags、SetEnvPrefix + SetEnvKeyReplace 等功能。
js, _ := nc.JetStream()
_, err := js.AddStream(&nats.StreamConfig{
Name: "orders",
Subjects: []string{"orders.>"},
Retention: nats.InterestPolicy,
})
_, err = js.AddConsumer("orders", &nats.ConsumerConfig{
Durable: "orders-processor-v1",
AckPolicy: nats.AckExplicit,
DeliverPolicy: nats.DeliverAll,
MaxAckPending: 100,
})
消费端必须手动 Ack,且不能漏 defer
JetStream 不像 RabbitMQ 那样自动 ack。每次 msg.Next() 或 msg.Ack() 后,必须显式调用 msg.Ack(),否则消息会在 AckWait(默认 30s)后重投。更危险的是:如果 handler panic 了,没 defer ack,这条消息就永远卡在 pending 状态。
- 务必在 handler 开头加
defer msg.Ack(),哪怕后续逻辑可能失败——失败时改用msg.Nak()或msg.Term() - 不要在 goroutine 里直接操作
msg,因为msg不是线程安全的;需要传参时用msg.Data复制内容 - 若业务需幂等,别依赖 JetStream 去 dedup——它不提供全局 message ID 去重,得靠应用层用
msg.Header.Get("Nats-Msg-Id")做判重
连接断开后如何不丢消息或重复投递
JetStream 消费者默认启用 ReplayPolicy: nats.ReplayInstant,网络抖动恢复后会从最新消息开始读,导致断连期间的消息丢失。要保证不丢,必须设 ReplayPolicy: nats.ReplayOriginal,并配合 DeliverPolicy: nats.DeliverByStartTime 或 DeliverByStartSequence。
-
MaxDeliver设为 3~5,避免死信无限重试;重试后仍失败,应发到 dead-letter subject(如"dlq.orders")人工介入 - 客户端 reconnect 时,
nats.go默认会重用旧 consumer,但前提是Durable名一致且 server 未删 consumer;建议定期检查js.ConsumerInfo("orders", "orders-processor-v1")是否返回 404 - 生产环境务必设
Heartbeat: 30 * time.Second,否则长时间空闲连接会被 NATS server 断开,触发不必要的重播
真正难的不是写通代码,而是把 ack、replay、durable name、stream retention 这几处联动关系理清楚——少配一个,监控上看不出问题,压测或故障时才暴露。
golang免费学习笔记(深入):立即使用
在学习笔记中,你将探索golang的核心概念和高级技巧!










