NATS JetStream 深度实战:从 Subject 路由拓扑到流存储与 Exactly-Once 投递语义的工程全解
如果说 Kafka 把消息系统做成了"分布式提交日志",那么 NATS 走的是另一条路:先有一个极致简单的发布订阅内核(Core NATS),再在其上叠加一层可持久化的流媒体层(JetStream)。这两层的解耦方式,决定了 JetStream 在工程上呈现出与 Kafka 截然不同的取舍。本文从协议、存储、共识、投递语义四个层次拆开 JetStream,并给出可直接上生产的代码与调优参数。
一、Core NATS:一层被严重低估的路由内核
理解 JetStream 之前必须理解 Core NATS。它不做持久化、不做队列、不做 ack,只做一件事:基于 subject 的内容无关路由。
subject 是点分隔的字符串(eu.orders.created.42),核心是两种通配符:
*:匹配单个 token,eu.orders.*.42能匹配eu.orders.created.42>:匹配剩余全部 token,且只能出现在末尾,eu.orders.>匹配其下所有消息
服务端为每个客户端维护一棵 subject 订阅树(sublist),本质是 token 级的基数树(radix trie)。发布消息时先在 sublist 上做一次匹配,得到订阅者集合,然后零拷贝 fan-out:同一条消息的 []byte 缓冲区通过引用计数共享给所有匹配的客户端写出队列,而不是 per-subscriber 深拷贝。这是 NATS 单机能轻松跑上千万 msg/s 的根本原因——它没有做"分区",也不需要复制数据。
集群路由靠两件事:
- Route 连接(
cluster { routes = [...] }):节点间全连接网状拓扑,转发订阅兴趣而非消息副本 - 兴趣传播(subscription interest propagation):每个节点维护本地订阅表,新增订阅时向邻接节点广播
RS+协议帧。收到消息时先查本地订阅,再查跨节点兴趣表,只向"确实有人订阅"的节点转发
这个设计的直接后果:没有订阅者的消息会被静默丢弃,且路由开销极低。但也带来一个反直觉的陷阱:
# 错误:Core NATS 的 queue group 是"负载均衡 + 丢弃"语义
nats sub orders.> --queue=workers
# 如果所有 worker 都掉线,消息直接蒸发,不会堆积
这就是为什么生产环境里绝大多数业务应该直接用 JetStream,而不是 Core NATS + 自己做持久化。
二、JetStream 的两层 Raft:元流分离的工程价值
部署 JetStream 集群时,最容易被忽略的设计是:它运行着数量众多的独立 Raft Group,而不是一个全局共识实例。
- Meta Group(元数据组):每个集群只有 1 个(3/5 节点),负责 stream / consumer 定义的创建、删除、修改,以及 account 级的元数据。这部分数据量小、变更频率低
- Stream Group(数据流组):每个 stream 拥有自己独立的 Raft Group,负责该 stream 上消息的写入与复制
这个划分带来的工程红利是实打实的:
| 维度 | 单组 Raft(传统做法) | JetStream 多层 Raft |
|---|---|---|
| 故障爆炸半径 | 全局不可用 | 仅影响故障 stream |
| 写入并行度 | 单一 leader 瓶颈 | 不同 stream 的 leader 分散在不同节点 |
| 成员变更成本 | 重配置全局状态 | 单个 stream 独立变化迁移 |
| raft log 体积 | 全局单份膨胀 | 按 stream 切割,缩短 compaction 时间 |
这也意味着一个必须记住的运维规则:Replicas=3 是 per-stream 配置,不是全局配置。我见过最常见的生产事故,是把所有 stream 的 leader 都挤到了同一个节点上——因为创建 stream 时没有显式指定 placement。
# 反例:默认策略,集群自由选主,热点可能堆积
nats stream add ORDERS --subjects "orders.>" --replicas 3
# 正例:显式指定 leader 与目标节点分布,把热点 stream 打散
nats stream add ORDERS \
--subjects "orders.>" \
--storage file \
--replicas 3 \
--retention limits \
--max-msgs -1 --max-bytes 100GB --max-age 7d \
--placement '{"tags":["az-a"],"preferred":"nats-1"}' \
--sync 30s --ack
nats stream add AUDIT \
--subjects "audit.>" \
--replicas 3 \
--placement '{"tags":["az-a"],"preferred":"nats-2"}'
注意 --sync 与 --ack 的区别:--sync=30s 意味着 leader 每 30 秒等待 follower 落盘确认;而 --ack 要求所有副本确认后才向生产者返回 PubAck。金融级场景必须带 --ack,但吞吐量会掉到普通模式的 1/3 左右——这是同步复制的必然代价,没有魔法。
三、filestore:为什么 JetStream 可以不需要 partition
Kafka 需要 partition,是因为它把"顺序写 + 并行度 + 消费位点"三件事耦合在同一个抽象里。JetStream 的解耦方式是:把顺序写落到按块(block)组织的 filestore,把并行度交给 subject,把消费位点交给 consumer。
filestore 在磁盘上的组织是这个样子:
orders-stream/
├── 1.blk # 消息数据块,默认 8MB(可按 block_size 调)
├── 1.blk.idx # 索引块:msg_seq -> file_offset, timestamp, crc
├── 2.blk
├── 2.blk.idx
└── meta.sum # 流式 CRC + 加密 salt + stream 状态
写入路径的关键点:
- 消息按到达顺序追加进当前 blk 文件,达到
block_size后 roll over 到下一个 - 索引块
.idx与数据块同步写,记录了每条消息的(seq, offset, len, timestamp, CRC32C),支持 O(1) 的按序号随机读,也支持按时间戳二分查找 - 每次 roll over 时对已完成块做全块 CRC 汇总写入 meta.sum,跳过立即 fsync——因此 JetStream 的写路径天然是批处理语义,页缓存充当了它的 WAL 缓冲
这里有两条经验:
block_size不是越大越好。调大到 64MB 能减少 fd 数量与元数据开销,但清理(pruning)粒度变粗——一个 64MB 的块里哪怕只剩 1 条有效消息,整个块也不会被删除,会引发磁盘空间的"幽灵占用"。高 churn 场景建议保持默认 8MB 或下调至 2MB- 删除操作是标记 + 异步 pruning,不会立刻还空间。若队列频繁删除大 object store 内容,
nats stream report会看到 "Deleted" 数量持续上涨,此时要检查 pruner 是否被max_bytes之外的策略卡住
// Go 代码:用 Filestore + 同步复制 + 显式 ack 构建生产级 Stream
import "github.com/nats-io/nats.go"
js, _ := nc.JetStream()
_, err := js.AddStream(&nats.StreamConfig{
Name: "ORDERS",
Subjects: []string{"orders.>"},
Storage: nats.FileStorage,
Replicas: 3,
MaxAge: 7 * 24 * time.Hour,
MaxBytes: 100 * 1024 * 1024 * 1024,
Duplicates: 2 * time.Minute, // 幂等窗口:2 分钟内同 MsgId 去重
Placement: &nats.Placement{
Tags: []string{"az-a"},
Preferred: "nats-1",
},
SyncInterval: 30 * time.Second,
})
Duplicates 这个字段是很多人不知道的宝贝:它开启了一个 per-subject 的消息 ID 去重窗口。
四、Exactly-Once:MsgId 去重 + 双 ack 的组合拳
NATS 官方从不提供"端到端 exactly-once"的宣称,但通过两个机制的组合,可以在"读-处理-写"pipeline 上达到有效 exactly-once。
第一层:生产者幂等。 每条消息带一个业务唯一的 Nats-Msg-Id,服务端在 Duplicates 窗口内若发现同一 subject 已存在该 ID,会返回上一条消息的 PubAck 而不是再次追加:
// 幂等发布:服务端去重
ack, err := js.Publish("orders.created", data,
nats.MsgId("ORD-2026-0930-00001"), // 业务主键,如数据库主键 / 幂等号
)
if err != nil { return err }
fmt.Printf("stored at stream=%s seq=%d duplicate=%v\n",
ack.Stream, ack.Sequence, ack.Duplicate)
// 注意:重发时 ack.Duplicate == true,但 ack.Sequence 是首次存储的位置
第二层:消费者双 ack。 这是与 Kafka offset commit 差异最大的地方。Kafka 的 commitSync 是"我处理完了",而 JetStream 提供两级:
sub, _ := js.PullSubscribe("orders.created", "order-processor",
nats.AckExplicit(), // 必须显式 ack,禁用自动 ack
nats.AckWait(30*time.Second),
nats.MaxAckPending(1000), // 未确认消息上限,超过则暂停投递
nats.MaxDeliver(5), // 重试 5 次后进死信
)
for {
msgs, err := sub.Fetch(50, nats.MaxWait(5*time.Second))
if err != nil { continue }
for _, m := range msgs {
meta, _ := m.Metadata()
if meta.NumDelivered > 1 {
// 重投递:业务逻辑必须幂等,或者依赖上游 MsgId 去重
log.Printf("redelivery seq=%d attempt=%d", meta.Sequence.Stream, meta.NumDelivered)
}
// ---- AckSync 是关键:阻塞等待服务端同步落盘后再返回 ----
// 相比 AckAsync,AckSync 牺牲延迟换取"确认不丢"
if err := processOrder(m.Data); err != nil {
m.Nak() // 负确认,立即重投
continue
}
m.AckSync(nats.Context(context.Background()))
}
}
这里的三个调优参数决定了你的吞吐与可靠性边界:
MaxAckPending:本质是消费者的流控窗口。设太大(>10000)会会导致故障恢复时大量消息瞬间重投,形成恢复期雪崩;设得太小(<100)又会因为频繁等待 ack 浪费带宽。经验值:目标 QPS × P99 处理时延 × 2AckWait:重投判定时间。必须 > P99 处理时延,否则会出现"处理还在跑,消息已被重投"的自攻击MaxDeliver:配合 dead letter 死信 Topic 兜住"永远处理不了"的消息。设为 -1 表示无限重试,生产上一定要设有限值
第三层:业务侧幂等。 前两层保证"消息不会重复投递给中间件",但消费者处理成功后、AckSync 返回前崩溃,消息仍会重投。因此最终防线是业务表上的唯一键:
CREATE TABLE orders (
order_id VARCHAR(64) PRIMARY KEY, -- 即 Nats-Msg-Id
payload JSONB NOT NULL,
processed_at TIMESTAMPTZ
);
-- 消费时用主键冲突做幂等,而不是先 SELECT 再 INSERT(有竞态)
INSERT INTO orders (order_id, payload, processed_at)
VALUES ($1, $2, now())
ON CONFLICT (order_id) DO NOTHING;
只有把这三层叠起来,才敢在系统边界上说"不重复、不丢失"。
五、保留策略:为什么 interest-based retention 值得单独拿出来讲
Kafka 只有时间/大小保留。JetStream 多了一个极其巧妙的策略:Interest-based Retention(基于消费兴趣回收)
nats stream add RPC_REQUESTS \
--subjects "rpc.>" \
--retention interest \ # 所有 consumer 都消费完,消息即从磁盘删除
--storage file --replicas 3
在这个模式下,消息一旦被所有订阅它的 consumer 确认,就会标记删除。这把 JetStream 变成了一个超轻量的 RPC / 任务队列:既享受了持久化带来的可靠性(没被消费前不会丢),又不需要再跑一个夜间 cron 去清理过期 topic。
对于 router 类型的跨服务调用,这是比 Kafka consumer group 语义干净得多的方案——因为 Kafka 里你不消费也会一直被收费(磁盘),而 interest 策略下"没人在乎的消息"自动消失。
六、几个踩过的坑
1. nats bench 的数字别直接当容量规划依据。 benchmark 默认跑在 --msgs=1000 且 --replicas=1。加上 --replicas 3 --ack 后性能通常是除以 3 到 5,而不是除以 1.5。
2. Durable name 是共享位点标识。 多个服务误用了同一个 durable name,会以为各自独立消费,实则共用一个位点、互相"偷"消息。要隔离就用 nats.Durable() 加服务前缀。
3. Meta group leader 迁移是全集群卡顿点。 当 meta leader 所在节点 IO 抖动时,nats stream add 这类操作会超时,但已有 stream 的读写不受影响(因为是独立的 Stream Group)。这是分层 Raft 的最大价值:故障隔离。
4. KV / Object Store 都是 JetStream 的上层封装。 Bucket() 底层是一个 max_msgs_per_subject=1 的 stream,利用了 LSM 式的 per-subject 最新值覆盖。做大文件(>1GB)object store 时记得调 --max-bytes,不然 chunk 数量会撑爆 stream。
结语
JetStream 的哲学可以总结为一句话:用 subject 树替代 partition,用多层 Raft 替代全局共识,用 consumer 位点替代 consumer group。它放弃了 Kafka 那种"一个 topic 撑起所有"的粗粒度抽象,换来了更细的故障隔离粒度与更轻的运维负担。对于中小规模(单集群 <50 万 msg/s)但要求亚毫秒 P99 与运维极简的场景,JetStream 的综合性价比明显高于 Kafka——尤其是当你同时还需要 Core NATS 那套基于 subject 的通配符路由能力时。反过来,如果你真的需要单 topic TB 级吞吐、严格的跨分区事务,Kafka 的地位依然不可替代。技术选型的本质,从来都是挑最贴合自身约束的那一套,而不是挑纸面上"最好"的那一套。

发表评论 取消回复