NATS消息队列

第一部分 · 快速上手 #

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.newtime.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(落盘持久)。
  • 复制因子 RR=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-ackAckSync 等服务端确认收到 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)提速;集群分担。
  • RabbitMQLazy Queue / Quorum Queue 尽快把消息落盘减内存压力防 OOM;设队列 max-length + overflow(drop-head / reject-publish / dead-letter);加消费者或限流上游。

模型差异:NATS 积压 = 日志多占磁盘 + 游标落后,老消息可按保留策略安全淘汰;RabbitMQ 积压 = 消息压在队列吃内存,核心手段是 lazy queue 落盘 + flow control 反压生产者。

15. 如何解决「消息乱序导致的业务问题」 #

乱序不一定要消灭,关键是别让乱序破坏业务正确性。两者思路一致,都靠业务层消解:

  • 带版本/时间戳 + 状态机:消息携带 versionoccurred_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 Publishserver 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,发布带 header Nats-Schedule(RFC3339 投递时刻)+ Nats-Schedule-Target(到点出现的目标主题),服务端负责定时、投递、清理。延迟不得超过 Stream 的 max_age@every/cron 等周期能力见 2.14。
    • 旧版/变通:消费时判断未到期 → Nak 带延迟重投;或外部调度器。
  • 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 起服务,nats CLI 收发,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.ioNATS 2.12 新特性jetstream Go 包RabbitMQ 延迟消息