Go 基础体系 · 第 71/113 篇。示例统一基于 Go 1.26.4;核心片段可能省略 package 与 import,完整程序可直接按文中结构运行。
Go NATS 与 JetStream:Pub/Sub、持久化、Consumer 和 KV
本文以 Go 1.26.4、NATS Server v2.11.8 和 github.com/nats-io/nats.go v1.47.0 为基准。Core NATS 提供低延迟的在线消息分发和 Request/Reply;JetStream 在同一协议与 Subject 空间上增加持久 Stream、可重放 Consumer、发布确认、KV 和 Object Store。二者的交付边界不同:订阅者不在线时,Core NATS 默认不会替它保存消息;需要重启恢复、确认和回放时,应明确使用 JetStream。
NATS 的 API 很短,但可靠性不能由一次 Publish 调用推断。生产系统需要同时定义连接恢复、消息落盘、发布确认、消费 Ack、幂等、顺序、积压和故障恢复。本篇用文章发布事件贯穿这些环节。
1. Core NATS、JetStream 与集群各负责什么
客户端与任意 NATS 节点建立长 TCP 连接,服务端根据 Subject 把消息路由给本地或远端订阅者。普通订阅获得一份副本;同一 Queue Group 中只有一个成员获得该消息,用于水平分摊。集群节点通过路由交换订阅兴趣,Leaf Node 和 Gateway 分别用于边缘接入与跨集群连接。
publisher -> NATS connection -> subject routing -> subscriber
-> queue group 中的一个成员
-> JetStream stream -> durable consumer
Core NATS 的成功写入只说明服务端接受了协议帧,不代表已有持久副本。JetStream Stream 按 Subject 捕获消息,配置存储类型、保留策略、容量、副本数;Consumer 保存读取位置、待确认集合和重投状态。JetStream 元数据由集群一致性组管理,消息副本则按 Stream 的 Replicas 分布。业务应先决定是否允许离线丢失,再选择层次,而不是给 Core NATS 自行补一个内存重试队列。
2. Subject、通配符与 Queue Group
Subject 是以点分隔的层级名称,例如 articles.published.v1。* 只匹配一个 token,> 匹配后续全部 token。版本放在 Subject 还是消息信封中都可以,但要形成稳定契约;删除、改名或扩大通配订阅都需要兼容评审。
sub, err := nc.QueueSubscribe("articles.published.v1", "search-indexer", func(msg *nats.Msg) {
// 同一 queue group 中只有一个在线成员处理本条 Core NATS 消息。
})
if err != nil {
return fmt.Errorf("subscribe article events: %w", err)
}
defer sub.Unsubscribe()
Queue Group 解决在线负载均衡,不保存组的游标,也不等于 durable consumer。普通订阅者 A 和 Queue Group B 会各得到一份;B 的多个成员只由其中一个得到。权限应限制发布和订阅 Subject,尤其谨慎授予 >。Subject 不应包含用户秘密或无限基数数据,否则权限、监控与订阅表都会难以治理。
3. 连接建立、重连与优雅排空
nats.Conn 可并发复用,应用通常创建一个或少量按安全域隔离的连接。握手交换 INFO、CONNECT、PING/PONG,客户端从 INFO 学到集群地址。断线后库会在已知节点间退避重连;重连成功会恢复订阅,但断线窗口中的 Core NATS 消息不会倒流回来。
nc, err := nats.Connect(
servers,
nats.Name("article-service"),
nats.Timeout(3*time.Second),
nats.MaxReconnects(-1),
nats.ReconnectWait(500*time.Millisecond),
nats.DisconnectErrHandler(func(_ *nats.Conn, err error) {
slog.Warn("nats disconnected", "error", err)
}),
nats.ReconnectHandler(func(conn *nats.Conn) {
slog.Info("nats reconnected", "url", conn.ConnectedUrl())
}),
)
if err != nil {
return fmt.Errorf("connect nats: %w", err)
}
defer nc.Close()
无限重连不等于无限缓存。客户端存在 pending buffer 上限,长时间断网时发布可能失败或内存上涨;业务要根据错误拒绝请求、写 Outbox 或降级。退出时先停止接收新业务,再 Drain:它会排空发布和订阅回调后关闭连接。给排空设置外部期限,不能让卡住的 handler 永远阻止发布。
4. Core Publish、Flush 与 Request/Reply 生命周期
Publish 主要把消息放入客户端写缓冲。需要确认此前协议数据已被服务端读取时调用 FlushTimeout,其实现利用 PING/PONG 建立栅栏;这仍不是持久化确认。Request/Reply 创建唯一 inbox,发布带 Reply Subject 的请求,并等待一个或多个响应。
ctx, cancel := context.WithTimeout(parent, 800*time.Millisecond)
defer cancel()
reply, err := nc.RequestWithContext(ctx, "articles.lookup.v1", []byte(articleID))
if err != nil {
return fmt.Errorf("request article lookup: %w", err)
}
fmt.Println(string(reply.Data))
超时只表示调用方没收到响应,处理方可能已经完成写操作,所以非幂等命令要携带请求 ID。Request/Reply 适合轻量 RPC,但复杂契约、流式接口和标准状态模型可能更适合 gRPC。Core 消息慢消费者会触发 pending 增长甚至被服务端断开,必须监控 slow consumer,而不是把订阅回调当作无界任务队列。
5. 创建 Stream:保留、容量与副本
Stream 通过 Subject 集合捕获消息。LimitsPolicy 按消息数、字节数和时间淘汰;InterestPolicy 在没有 Consumer 兴趣后删除;WorkQueuePolicy 让一条消息由一个消费路径处理,并对重叠 Consumer 有限制。事件回放通常使用 Limits,工作队列才考虑 WorkQueue。
js, err := jetstream.New(nc)
if err != nil {
return fmt.Errorf("create jetstream context: %w", err)
}
stream, err := js.CreateOrUpdateStream(ctx, jetstream.StreamConfig{
Name: "ARTICLES",
Subjects: []string{"articles.*.v1"},
Storage: jetstream.FileStorage,
Retention: jetstream.LimitsPolicy,
MaxAge: 7 * 24 * time.Hour,
MaxBytes: 20 << 30,
Replicas: 3,
})
if err != nil {
return fmt.Errorf("configure article stream: %w", err)
}
_ = stream
配置更新不是都能无损完成,应由部署工具声明式管理并检查服务端返回。Replicas: 3 能容忍有限节点故障,但前提是节点和存储故障域独立。文件存储降低进程重启丢失风险,不等于备份;误删、凭据泄漏和跨集群灾难仍需快照、复制或源流方案。
6. 发布确认、去重与不确定结果
JetStream Publish 成功会返回包含 Stream 和 Sequence 的 PubAck,说明领导者接受并按当前复制策略提交。若调用超时,结果是未知:消息可能未写入,也可能已写入但确认丢失。生产者应携带稳定消息 ID,重试时保持不变,让 Stream 在配置的 duplicate window 内去重。
msg := nats.NewMsg("articles.published.v1")
msg.Header.Set("Nats-Msg-Id", event.ID)
msg.Header.Set("Content-Type", "application/json")
msg.Data = payload
ack, err := js.PublishMsg(ctx, msg)
if err != nil {
return fmt.Errorf("publish event %q: %w", event.ID, err)
}
slog.Info("event persisted", "stream", ack.Stream, "sequence", ack.Sequence)
去重窗口不是永久幂等表,窗口过后同一 ID 仍可能写入;它也无法把数据库业务提交与发布原子化。需要端到端可靠时,把业务变更与 Outbox 写入同一数据库事务,由 Relay 发布。异步发布提高吞吐,但必须消费每个 future 的结果、限制 pending 数量,并在关闭前等待确认。
7. Consumer 类型、起点与拉取消费
Consumer 可以 ephemeral 或 durable,可以 push 或 pull。Durable 保存进度,适合服务重启;pull consumer 由客户端决定每批数量和等待时间,背压更直观。起点可选全部、最新、指定时间或序列,创建前要明确“新部署是否处理历史”。
consumer, err := stream.CreateOrUpdateConsumer(ctx, jetstream.ConsumerConfig{
Durable: "SEARCH_INDEXER",
FilterSubject: "articles.published.v1",
AckPolicy: jetstream.AckExplicitPolicy,
AckWait: 30 * time.Second,
MaxDeliver: 8,
MaxAckPending: 256,
BackOff: []time.Duration{
500 * time.Millisecond,
2 * time.Second,
10 * time.Second,
time.Minute,
},
})
if err != nil {
return fmt.Errorf("configure search consumer: %w", err)
}
同名 durable 的配置必须与部署预期一致。多个实例绑定同一 pull durable 可以共享工作,但并发处理会改变完成顺序。MaxAckPending 是服务端在途闸门;批量大小和本地 worker 数还要受数据库连接池、外部 API 配额约束。
8. Ack、Nak、InProgress 与重投
显式 Ack 的正确边界是本地副作用成功提交之后。处理时间超过 AckWait,服务端会重投;可预计的长任务可发送 InProgress 延长等待。瞬时错误可 NakWithDelay,永久无效消息可 Term 停止继续投递并进入告警/归档流程。
func handle(ctx context.Context, msg jetstream.Msg, store *Store) error {
metadata, err := msg.Metadata()
if err != nil {
return fmt.Errorf("read message metadata: %w", err)
}
if err := store.ApplyEvent(ctx, msg.Data()); err != nil {
if errors.Is(err, ErrInvalidEvent) {
return msg.TermWithReason("invalid event")
}
return msg.NakWithDelay(backoff(metadata.NumDelivered))
}
if err := msg.Ack(); err != nil {
return fmt.Errorf("ack message: %w", err)
}
return nil
}
本地提交后 Ack 丢失会重复,因此 handler 必须幂等。反过来先 Ack 再提交会在进程崩溃时永久丢失。达到 MaxDeliver 后消息不应悄悄消失:监控 advisory,保存 event ID、Consumer、投递次数和最后错误,并提供受审计的重放流程。
9. 幂等、顺序与并行处理
常用幂等方式是在同一数据库事务中插入 (consumer, event_id) 唯一键并更新业务表。重复键表示该消费者已提交,可直接 Ack。对于实体状态事件,再带 aggregate_id 和递增 version,只允许 version = current + 1,旧版本忽略,缺口则延迟重试或触发回补。
JetStream 为 Stream 分配全局序列,但多个 Subject、并发 Consumer、重投和网络延迟会让业务完成顺序不同。必须严格按文章 ID 有序时,可将 ID 映射到固定分片 Subject,每个分片串行处理;代价是并行度受分片数限制且热点实体仍会阻塞。不要用全局单 Consumer 换取不必要的全局顺序。
“Exactly Once”需要限定范围:发布去重加双重 Ack 能减少重复传输,无法自动让数据库、邮件或支付只产生一次业务效果。业务层仍依赖幂等记录、唯一约束和可查询状态。
10. 背压、积压与容量计算
pull consumer 应使用有界批次:获取 N 条,交给固定 worker,等待空位后再取。若平均处理速率低于到达速率,积压必然增长;增加 worker 只有在下游还有容量时有效。可用 积压时间约等于 pending / 稳态处理速率 估算恢复窗口。
messages, err := consumer.Fetch(32, jetstream.FetchMaxWait(2*time.Second))
if err != nil {
return fmt.Errorf("fetch article events: %w", err)
}
for msg := range messages.Messages() {
select {
case jobs <- msg:
case <-ctx.Done():
return ctx.Err()
}
}
if err := messages.Error(); err != nil {
return fmt.Errorf("consume batch: %w", err)
}
本地 channel 必须有容量依据,并与 MaxAckPending 配合。遇到下游限流时降低 fetch 频率或暂停拉取,不要让 NATS 重投、本地重试和 HTTP 重试叠加成风暴。消息大小上限、Stream MaxBytes、磁盘写入速率和副本网络流量要一起做容量测试。
11. 持久化、恢复与集群故障
FileStorage 把消息写到磁盘,MemoryStorage 适合可丢失或可重建数据。节点退出时 Stream leader 转移;在法定副本不可用时,宁可暂时拒绝需要一致确认的写入,也不能假装持久成功。恢复演练应覆盖单节点重启、leader 故障、磁盘满、网络分区以及整个可用区丢失。
nats stream info ARTICLES
nats consumer info ARTICLES SEARCH_INDEXER
nats stream report
nats server check jetstream --server nats://127.0.0.1:4222
磁盘高水位会影响写入和复制,需在触顶前告警。备份必须验证恢复到隔离环境,而不是只确认快照文件存在。跨域灾备可用 mirror/source 或应用级重放,但要明确 RPO、RTO、Subject 权限和故障切换后的双写冲突。
12. KV、Watch 与能力边界
JetStream KV 基于 Stream,提供 revision、CAS、历史和 Watch,适合配置快照、轻量协调和服务状态。它不是支持多行事务、复杂查询和任意一致性约束的关系数据库。Watch 可能先交付历史再进入实时阶段,消费者必须处理删除标记、重连和重复 revision。
bucket, err := js.CreateOrUpdateKeyValue(ctx, jetstream.KeyValueConfig{
Bucket: "ARTICLE_FLAGS",
History: 5,
TTL: 24 * time.Hour,
Replicas: 3,
})
if err != nil {
return fmt.Errorf("configure flags bucket: %w", err)
}
revision, err := bucket.Update(ctx, "editor.enabled", []byte("true"), expectedRevision)
CAS 失败应重新读取并按业务决定合并,而不是无限覆盖。锁和选主还要处理租约超时、持有者暂停与 fencing token;仅凭“键存在”无法阻止旧持有者继续写外部系统。
13. 诊断、监控与测试
连接层观察连接数、重连、RTT、慢消费者和 pending bytes;JetStream 观察 Stream 消息/字节/副本、Consumer pending、ack pending、redelivery、等待拉取数和 leader 变化。业务指标包括发布确认延迟、端到端事件年龄、重复率、永久失败数和处理耗时。日志至少带 event ID、Subject、Stream sequence、Consumer 与投递次数,但不记录敏感 payload。
nats bench articles.load.v1 --pub 4 --sub 4 --msgs 100000 --size 1024
go test ./...
go test -race ./...
go test -run TestJetStreamRedelivery -count=20 ./integration/...
集成测试应启动真实 v2.11.8 Server,验证发布确认、进程在提交后 Ack 前崩溃、Consumer 重建、积压恢复和 Drain。性能测试使用真实消息大小与下游耗时,分别记录 Core 和 JetStream,不能用无持久化的本机 benchmark 推导三副本生产吞吐。
14. TLS、账户权限与生产部署
生产跨主机连接启用 TLS,按服务使用 NKey/JWT、用户凭据或 mTLS 身份。Account 隔离租户和 Subject 空间,export/import 明确跨账户共享。权限分别限制 publish、subscribe、响应 inbox 和管理 API;客户端凭据以文件或密钥系统注入并支持轮换。
authorization {
users = [
{ user: article_relay, password: $ARTICLE_RELAY_PASSWORD,
permissions: { publish: ["articles.*.v1"], subscribe: ["_INBOX.>"] } }
]
}
jetstream { store_dir: "/var/lib/nats/jetstream", max_file_store: 200GB }
节点使用独立持久盘和故障域,限制监控端口访问,设置文件描述符与磁盘告警。滚动升级先确认副本健康,一次只排空一个节点,等待集群恢复再继续。NATS 适合低延迟服务通信、事件和边缘拓扑;若主要需求是超长历史、按分区大规模扫描和成熟数据湖生态,应同时评估 Kafka。若只需复杂企业路由和每队列治理,可比较 RabbitMQ。选型最终取决于故障语义和运维能力,而不是客户端 API 行数。
系列导航与关联阅读
- 系列入口:Go 完整技术体系学习路线:从语法、并发到框架、中间件与 AI
- 上一篇:Go RocketMQ 实战:Producer、Consumer、顺序与延迟消息
- 下一篇:Go NSQ 基础:Topic、Channel、消费者与存量系统使用边界
- 延伸:Go 可靠消息统一设计:Outbox、幂等、重试、顺序与死信
- 延伸:Go WebSocket 与 SSE:实时通信、心跳、背压和断线恢复
- 延伸:Go 服务发现与配置中心:etcd、Consul、Nacos 的正确边界
官方资料
本文依据 Go 官方规范、标准库文档和 Go 官方博客重新梳理;正文与示例由 WR BLOG 编写。

评论
0 条讨论