第一部分 · 快速上手 #
1. NATS 是什么 #
NATS 是一个轻量、高性能的连接技术(connective technology),为分布式系统统一处理「寻址、发现、消息交换」。一句话:让发消息的人和收消息的人彻底解耦——生产者只把消息发到某个主题(subject),完全不关心谁在收、有几个在收。
标志性特点:基于主题寻址而非 IP:端口(天然 M:N、位置无关)、默认安全(TLS + 账户隔离 + JWT)、JetStream 内建持久化/KV/对象存储。
| 中间件 | 底层模型 | 一句话 |
|---|---|---|
| NATS | 主题路由 + JetStream 日志存储 | 云原生时代的轻量连接层 |
| Kafka | 分区日志 | 大数据流水线首选,但偏重 |
| RabbitMQ | Exchange→Queue(AMQP,Erlang) | 路由灵活的传统企业消息 |
2. 安装 #
2.1 Docker(最快) #
# -js 开启 JetStream;4222=客户端端口,8222=监控端口
docker run --name nats -p 4222:4222 -p 8222:8222 nats:latest -js
# 持久化到宿主机(重启不丢 JetStream 数据)
docker run --name nats -p 4222:4222 -p 8222:8222 \
-v /data/nats:/data nats:latest -js -sd /data
启动后浏览器开 http://localhost:8222 看监控页。
2.2 二进制 / 包管理器 #
brew install nats-server && nats-server -js # macOS
nats-server -js -m 8222 # Linux:下载 release 解压,-m 开监控端口
2.3 公共测试服务器 #
不想装也能玩:官方 nats://demo.nats.io:4222 直接连即可练手。
3. 用 CLI 先跑通(不写代码) #
官方 nats CLI 是学习和排障利器。开两个终端:
brew install nats-io/nats-tools/nats # 或下载 release
# 终端 A:订阅 # 终端 B:发布
nats sub "orders.>" nats pub orders.new "hello"
# 终端 A 立刻收到:[#1] Received on "orders.new": hello
请求/应答、队列组、JetStream:
nats reply "svc.time" --command="date" # 起应答服务
nats request "svc.time" "" # 发请求拿回复
nats sub "tasks" --queue workers # 多开几个体验组内分摊
nats stream add ORDERS --subjects "orders.>" # 建流
nats consumer add ORDERS proc # 建消费者
nats stream info ORDERS # 看条数/存储水位/消费者
4. Go 客户端:Core NATS #
go get github.com/nats-io/nats.go
package main
import (
"fmt"
"time"
"github.com/nats-io/nats.go"
)
func main() {
nc, err := nats.Connect("nats://localhost:4222",
nats.MaxReconnects(-1), nats.ReconnectWait(2*time.Second)) // 无限重连
if err != nil {
panic(err)
}
defer nc.Drain() // 优雅关闭:处理完缓冲再断开
nc.Subscribe("orders.*", func(m *nats.Msg) { // 普通订阅:都收
fmt.Printf("收到 %s: %s\n", m.Subject, m.Data)
})
nc.QueueSubscribe("tasks", "workers", func(m *nats.Msg) { // 队列组:组内一个收
fmt.Printf("本实例处理: %s\n", m.Data)
})
nc.Publish("orders.new", []byte("hello"))
if reply, err := nc.Request("svc.time", nil, time.Second); err == nil { // 请求/应答
fmt.Println("应答:", string(reply.Data))
}
nc.Flush()
time.Sleep(time.Second)
}
5. Go 客户端:JetStream(持久化) #
用现代 jetstream 子包(以 pull 消费者为主,官方当前推荐):
import "github.com/nats-io/nats.go/jetstream"
ctx := context.Background()
js, _ := jetstream.New(nc)
// 建/更新流:声明收哪些主题、怎么存、留多久
stream, _ := js.CreateOrUpdateStream(ctx, jetstream.StreamConfig{
Name: "ORDERS",
Subjects: []string{"orders.>"},
Retention: jetstream.LimitsPolicy,
Storage: jetstream.FileStorage, // 落盘
MaxAge: 7 * 24 * time.Hour,
})
// 发布:等服务端 ack(已落盘);MsgID 启用发布去重
ack, _ := js.Publish(ctx, "orders.new", []byte("o1"), jetstream.WithMsgID("order-1001"))
fmt.Println("已持久化 seq =", ack.Sequence)
// 持久消费者 + 消费
cons, _ := stream.CreateOrUpdateConsumer(ctx, jetstream.ConsumerConfig{
Durable: "processor", // 多实例共享此名 = 分布式队列
AckPolicy: jetstream.AckExplicitPolicy, // 必须逐条显式 ack
FilterSubject: "orders.new",
AckWait: 30 * time.Second, // 30s 不 ack 就重投
MaxDeliver: 5,
})
cc, _ := cons.Consume(func(msg jetstream.Msg) {
if err := handle(msg.Data()); err != nil {
if permanent(err) {
msg.Term() // 永久错误:别再投
} else {
msg.Nak() // 可重试:稍后重投
}
return
}
msg.Ack() // 成功:确认
})
defer cc.Stop()
因「至少一次」会重投,消费逻辑必须幂等(详见第三部分 §18)。
第二部分 · 核心模块原理 #
6. 主题寻址与 interest 路由 #
6.1 Subject(主题)= 消息的逻辑地址 #
大小写敏感、点分层级:orders.new、time.us.east。订阅支持通配符(只在订阅端用):
| 通配符 | 含义 | orders.us.new 能否命中 |
|---|---|---|
* |
匹配恰好一段 | orders.*.new ✅;orders.* ❌ |
> |
匹配其后所有段 | orders.> ✅;orders.us.> ✅ |
6.2 interest-based 路由:服务端怎么决定投给谁 #
服务端在内存里维护一张订阅兴趣表:谁订了哪些主题、属于哪个队列组。每条 PUB 进来只做一次主题匹配,只投给有兴趣的订阅者——没人订的主题,消息发出即丢(Core NATS 不存储)。
flowchart LR
P[发布者] -->|"PUB orders.new"| SRV
subgraph SRV["NATS 服务端(内存兴趣表)"]
T["orders.* → 订阅者A, 订阅者B<br/>orders.new → 订阅者C<br/>billing.> → 订阅者D"]
end
SRV -->|匹配命中| A[订阅者A]
SRV -->|匹配命中| B[订阅者B]
SRV -->|匹配命中| C[订阅者C]
SRV -. 不匹配,不投 .-> D[订阅者D]
关键认知:没有订阅者 = 消息凭空消失。这是 Core NATS「最多一次」语义的根源;要「不丢」需上 JetStream(§9)。
7. 简单的文本线协议 #
NATS 客户端与服务端之间是一套纯文本协议(像 HTTP 一样可读),几乎任何语言都能几十行实现客户端,也是看清「底层在干嘛」的窗口。核心操作(均以 CRLF ␍␊ 结尾):
| 操作 | 方向 | 含义 | 线格式 |
|---|---|---|---|
INFO |
S→C | 服务端自报家门 | INFO {json}␍␊ |
CONNECT |
C→S | 客户端鉴权与参数 | CONNECT {json}␍␊ |
PUB |
C→S | 发布 | PUB <subject> [reply-to] <字节数>␍␊<payload>␍␊ |
HPUB |
C→S | 带 header 发布 | HPUB <subject> [reply] <hdr长> <总长>␍␊<headers>␍␊␍␊<payload>␍␊ |
SUB |
C→S | 订阅 | SUB <subject> [queue] <sid>␍␊ |
UNSUB |
C→S | 退订(可设收够 N 条自动退) | UNSUB <sid> [max_msgs]␍␊ |
MSG |
S→C | 投递给订阅者 | MSG <subject> <sid> [reply-to] <字节数>␍␊<payload>␍␊ |
PING/PONG |
双向 | 心跳保活/探活 | PING␍␊ / PONG␍␊ |
+OK/-ERR |
S→C | 确认 / 报错 | +OK␍␊ / -ERR <原因>␍␊ |
sid:客户端为每个订阅生成的唯一编号;MSG投递时带回,据此派发到对应回调。reply-to:可选「回信地址」,请求/应答靠它实现。<字节数>:payload 长度预先声明,服务端按字节精确读取,零拷贝解析。
sequenceDiagram
participant C as 客户端
participant S as NATS 服务端
S-->>C: INFO {server_id, version, ...}
C->>S: CONNECT {user, pass, ...}
C->>S: SUB orders.* 90 (sid=90)
C->>S: PUB orders.new 5␍␊hello␍␊
S-->>C: MSG orders.new 90 5␍␊hello␍␊
S-->>C: PING
C->>S: PONG
telnet demo.nats.io 4222连上即见服务端推来的INFO,手敲SUB foo 1/PUB foo 5␍␊hello就能收发。
8. 通信与交互模式 #
NATS 的底层通信原语其实只有 pub/sub 一种,其余都是在它(或 JetStream)之上构建的用法。下面给出完整全景,并附每种怎么用 NATS 实现。
flowchart TB
ROOT["NATS 通信原语:Publish-Subscribe(一种)"]
ROOT --> CORE["Core NATS 基本模式"]
ROOT --> JS["JetStream 高层模式(持久化之上)"]
CORE --> M1["① 发布/订阅 扇出"]
CORE --> M2["② 队列组 负载均衡<br/>(pub/sub 的变体)"]
CORE --> M3["③ 请求/应答"]
M3 --> M3b["③b Scatter-Gather 一问多答"]
M3 --> M3c["③c Services(micro) 微服务框架"]
JS --> J1["④ Streaming 持久流/重放"]
JS --> J2["⑤ Work Queue 工作队列"]
JS --> J3["⑥ KV 键值+变更通知"]
JS --> J4["⑦ Object Store 对象存储"]
8.1 Core NATS 的三种基本模式 #
flowchart TB
subgraph PS["① 发布/订阅:一对多扇出"]
P1[发布者] --> S1[订阅者A] & S2[订阅者B] & S3[订阅者C]
end
subgraph QG["② 队列组:组内只投一个(负载均衡)"]
P2[发布者] --> Q{queue: workers}
Q -->|服务端随机挑一个| W1[实例1]
Q -.本次空闲.-> W2[实例2]
end
subgraph RR["③ 请求/应答:等回复"]
R1[请求方] -->|"PUB + reply-to=_INBOX.xyz"| R2[应答方]
R2 -->|"PUB 到 _INBOX.xyz"| R1
end
- ① 发布/订阅:每个订阅者各收一份,天然广播。
- ② 队列组:同名队列组的多订阅者,每条消息只随机投给一个——内建负载均衡,「不漏不重」,无需额外组件。严格说它是 pub/sub 的负载均衡变体,而非独立模式。
- ③ 请求/应答:请求方带临时回信主题
_INBOX.<随机>并临时订阅,应答方发回该主题。底层复用 pub/sub,常配超时。
// ① 发布/订阅
nc.Subscribe("orders.*", func(m *nats.Msg) { fmt.Println(string(m.Data)) })
nc.Publish("orders.new", []byte("hi"))
// ② 队列组:多实例只投一个
nc.QueueSubscribe("tasks", "workers", func(m *nats.Msg) { /* 处理 */ })
// ③ 请求/应答
nc.Subscribe("svc.time", func(m *nats.Msg) { m.Respond([]byte(time.Now().String())) })
reply, _ := nc.Request("svc.time", nil, time.Second)
8.2 请求/应答的扩展:Scatter-Gather(一问多答) #
普通 request/reply 只取第一个回复;Scatter-Gather(也叫 Request-Many) 发一个请求,让多个响应者都回复,请求方在一段时间窗口内收集多份——用于服务发现、并行查询、「问所有人再聚合」。
实现要点:自建一个临时 inbox 订阅,用 PublishRequest 把回信地址指过去,在超时窗口内循环收集:
inbox := nats.NewInbox()
sub, _ := nc.SubscribeSync(inbox)
nc.PublishRequest("svc.discover", inbox, nil) // 请求发到 svc.discover,回信到 inbox
deadline := time.Now().Add(200 * time.Millisecond)
for time.Now().Before(deadline) {
msg, err := sub.NextMsg(time.Until(deadline))
if err != nil {
break // 窗口内收完了
}
fmt.Println("收到一个响应:", string(msg.Data))
}
8.3 请求/应答的框架化:Services(micro) #
官方在 request/reply 之上封装了 micro 微服务框架:自动注册服务、暴露多端点、内建服务发现与调用统计——把「裸 request/reply」升级成规范的服务化交互。
import "github.com/nats-io/nats.go/micro"
srv, _ := micro.AddService(nc, micro.Config{Name: "calc", Version: "1.0.0"})
srv.AddEndpoint("add", micro.HandlerFunc(func(r micro.Request) {
r.Respond([]byte("result")) // 调用方仍用 nc.Request("add", ...) 调用
}))
8.4 JetStream 高层模式(持久化之上) #
这些是 Core 三种所没有的,依赖 JetStream(原理见 §10)。这里只给「模式视角」的最小用法:
js, _ := jetstream.New(nc)
// ④ Streaming:持久流 + 从历史任意位置重放(事件溯源)
stream, _ := js.CreateOrUpdateStream(ctx, jetstream.StreamConfig{
Name: "ORDERS", Subjects: []string{"orders.>"},
})
stream.CreateOrUpdateConsumer(ctx, jetstream.ConsumerConfig{
DeliverPolicy: jetstream.DeliverAllPolicy, // 从头重放全部历史
})
// ⑤ Work Queue:消费成功即删,一条只被处理一次(分布式任务队列)
js.CreateOrUpdateStream(ctx, jetstream.StreamConfig{
Name: "TASKS", Subjects: []string{"tasks.>"},
Retention: jetstream.WorkQueuePolicy,
})
// ⑥ KV:键值存储 + 变更通知(共享配置/状态)
kv, _ := js.CreateKeyValue(ctx, jetstream.KeyValueConfig{Bucket: "config"})
kv.Put(ctx, "color", []byte("blue"))
entry, _ := kv.Get(ctx, "color")
w, _ := kv.Watch(ctx, "color") // 监听变更
for u := range w.Updates() {
if u != nil {
fmt.Println("配置变了:", string(u.Value()))
}
}
// ⑦ Object Store:大文件分块存取
obj, _ := js.CreateObjectStore(ctx, jetstream.ObjectStoreConfig{Bucket: "files"})
obj.PutBytes(ctx, "a.bin", data)
got, _ := obj.GetBytes(ctx, "a.bin")
小结:「三种」只是 Core 层的入门切分(且队列组是 pub/sub 的变体)。加上 Scatter-Gather、Services,以及 JetStream 的 Streaming / Work Queue / KV / Object Store,NATS 支持的交互模式远不止三种——但它们都长在「pub/sub + JetStream 持久化」这两块基石上。
9. Core NATS vs JetStream #
| Core NATS | JetStream(-js 内建) |
|
|---|---|---|
| 存储 | 不存,转发即忘 | 持久化(内存/文件) |
| 订阅者不在线 | 丢失 | 上线后可补投/重放 |
| 投递保证 | 最多一次 | 至少一次(默认)/ 精确一次(可选) |
| 适用 | 实时、可容忍丢 | 必达、需重放、需历史 |
一句话:怕丢、要重放用 JetStream,纯实时不怕丢用 Core NATS。
10. JetStream 内部原理 #
10.1 Stream 与 Consumer #
- Stream(流):消息的持久化容器,按主题把消息追加写入一条 append-only 日志(类似 Kafka 分区),按保留策略留存。
- Consumer(消费者):流上的一个游标视图,记录「消费到哪了」。一个流可挂多个消费者,各自独立从不同位置、按不同过滤读,互不影响。
flowchart LR
PUB[发布者] -->|"orders.*"| ST["Stream: ORDERS<br/>append-only log: seq 1,2,3,4,5…"]
ST --> C1["Consumer A: 实时处理(从最新)"]
ST --> C2["Consumer B: 审计(从头重放)"]
ST --> C3["Consumer C: 队列组(多实例分摊)"]
记住这个**「日志 + 游标」模型**——它和 RabbitMQ 的「队列 + 出队即删」模型是后面所有差异的根。
10.2 存储与高可用 #
- 存储:
Memory(快、重启即失)或File(落盘持久)。 - 复制因子 R:
R=1无容错;R=3容忍挂 1 台(推荐);R=5容忍挂 2 台。靠 Raft 在集群复制。
10.3 保留策略(消息何时删) #
| 策略 | 行为 | 场景 |
|---|---|---|
| Limits(默认) | 按上限(时长/大小/条数)保留,到顶按 Discard 淘汰 | 通用事件流、可重放 |
| WorkQueue | 消息被成功消费后即删 | 任务队列(一条只处理一次) |
| Interest | 只在还有消费者关心时保留 | 节省空间的事件分发 |
配合 Discard:到上限时 old(删最旧)或 new(拒新)。
10.4 消费者类型 & 确认机制 #
- Pull(现代默认):客户端按需拉取(带批量/流控),易水平扩展,多 worker 共享一个 durable 即分布式队列。Push:服务端主动推到投递主题。
- Ack 回执(至少一次的核心):
stateDiagram-v2
[*] --> 待确认
待确认 --> 完成: Ack ✅ 成功
待确认 --> 重投: Nak ⏳ 可重试错误
待确认 --> 重投: AckWait 超时未回执
待确认 --> 待确认: InProgress 🔄 处理中,续期
重投 --> 待确认: 退避后再投
重投 --> 丢弃: 重投达 MaxDeliver 上限
待确认 --> 丢弃: Term ☠️ 永久错误
完成 --> [*]
丢弃 --> [*]
10.5 投递语义与去重 #
- 至少一次(默认):服务端对发布 ack,但极端故障下消费端可能重复 → 消费端需幂等。
- 精确一次(可选):① 发布带
Nats-Msg-Id,服务端在去重窗口(默认 2 分钟)丢弃重复发布;② double-ack(AckSync等服务端确认收到 ack)防 ack 丢失导致重投。
10.6 派生能力 #
同一存储引擎还派生 KV 键值存储、Object Store 对象存储、流镜像/源(容灾)、Subject 映射、消息调度(延迟,2.12+)——本质都是「特殊配置的 Stream」。
11. 集群与拓扑 #
flowchart TB
subgraph SC["Supercluster(跨地域,经 Gateway 连接)"]
subgraph CL1["集群 A(全互联 mesh)"]
A1[server] <--> A2[server] <--> A3[server]
A1 <--> A3
end
subgraph CL2["集群 B"]
B1[server] <--> B2[server]
end
CL1 <===>|Gateway| CL2
end
LEAF["Leaf Node(边缘/IoT)"] -.-> CL1
- Cluster:多台 server 全互联 + gossip 互传订阅兴趣,客户端连任意一台即可全局通达;某台挂了自动重连其它台。
- Supercluster:用 Gateway 连接多集群(跨数据中心),按需路由、就近优化。
- Leaf Node:边缘/内网轻量节点,把本地流量安全桥接到上游,断连可缓冲。
对应用透明:从单机到全球超级集群,nats.Connect 与 pub/sub 代码一字不改。
12. 安全 #
TLS 加密;鉴权支持 用户名密码 / token / NKey(Ed25519 签名) / 去中心化 JWT + Account;不同 Account 主题命名空间天然隔离,互通用 Export/Import 显式授权——适合多租户。
第三部分 · 十大核心问题:NATS vs RabbitMQ #
这部分用消息中间件的经典难题,横向对比两者的底层解法。
13. 先理解两者的底层模型差异(理解后面全部问题的钥匙) #
flowchart LR
subgraph RMQ["RabbitMQ:智能 Broker · 队列模型(AMQP / Erlang)"]
PR[Producer] --> EX{Exchange<br/>direct/topic/fanout/headers}
EX -->|binding| Q1[(Queue 1<br/>消息存这, 出队即删)]
EX -->|binding| Q2[(Queue 2)]
Q1 --> CR1[Consumer]
end
subgraph NATS["NATS:简单路由 · 日志模型"]
PN[Producer] -->|subject| LOG[("JetStream Stream<br/>append-only 日志, 按保留策略留存")]
LOG --> CN1[Consumer 游标A]
LOG --> CN2[Consumer 游标B]
end
| 维度 | RabbitMQ | NATS |
|---|---|---|
| 路由 | Exchange + binding 在 broker 内做富路由 | Subject 字符串匹配,极简极快 |
| 存储单元 | Queue:消息进队列、被消费即删 | Stream:append-only 日志、消费不删(按保留策略) |
| 消费 | broker 把消息推/派发给消费者,消费即出队 | 消费者持游标主动拉/被推,多消费者各读各的 |
| 一条消息多方消费 | 需绑多个 Queue(每队列各一份) | 多 Consumer 共享同一份日志即可 |
| 实现 | Erlang,状态存 Mnesia | Go,JetStream 用 Raft |
核心一句:RabbitMQ 是「队列:投进去、取出来就没了」;NATS JetStream 是「日志:追加进去、谁读到哪自己记」(更像 Kafka)。下面每个问题的差异,几乎都能回溯到这一点。
14. 如何解决消息积压(瞬时 & 长期) #
瞬时积压(突发流量,消费者临时跟不上):
- NATS:JetStream 把消息落盘进 Stream 当缓冲,消费者按自己节奏拉;pull consumer 用
MaxAckPending(在途未 ack 上限)做背压,不会被冲垮。(Core NATS 则触发 slow consumer 保护:连接 pending 超限直接丢消息并告警,保护 broker。) - RabbitMQ:消息堆在 Queue 里;用 prefetch(QoS) 限制每消费者在途条数 + 加并发消费者;内存/磁盘到水位触发 credit-based flow control,阻塞生产者给消费侧喘息。
长期积压(消费能力长期 < 生产):
- NATS:靠 Limits(MaxAge/MaxBytes/MaxMsgs)+ Discard 给日志封顶防磁盘爆;监控 consumer
num_pending(lag);扩消费者(共享 durable 的 pull group)提速;集群分担。 - RabbitMQ:Lazy Queue / Quorum Queue 尽快把消息落盘减内存压力防 OOM;设队列
max-length+overflow(drop-head / reject-publish / dead-letter);加消费者或限流上游。
模型差异:NATS 积压 = 日志多占磁盘 + 游标落后,老消息可按保留策略安全淘汰;RabbitMQ 积压 = 消息压在队列吃内存,核心手段是 lazy queue 落盘 + flow control 反压生产者。
15. 如何解决「消息乱序导致的业务问题」 #
乱序不一定要消灭,关键是别让乱序破坏业务正确性。两者思路一致,都靠业务层消解:
- 带版本/时间戳 + 状态机:消息携带
version或occurred_at,消费端只接受比当前更新的,丢弃过期的(last-write-wins)。 - NATS:每条消息天然有 stream sequence + timestamp,消费端可直接拿来判序、去重旧更新。
- RabbitMQ:消息自带业务时间戳/版本号,消费端幂等 + 版本比较。
中间件只能「尽量不产生乱序」(§16),乱序带来的业务影响最终靠幂等 + 版本号在消费端兜住。
16. 如何保证消息顺序性 #
乱序根因:多消费者并发 + 失败重投/requeue。要严格有序,必须牺牲并发,走「单一有序通道 + 单消费者」:
| NATS | RabbitMQ | |
|---|---|---|
| 有序前提 | 单 Stream 内按写入 seq 严格有序 | 单 Queue 内 FIFO |
| 单消费者 | 单个 Consumer 顺序消费 / OrderedConsumer 保证客户端按序收 |
Single Active Consumer(一队列同时只一个消费者活跃) |
| 按 key 分区有序 | 用 subject 分流(如 orders.{id}),同 key 落同 stream 单消费者 |
Consistent Hash Exchange 插件把同 key 路由到固定队列 |
| 破坏顺序的操作 | queue group / 多 worker 并发;Nak 重投插队 | 多消费者轮询;nack-requeue 把消息放回队尾 |
共同规律:顺序 = 单分区 + 单消费者,与吞吐天然矛盾。NATS 用 subject 当分区键,RabbitMQ 用队列(+一致性哈希)当分区键。
17. 如何保证消息不丢失(生产→Broker→消费 三段都要保) #
flowchart LR
P[生产者] -->|① 发送确认| B[(Broker 存储)]
B -->|② 持久化+副本| B
B -->|③ 消费确认| C[消费者]
| 环节 | NATS | RabbitMQ |
|---|---|---|
| ① 生产→Broker | JetStream Publish 等 server ack(确认落盘),失败重试 |
Publisher Confirms(broker 持久化后回 ack);旧 AMQP 事务(慢,少用) |
| ② Broker 存储 | File 存储 + R=3(Raft 复制) | durable 队列 + persistent 消息(deliveryMode=2)+ Quorum/镜像队列 |
| ③ Broker→消费 | 消费者显式 ack,ack 前不删,超时重投 | 消费者手动 ack,未 ack 重新入队 |
机制几乎一一对应:publish confirm ≈ publish ack;durable+persistent ≈ file storage;quorum/Raft ≈ R=3;manual ack 两边都有。任一段用了「不确认 / 不持久 / autoAck」就会丢(详见 §22)。
18. 如何做消息分发 #
- RabbitMQ(路由在 Exchange):Producer 发到 Exchange,按类型路由——
fanout广播到所有绑定队列、direct精确匹配 routing key、topic通配匹配、headers按 header;进入 Queue 后,多个消费者默认 round-robin 轮询派发(受 prefetch 调节)。 - NATS(路由在 subject):服务端按 subject 匹配兴趣表投递;扇出靠多订阅者各收一份,负载均衡靠 queue group(Core)或多 worker 共享一个 pull consumer(JetStream work-sharing)。
flowchart TB
subgraph R["RabbitMQ:Exchange 决定路由"]
e{"topic exchange"} -->|"order.*"| q1[("队列1")]
e -->|"*.vip"| q2[("队列2")]
q1 --> rr["消费者们轮询<br/>(prefetch 调节)"]
end
subgraph N["NATS:subject 决定路由"]
s["subject: orders.vip"] --> sub1["订阅者<br/>(扇出,各收一份)"]
s --> grp{"queue group<br/>组内只投一个"}
end
RabbitMQ 路由能力更「富」(四种 exchange + binding),代价是配置复杂;NATS 路由是字符串匹配,简单到极致、极快。
19. 如何防止重复消费 #
「至少一次」必然可能重复,消费端幂等是根本;broker 能帮一部分:
| NATS | RabbitMQ | |
|---|---|---|
| Broker 端去重 | 发布去重(Nats-Msg-Id + 去重窗口默认 2min)防生产重发;double-ack 防 ack 丢失重投 |
无内建去重 |
| 消费端(必做) | 业务幂等:去重表 / 唯一键 / 状态机 | 业务幂等:消息 message-id + 去重存储 |
// 消费端幂等通用范式:抢插去重键,已存在则直接 ack 跳过
res, _ := db.Exec(`INSERT INTO processed(consumer, msg_id) VALUES($1,$2)
ON CONFLICT DO NOTHING`, "proc", msgID)
if n, _ := res.RowsAffected(); n == 0 {
msg.Ack(); return // 重复消息,幂等跳过
}
NATS 多了一层 broker 端发布去重窗口;RabbitMQ 完全依赖业务幂等。无论哪个,最终都要消费端幂等兜底。
20. 如何实现事务消息,原理是什么 #
先分清两种「事务」:
(a) Broker 内多消息原子(发多条要么全成):RabbitMQ 有 AMQP tx.select/commit(性能差,基本不用,实务用 publisher confirms 替代);NATS 无此事务。
(b) 「业务库写入 + 消息发送」一致性(最常说的「事务消息」):两者都没有 RocketMQ 式的半消息两阶段事务。通用且推荐的方案是 事务发件箱(Transactional Outbox):
sequenceDiagram
participant S as 业务服务
participant DB as 业务库(业务表 + outbox 表)
participant R as 投递器(Relay)
participant MQ as NATS / RabbitMQ
Note over S,DB: ① 同一本地事务,原子提交
S->>DB: BEGIN
S->>DB: 写业务状态
S->>DB: INSERT outbox(event)
S->>DB: COMMIT
loop ② 后台异步
R->>DB: 取未发送的 outbox 行
R->>MQ: 发布
MQ-->>R: ack
R->>DB: 标记已发送
end
MQ-->>MQ: ③ 至少一次投递 → 消费端幂等去重
原理:把「跨两个系统的分布式事务」降级为「一个本地 DB 事务(业务写 + 事件登记原子)+ 至少一次投递 + 消费端幂等」。这样保证「状态变了 ⇔ 事件最终必发」,规避了两步之间崩溃丢事件的问题。
补充:NATS 还能用 JetStream 的
ExpectLastSubjectSequence(乐观并发)、KV 的 CAS 做轻量原子保证;RabbitMQ 侧通常配 publisher confirms + outbox。
21. 如何实现延迟消息 #
- NATS:
- 2.12+ 原生支持(JetStream Message Scheduler):Stream 开
AllowMsgSchedules,发布带 headerNats-Schedule(RFC3339 投递时刻)+Nats-Schedule-Target(到点出现的目标主题),服务端负责定时、投递、清理。延迟不得超过 Stream 的max_age。@every/cron 等周期能力见 2.14。 - 旧版/变通:消费时判断未到期 →
Nak带延迟重投;或外部调度器。
- 2.12+ 原生支持(JetStream Message Scheduler):Stream 开
- RabbitMQ:
- TTL + DLX(死信交换机):消息设 TTL,过期后死信到 DLX→目标队列。坑:队头阻塞——只能从队头死信,短 TTL 排在长 TTL 后会被拖延;适合固定几档延迟。
- delayed-message-exchange 插件:按各自延迟精确、有序投递;但调度状态存在收消息节点的 Mnesia,到点前该节点故障可能丢消息,单节点约 10 万 in-flight 软上限。
| NATS 2.12+ | RabbitMQ TTL+DLX | RabbitMQ 延迟插件 | |
|---|---|---|---|
| 任意 per-message 延迟 | ✅ header 驱动 | ⚠️ 队头阻塞 | ✅ 有序 |
| 可靠性 | 随 Stream 持久化/复制 | 随 quorum 队列可达数百万、集群安全 | 节点故障可能丢 |
22. 如果消息丢失了,可能的原因 #
| 原因 | NATS | RabbitMQ |
|---|---|---|
| 用了非持久路径 | 用了 Core 而非 JetStream | 队列非 durable / 消息非 persistent |
| 生产没等确认 | 没等 Publish 的 ack 就当成功 |
没开 Publisher Confirms |
| 消费端 autoAck | — | autoAck=true,收到即 ack,处理中崩溃即丢 |
| ack 早于处理 | ack 后才真正处理,中途崩溃 | 同左 |
| Broker 单点无副本 | R=1 且该节点磁盘损坏 |
非 quorum/镜像,节点宕机 |
| 积压被淘汰 | Limits + Discard old 删了老消息 / MaxAge 过期 |
max-length overflow drop-head / 队列 TTL 过期 |
| 消费太慢被保护性丢弃 | Core NATS slow consumer 超 pending 限被丢 | — |
| 延迟调度节点故障 | — | 延迟插件 Mnesia 状态随节点丢失 |
23. 如果消息被重复消费了,可能的原因 #
| 原因 | 说明(NATS / RabbitMQ 通用) |
|---|---|
| ack 丢失 / 超时 | 处理成功但 ack 没到 broker(NATS AckWait 超时 / RabbitMQ 连接断开未 ack)→ 重投 |
| 处理耗时 > ack 超时 | 实际还在处理却被判失败重投。NATS 应周期 InProgress 续期;RabbitMQ 调大超时/加心跳 |
| 消费端在 ack 前崩溃 | 消息已处理一部分副作用,重启后又投一次 |
| 生产端重试重发 | 网络抖动下生产者重发,又没带去重 ID(NATS 可用 Nats-Msg-Id 去重窗口消解) |
| 消费者重连/再均衡 | 未 ack 的在途消息在重连后被重新投递 |
治本三件套:消费端幂等 + 合理 AckWait/超时 + InProgress 续期 + 发布去重 ID。
小结 #
- 快速上手:
docker run nats -js起服务,natsCLI 收发,nats.go写 Core,jetstream子包写持久化(pull 消费者 + 显式 ack + 幂等)。 - 底层原理:基于主题寻址 + 内存兴趣表路由,文本线协议收发,Core「最多一次」/JetStream「至少一次~精确一次」;JetStream 是**「日志 + 游标」**模型,靠 Raft 复制、Ack 机制、MsgID 去重保证可靠。
- NATS vs RabbitMQ 的根:日志模型(消费不删、多游标) vs 队列模型(出队即删、富 Exchange 路由)。由此分化出积压(落盘缓冲+背压 vs lazy queue+flow control)、顺序(subject 分区 vs 队列+一致性哈希)、延迟(2.12 原生调度 vs TTL+DLX/插件)等几乎所有差异。
- 不变的真理:可靠投递三段都要确认+持久化;至少一次必然可能重复 → 消费端幂等是最后防线;跨库与消息的一致性用 事务发件箱而非「事务消息」。
延伸阅读:docs.nats.io、NATS 2.12 新特性、
jetstreamGo 包、RabbitMQ 延迟消息。