NATS Core、JetStream 与分布式事件通信实战
从一条“发布成功但消费者没收到”的消息开始
开发环境里最容易出现这样的判断:发布端没有报错,日志打印了 published,于是大家认定消息已经可靠送达。几秒后消费者重启,消息却再也找不到。继续排查才发现,发布端使用的是 Core NATS,消费者发布时不在线;团队把“写入客户端 socket 成功”误当成了“消息已进入持久化队列”。
NATS 的 API 很轻,真正需要先建立的是两套不同的心智模型:
Core NATS 负责在线的 subject 路由。消息只发给当时存在且匹配的订阅者,不保存历史,也没有服务端 ack、重投和消费位点。JetStream 订阅普通 NATS subject 并把消息写入 stream。consumer 在 stream 之上保存投递进度、ack 状态和重投状态,允许离线恢复与历史重放。
这不是“同一个队列开不开持久化”的区别。Core 的可靠性来自在线连接、请求超时、客户端重连与业务补偿;JetStream 的可靠性来自存储确认、stream 副本、consumer 状态和应用幂等。架构设计一旦把两者混用,后面的监控、容量和故障恢复都会失真。
先认识 subject、订阅和 queue group
NATS 不要求预先创建 topic。生产者把消息发到 subject,服务端依据订阅兴趣图做路由。subject 由点号分隔 token,例如:
acme.order.dev.created
acme.order.dev.cancelled
acme.payment.dev.succeeded订阅端可以使用 * 匹配一个 token,使用 > 匹配剩余层级:
acme.order.dev.* # 匹配 created、cancelled,但不匹配 created.audit
acme.order.dev.> # 匹配该前缀下所有层级subject 区分大小写,拼写错误不会触发“topic 不存在”。它会成为一个没有订阅者的新路由,因此 subject 命名本身就是接口契约。建议把组织、系统、环境、实体和动作放在稳定位置,把高基数字段放进消息体,不要把用户 ID、订单 ID 直接展开成无限多 subject。
普通订阅是广播:同一 subject 上每个订阅者都收到一份。queue group 则是竞争消费:同一 subject、同一 queue group 中仅一个成员收到该消息。官方的 queue group 说明 还强调,扩缩容不需要修改服务端配置,成员加入和退出即可改变并发度。
# 终端 A:普通订阅者,所有副本都会收到
nats sub acme.order.dev.created
# 终端 B、C:同属 workers,每条消息只交给其中一个
nats sub acme.order.dev.created --queue acme.order.dev.workersqueue group 只改变在线分发,不增加持久化。所有成员同时离线时,Core 消息仍会消失。
request/reply 也是 Core 语义:请求方创建临时 inbox,把 reply subject 附在请求中;处理方响应该 inbox。没有响应者时,支持该能力的客户端会收到 no responders,而不是永远等待。因此请求调用必须同时设置超时和“没有服务实例”的错误处理。
把单节点环境从零跑起来
NATS Server 当前稳定基线可以固定在 2.14.3。官方 Docker 说明 定义了三个常见端口:4222 是客户端连接,6222 是集群 route,8222 是 HTTP 监控。单机开发只需要映射前两个中的客户端端口和本机监控端口,不能把监控端点直接暴露到公网。
先确认 Docker 和端口可用:
docker version
Get-NetTCPConnection -LocalPort 4222,8222 -ErrorAction SilentlyContinue安装 nats CLI 时,优先从 NATS CLI 发布页 取得与操作系统匹配的二进制;使用 Go 工具链也可以执行:
go install github.com/nats-io/natscli/nats@latest
nats --versionCLI 版本应进入开发机清单或工具镜像,不应由每个人长期追随 latest。下面的服务端镜像则明确固定版本:
不使用容器时,可从 nats-server 官方发布页 下载对应平台压缩包,核对发布页给出的 SHA-256 后把二进制放入受控工具目录:
unzip nats-server-v2.14.3-linux-amd64.zip
install -m 0755 nats-server-v2.14.3-linux-amd64/nats-server "$HOME/.local/bin/nats-server"
nats-server --version
nats-server --jetstream --store_dir "$HOME/.local/state/nats-dev" --http_port 8222原生进程适合调试配置加载、系统服务和文件权限;容器适合项目隔离与可重复清理。两者的消息语义没有差别,差别在数据目录、进程托管、网络暴露和升级回滚方式。生产环境无论使用 systemd、容器平台还是 Kubernetes,都要把配置、身份、数据盘、健康探针和优雅停止显式管理,不能长期运行一条手工命令。
docker pull nats:2.14.3
docker run -d --name nats-dev \
-p 127.0.0.1:4222:4222 \
-p 127.0.0.1:8222:8222 \
-v nats-dev-data:/data \
nats:2.14.3 \
-js -sd /data -m 8222关键参数并非启动装饰:
-js 启用 JetStream;没有它时,Core 可以工作,但所有 stream 命令都会失败。-sd /data 把 JetStream 文件存储放到已挂载 volume。省略 volume 后,删除容器也会删除消息和 consumer 状态。-m 8222 启用 HTTP 监控。它用于观测,不是带鉴权的管理控制台。
127.0.0.1: 限制宿主机监听地址,避免开发机在办公网络暴露无认证服务。
日志里应该同时出现客户端监听、HTTP 监控和 JetStream 启动信息:
docker logs nats-dev --tail 80
curl -fsS http://127.0.0.1:8222/healthz
curl -fsS http://127.0.0.1:8222/varz
nats --server nats://127.0.0.1:4222 server pinghealthz 返回成功只能证明进程健康;varz 展示版本、连接数和流量;真正的 JetStream 能力还要由 nats account info 或 stream 操作确认。
用 Compose 固化开发依赖
仓库中的开发依赖更适合写成 Compose,使端口、volume 和配置成为可审查文件:
services:
nats:
image: nats:2.14.3-alpine
container_name: acme-nats-dev
command:
- --jetstream
- --store_dir=/data
- --http_port=8222
ports:
- "127.0.0.1:4222:4222"
- "127.0.0.1:8222:8222"
volumes:
- nats-data:/data
healthcheck:
test: ["CMD-SHELL", "wget -q -O /dev/null http://127.0.0.1:8222/healthz || exit 1"]
interval: 5s
timeout: 2s
retries: 12
restart: unless-stopped
volumes:
nats-data:这里显式选择 Alpine 变体,是因为健康检查需要容器内的 wget。默认 Linux 共享标签可能指向只包含 nats-server 的 scratch 镜像,不能假定其中存在 shell、wget 或 curl。生产镜像若坚持使用 scratch,应由编排平台的外部 HTTP 探针访问 /healthz,不要为了探针临时向运行镜像塞入调试工具。
docker compose up -d
docker compose ps
docker compose logs --tail 80 nats健康检查失败时先在容器内确认监控端口是否监听,再检查 Compose 的 command 是否覆盖了镜像默认参数。不要直接删除 volume,因为它正是后续恢复实验的证据。
正向实验:先证明 Core 的在线路由
打开终端 A:
nats --server nats://127.0.0.1:4222 sub acme.order.dev.created终端 B 发布:
nats --server nats://127.0.0.1:4222 pub acme.order.dev.created '{"eventId":"evt-001","orderId":"ord-1001"}'终端 A 应看到 subject 和 JSON 载荷。此时能证明的是:客户端连通、subject 拼写一致、订阅在发布时在线。它不能证明离线保存、重投、重复抑制或消费恢复。
再启动两个 queue subscriber:
nats sub acme.order.dev.work --queue acme.order.dev.workers
nats sub acme.order.dev.work --queue acme.order.dev.workers连续发布十条:
nats pub acme.order.dev.work "job {{Count}}" --count 10两端收到数量未必严格各半,但总数应等于发布数,每条只落到一个在线成员。若总数减少,先查看客户端断连和 slow consumer;若两端都收到全部消息,通常是 queue group 名不同或根本使用了普通订阅。
反向实验:证明 Core 不会替离线消费者保存消息
停止所有 acme.order.dev.offline 订阅者。发布一条带唯一 ID 的消息。发布完成后再启动订阅。
nats pub acme.order.dev.offline '{"eventId":"evt-core-lost"}'
nats sub acme.order.dev.offline订阅端不会补收 evt-core-lost。这不是 broker 丢数据,而是 Core 的既定语义。实验完成后按 Ctrl+C 退出订阅即可,没有服务端对象需要清理。
如果业务要求“消费者停机后恢复处理”,应改用 JetStream;如果业务只是在线服务发现、低延迟 request/reply,且调用方已有超时和补偿,Core 反而更简单。
JetStream 的三个状态对象
JetStream 不是一个统一的“持久队列”开关,而是三层状态:
subject 仍然是消息入口和路由名称。stream 捕获一个或多个 subject,保存消息及其 stream sequence,并执行保留、淘汰、复制和存储策略。consumer 在 stream 上定义过滤条件、起始位置、投递方式、ack 策略、待确认上限和重投策略,并保存自己的消费序列。
一条消息可以被一个 stream 捕获,再由多个 consumer 独立读取。删除 consumer 不等于删除 stream 消息;删除 stream 则会同时删除它保存的消息和关联 consumer。排障时必须先问“消息是否进入 stream”,再问“consumer 是否看到、是否 ack”。
正向实验:创建可恢复的 stream 和 durable consumer
先创建 stream。官方 stream 配置说明 将 limits、workqueue、interest 定义为三种保留策略;本实验使用 limits,让消息是否删除只由数量、字节和时间限制决定,便于重放。
nats stream add ORDERS_DEV \
--subjects "acme.order.dev.events.>" \
--storage file \
--retention limits \
--discard old \
--max-age 1h \
--max-bytes 256MB \
--replicas 1 \
--defaults逐项理解这些值:
subjects 决定哪些消息被捕获;范围过宽会混入别的项目,范围过窄会出现“发布成功但 stream 没增长”。storage file 把消息写入文件存储;memory 延迟低,但进程退出后不能承担恢复承诺。retention limits 允许多个 consumer 独立重放。
discard old 在任一上限达到时淘汰最老消息。改成 discard new 后,新发布会收到容量错误,更适合不能静默丢旧数据的审计流。max-age、max-bytes 与 max-msgs 共同生效,最先触达的限制决定删除。replicas 1 只适用于单节点开发,不提供节点故障冗余。
用 JetStream publish 等待服务端存储确认:
nats pub acme.order.dev.events.created \
'{"eventId":"evt-js-001","orderId":"ord-1001"}' \
--header "Nats-Msg-Id:evt-js-001" \
--jetstream
nats stream info ORDERS_DEV预期证据是 publish 返回 stream 名和 sequence,stream info 的 message count 增加。只使用普通 Core publish 时,消息仍可被 stream 捕获,但发布者没有收到“已经写入 stream”的确认,无法区分无匹配 stream 与持久化成功。
创建 durable pull consumer:
nats consumer add ORDERS_DEV ORDER_WORKER_DEV \
--filter "acme.order.dev.events.created" \
--pull \
--deliver all \
--ack explicit \
--wait 10s \
--max-deliver 5 \
--max-pending 100 \
--replay instant \
--defaults
nats consumer next ORDERS_DEV ORDER_WORKER_DEV --count 1 --ack
nats consumer info ORDERS_DEV ORDER_WORKER_DEVAckExplicit 要求每条消息单独确认,也是 pull consumer 的可靠默认选择。AckWait 是一次处理的时间预算,超时会触发重投;MaxAckPending 达到上限后服务端暂停继续投递,防止消费者无限占用内存;MaxDeliver 到达后停止继续尝试,但消息不会自动变成传统 MQ 的 DLQ。
成功确认后,consumer info 中 acknowledgement floor 前移,ack pending 回到零。若消息还在 stream 中,这是 LimitsPolicy 的正常表现,不代表 ack 失败。
反向实验:制造 ack 超时与重投
再发布一条唯一消息:
nats pub acme.order.dev.events.created \
'{"eventId":"evt-js-redelivery","orderId":"ord-1002"}' \
--jetstream第一次拉取时故意不 ack:
nats consumer next ORDERS_DEV ORDER_WORKER_DEV --count 1 --no-ack等待超过 AckWait 后再次拉取:
nats consumer next ORDERS_DEV ORDER_WORKER_DEV --count 1 --ack
nats consumer info ORDERS_DEV ORDER_WORKER_DEV第二次应看到同一 event ID,消息元数据的 delivery count 增加;consumer info 的 redelivered 计数也会变化。若没有重投,检查第一次命令是否自动 ack、consumer 是否确实为 explicit ack,以及 AckWait 是否尚未到期。
这个实验揭示了两个生产事实:
JetStream 提供的是至少一次处理基础,超时、断连或 ack 丢失都可能造成重复。ack 只能证明 broker 收到了确认,不能替业务证明数据库事务已经提交。业务处理必须以 eventId 或业务幂等键去重。
应用处理时间可能超过 AckWait 时,不要盲目把窗口改得很大。客户端可以发送 in-progress ack 延长处理时间;更稳妥的做法是拆小任务、限制并发,并让耗时分布与 AckWait 对齐。
去重、双重确认和“恰好一次”的真实边界
JetStream stream 可以在一个滑动窗口内记住 Nats-Msg-Id。重复发布相同 ID 时,服务端返回确认但不会再次写入。官方的 JetStream 深入说明 将发布去重与 double ack 组合描述为 exactly-once 能力,但它不等于跨数据库、外部 HTTP 调用和消息系统的全局事务。
nats pub acme.order.dev.events.created '{"eventId":"evt-dedup-001"}' \
--header "Nats-Msg-Id:evt-dedup-001" --jetstream
nats pub acme.order.dev.events.created '{"eventId":"evt-dedup-001"}' \
--header "Nats-Msg-Id:evt-dedup-001" --jetstream
nats stream info ORDERS_DEV在 duplicate window 内,stream message count 只应增加一次。超过窗口后,同一 ID 可能再次进入 stream。因此生产者仍需稳定生成 ID,消费者仍需维护幂等记录;不能用一个很长的去重窗口替代业务状态机,因为窗口本身消耗内存且不是永久唯一索引。
对消费确认,普通 ack 发送后客户端可能在 broker 已处理 ack、响应却丢失时产生不确定性。支持 double ack 的客户端会等待服务端确认 ack 已被处理,缩小这个窗口。即便如此,应用数据库在 ack 前后崩溃仍可能造成“业务已提交但消息重投”或“先 ack 后业务未提交”。常用落地是本地事务写业务状态与 inbox 表,提交成功后再 ack;重投时 inbox 唯一键阻止重复副作用。
保留、清理和死信等价能力
三种 retention policy 对架构含义完全不同:
LimitsPolicy:事件日志与独立重放
消息由 age、bytes、count 等限制淘汰,不因 ack 删除。多个 consumer 可以保留各自进度,适合事件回放、审计窗口和多个下游。容量预算必须按“写入速率 × 平均消息大小 × 保留时长 × 副本数”估算,再加索引、文件和增长余量。
WorkQueuePolicy:共享工作队列
同一 subject filter 不能有重叠 consumer;消息被成功 ack 后从 stream 删除。达到 MaxDeliver 的消息仍留在 stream,需要运维或应用根据 advisory/API 找出并转存。它不是自动 DLQ。
InterestPolicy:只为已有兴趣保留
消息要被所有相关 consumer ack 后才删除。没有 consumer 表示没有兴趣,消息可能在发布后立即删除。因此必须先创建 consumer,再开放生产者流量。
NATS 没有一个自动生成的 ActiveMQ.DLQ。常见等价设计是:业务达到最大尝试次数后,将原消息、失败原因、首次时间、最后一次错误和原 stream sequence 发布到 acme.order.dev.dead,由独立 DEAD stream 保存;只有死信写入得到 JetStream pub ack 后,才终止原消息继续处理。这样做仍需防止“死信已写入但原消息未确认”产生重复死信,死信 event ID 必须可幂等。
顺序语义不能只看 subject
单个 stream 为消息分配单调递增 sequence,但端到端业务顺序还会被并发 consumer、重投、网络延迟和多 subject 打破。若订单事件必须按订单串行处理,可以采用:
以稳定键映射到有限数量的 shard subject,例如 orders.shard.00 到 orders.shard.31。每个 shard 只允许一个顺序处理通道,扩容通过增加 shard,而不是让任意并发消费者抢同一序列。消息体携带业务 version,消费者拒绝回退版本,并对缺号进入补偿路径。
不要为每个订单创建 stream 或 consumer。对象数量和元数据会迅速膨胀,导致控制面、内存与恢复时间不可控。
项目接入:连接生命周期与幂等处理
下面使用 Node.js 展示最小接入;npm install nats 会把解析出的版本写入 lockfile,团队应提交 lockfile 并通过依赖更新流程升级。
mkdir nats-demo && cd nats-demo
npm init -y
npm install nats// worker.mjs
import { connect, StringCodec } from "nats";
const sc = StringCodec();
const nc = await connect({
servers: process.env.NATS_URL ?? "nats://127.0.0.1:4222",
name: "order-worker-dev",
maxReconnectAttempts: -1,
reconnectTimeWait: 1000,
});
nc.closed().then((err) => {
if (err) console.error("NATS connection closed", err);
});
const sub = nc.subscribe("acme.order.dev.notifications", {
queue: "acme.order.dev.notification-workers",
});
const stop = () => void nc.drain();
process.once("SIGINT", stop);
process.once("SIGTERM", stop);
for await (const msg of sub) {
const body = JSON.parse(sc.decode(msg.data));
console.log("received", body.eventId);
}
await nc.closed();Core worker 适合可丢、可重建的在线通知。订单、计费这类需要离线恢复的任务,应绑定已经评审过的 durable consumer,不能让每个应用副本在启动时随意创建新 consumer。真正的提交边界也不是“JSON 解析成功”,而是业务数据库事务已经提交。
先在业务数据库建立 inbox 唯一键。它与业务更新处于同一个本地事务,负责把 JetStream 的至少一次投递收敛为业务上的一次生效:
CREATE TABLE message_inbox (
consumer_name VARCHAR(128) NOT NULL,
event_id VARCHAR(128) NOT NULL,
received_at TIMESTAMP NOT NULL DEFAULT CURRENT_TIMESTAMP,
PRIMARY KEY (consumer_name, event_id)
);下面的 worker 把项目数据库适配器抽象为 db.transaction() 和 tx.query()。关键顺序是:写 inbox、更新业务、提交数据库、最后 ack。暂时性异常使用延迟 nak;不可恢复的 schema 或权限错误只有在失败记录已经持久化后才能 term,因为 term 只停止继续投递,并不会自动把消息搬入传统 DLQ。
// jetstream-worker.mjs
import { connect, StringCodec } from "nats";
import { db, isPermanentMessageError, saveFailedMessage } from "./persistence.mjs";
const stream = process.env.NATS_STREAM ?? "ORDERS_DEV";
const durable = process.env.NATS_CONSUMER ?? "ORDER_WORKER_DEV";
const sc = StringCodec();
const nc = await connect({
servers: process.env.NATS_URL ?? "nats://127.0.0.1:4222",
name: durable,
maxReconnectAttempts: -1,
});
const js = nc.jetstream();
const consumer = await js.consumers.get(stream, durable);
const messages = await consumer.consume();
let stopping = false;
const stop = () => {
stopping = true;
messages.stop();
};
process.once("SIGINT", stop);
process.once("SIGTERM", stop);
for await (const msg of messages) {
try {
const event = JSON.parse(sc.decode(msg.data));
await db.transaction(async (tx) => {
const inserted = await tx.query(
`INSERT INTO message_inbox(consumer_name, event_id)
VALUES ($1, $2)
ON CONFLICT (consumer_name, event_id) DO NOTHING
RETURNING event_id`,
[durable, event.eventId],
);
if (inserted.rowCount === 0) return;
await tx.query(
`UPDATE orders SET projection_version = $1 WHERE order_id = $2`,
[event.version, event.orderId],
);
});
// 仅用于故障实验:制造“数据库已提交、ack 尚未发送”的崩溃窗口。
if (process.env.CRASH_AFTER_COMMIT === event.eventId) process.exit(23);
msg.ack();
} catch (error) {
if (isPermanentMessageError(error)) {
await saveFailedMessage(msg, error);
msg.term();
} else {
msg.nak(5_000);
}
}
}
if (stopping) await nc.drain();
await nc.closed();persistence.mjs 必须使用项目真实的连接池和事务 API,不能把上述 inbox 插入与业务更新拆成两个事务。多个 worker 同时收到同一 event ID 时,只有一个事务能插入唯一键并更新业务;其余事务识别重复后直接返回,随后仍要 ack,避免同一消息永久重投。
正向验证先启动 worker,再发布唯一事件:
nats pub acme.order.dev.events.created \
'{"eventId":"evt-app-001","orderId":"ord-1001","version":1}' \
--jetstream
nats consumer info ORDERS_DEV ORDER_WORKER_DEV数据库中应同时出现一条 message_inbox 记录和一次业务更新,consumer 的 acknowledgement floor 前移、ack pending 回到零。只看到应用日志不算完成,必须把数据库结果和 consumer 状态一起保留为证据。
反向验证使用测试事件制造提交后崩溃:先以 CRASH_AFTER_COMMIT=evt-app-crash 启动 worker,发布同名事件,确认进程以 23 退出且数据库事务已经提交;随后取消该环境变量并重启 worker。消息应再次投递,inbox 唯一键阻止第二次业务更新,worker 最终 ack;consumer info 应显示 redelivery 增加且 ack pending 回到零。若业务数据被更新两次,问题在事务或唯一键,而不在 JetStream。
连接对象应由应用生命周期统一持有。停止实例时先停止继续拉取,等待当前数据库事务结束并发送 ack/nak,再 drain 连接;直接结束进程会让未 ack 消息按策略重投,这是可靠性机制,不是异常。
项目配置建议拆开连接、安全和消息对象:
messaging:
nats:
servers:
- ${NATS_URL:nats://127.0.0.1:4222}
credentials-file: ${NATS_CREDS_FILE:}
connect-timeout: 2s
request-timeout: 3s
subject-prefix: acme.order.dev
stream: ORDERS_DEV
consumer: ORDER_WORKER_DEV
max-inflight: 100启动时只验证连接不够,还应检查 stream 的 subject、retention、storage、replicas 与 consumer 的 filter、ack policy、deliver policy 是否符合声明。对象已经存在但配置不一致时应失败启动或报警,不能静默沿用旧配置。
安全配置:账户比密码更重要
默认 NATS Server 没有认证和授权,只适合绑定 localhost 的一次性开发环境。共享环境至少要启用 TLS、应用身份和 subject 权限。官方 认证说明 定义了 user/password、token、NKey 与 JWT/operator 等模式;其中 account 形成独立 subject 与 stream 命名空间,跨 account 只有显式 import/export 才能通信。
小规模内部环境可先使用配置文件账户:
port: 4222
http: 127.0.0.1:8222
jetstream {
store_dir: "/data"
max_file_store: 20GB
max_memory_store: 512MB
}
accounts {
ORDER_APP {
jetstream: enabled
users: [
{
user: order_service
password: $ORDER_SERVICE_PASSWORD
permissions: {
publish: ["acme.order.prod.>"]
subscribe: ["acme.order.prod.>", "_INBOX.>"]
}
}
]
}
}真实密码不要写在仓库中的配置文件。通过密钥系统生成运行时文件或环境注入,并限制读取权限。生产环境更适合 operator/account JWT 与 .creds 文件,使账户委派、吊销和最小权限不依赖修改所有 server;凭证文件不能进仓库、镜像、日志、截图或工单正文。
TLS 需要同时覆盖客户端、route、gateway 和 leaf connection。官方 TLS 文档 提醒:集群与 gateway 连接始终验证证书,advertise 的主机名必须出现在证书 SAN 中。关闭证书校验只会把“能连通”换成中间人风险。
监控端点不会自动继承普通客户端的用户名密码。共享环境应绑定管理网、通过受控代理暴露或使用 HTTPS/mTLS,并限制 /connz、/subsz、/jsz 等包含拓扑和 subject 信息的响应。
从单节点走向生产拓扑
三节点集群:同一故障域内的基础形态
Core NATS 集群通过 route 共享订阅兴趣与连接负载,客户端连接任一节点都能与其他节点上的订阅者通信。JetStream 额外使用 Raft 管理 stream 和 consumer 复制。典型生产起点是三个节点、每个节点独占本地 SSD 和存储目录,stream 设置 replicas=3。
优点是可以容忍单节点故障并自动选主;代价是每条持久消息写入多个副本,网络与磁盘延迟进入确认路径。两个副本不能在网络分区中同时兼顾可用性与多数派安全,三个节点才有清晰多数派。官方 JetStream 配置 明确建议本地快速 SSD,避免 NAS/NFS,并要求每台 JetStream server 使用独立存储目录。
下面以三台独立主机 nats-1.internal、nats-2.internal、nats-3.internal 为例。三台机器都监听客户端端口 4222、route 端口 6222 和仅管理网可达的监控端口 8222。nats-1 的配置如下:
server_name: nats-1
port: 4222
http_port: 8222
jetstream {
store_dir: "/var/lib/nats/jetstream"
max_file_store: 100GB
max_memory_store: 2GB
}
cluster {
name: acme-orders
listen: 0.0.0.0:6222
routes: [
nats-route://nats-2.internal:6222
nats-route://nats-3.internal:6222
]
}nats-2 和 nats-3 不能复制同一个 server_name,也不能共享存储目录。它们只替换节点身份和 route 种子:
| 节点 | server_name | routes |
|---|---|---|
nats-1.internal | nats-1 | nats-2:6222、nats-3:6222 |
nats-2.internal | nats-2 | nats-1:6222、nats-3:6222 |
nats-3.internal | nats-3 | nats-1:6222、nats-2:6222 |
在每台主机创建独立数据目录并由 NATS 运行用户持有,再交给服务管理器启动:
sudo install -d -o nats -g nats -m 0750 /var/lib/nats/jetstream
sudo -u nats nats-server -t -c /etc/nats/nats.conf
sudo systemctl restart nats
sudo systemctl status nats --no-pagernats-server -t 只检查配置语法,不能证明 route、Raft 或磁盘可用。三台节点启动后,从管理网分别查看 /routez,每台都应看到另外两个节点;随后创建三副本 stream:
export NATS_URL='nats://nats-1.internal:4222,nats://nats-2.internal:4222,nats://nats-3.internal:4222'
curl -fsS http://nats-1.internal:8222/routez
curl -fsS http://nats-2.internal:8222/routez
curl -fsS http://nats-3.internal:8222/routez
nats --server "$NATS_URL" stream add ORDERS_HA \
--subjects 'acme.order.ha.>' \
--storage file \
--retention limits \
--replicas 3 \
--max-age 24h \
--max-bytes 1GiB \
--discard old \
--defaults
nats --server "$NATS_URL" stream info ORDERS_HAstream info 应列出一个 leader 和两个 current replica。这里的 100GB 是单节点上限,不是业务容量承诺;stream 的 1GiB 逻辑上限会在三个节点各保存一份,磁盘预算还要包含块、索引、临时追平和其他 stream。生产配置还必须合并前文的客户端、route TLS 和账户权限,监控端口不能暴露到业务网或公网。
单节点故障:允许短暂选主,不允许丢失确认语义
先记录 ORDERS_HA 的 leader、每个 replica 的 lag 和当前最后序号。在隔离演练环境停止 leader 所在节点:
nats --server "$NATS_URL" stream info ORDERS_HA
# 在上一步显示的 leader 主机执行
sudo systemctl stop nats客户端可能在选主窗口内遇到一次超时或重连;新 leader 产生后,向仍存活的节点发布唯一事件并要求 JetStream ack。下面假设最初停止的是 nats-1;实际演练必须按 stream info 的结果设置存活节点地址:
export SURVIVOR_URLS='nats://nats-2.internal:4222,nats://nats-3.internal:4222'
nats --server "$SURVIVOR_URLS" \
pub acme.order.ha.created \
'{"eventId":"evt-ha-one-down","orderId":"ord-ha-1"}' \
--jetstream
nats --server "$SURVIVOR_URLS" stream info ORDERS_HA合格证据是发布最终收到 stream/sequence 确认、leader 已迁移、一个 replica 显示 offline,而不是三台进程都 ready。恢复故障节点后继续观察,直到它重新变为 current 且 lag 回到零;若追平期间发布确认延迟持续恶化,应检查磁盘吞吐、route RTT 和待复制字节。
两节点故障:失去多数派必须明确失败
三副本 stream 同时失去两个节点后只剩一个副本,无法形成多数派。这个实验只能在隔离环境进行:保持第一台节点停止,再停止当前 leader,随后向剩余节点执行 JetStream 发布。下面继续假设新 leader 是 nats-2、最后存活节点是 nats-3;节点角色与假设不一致时必须替换地址,不能照抄主机名。
# 在当前 leader 主机执行
sudo systemctl stop nats
export LAST_NODE_URL='nats://nats-3.internal:4222'
nats --server "$LAST_NODE_URL" \
pub acme.order.ha.created \
'{"eventId":"evt-ha-no-quorum","orderId":"ord-ha-2"}' \
--jetstream预期是发布超时或返回 stream/leader 不可用,不能获得成功的持久化确认。剩余节点上的 Core NATS 仍可能接受在线 publish/subscribe,但这不代表 ORDERS_HA 可写,也不能把 Core 的发送成功日志当作补偿证据。恢复任意一个已停止节点形成多数派后,等待 leader 和 replica 状态稳定,再重新发布一个新的 event ID;不要盲目重试 evt-ha-no-quorum,除非生产者使用稳定消息 ID 或业务幂等键处理“请求失败但结果未知”的窗口。
判断集群是否健康不能只看三台进程存活。应同时看:
/routez 中 route 是否形成预期拓扑,RTT 与 pending bytes 是否异常。/jsz、stream info 中 leader、replica lag、offline replica 是否正常。发布确认延迟是否随磁盘或跨区 RTT 增长。
consumer num_pending、num_ack_pending、redelivered 是否持续增长。
Leaf node:边缘自治与受控汇聚
Leaf node 主动连接 hub,本地客户端可以使用本地低延迟服务,只有存在兴趣的流量跨连接传输。它适合门店、工厂、边缘网络或无法让中心主动入站的环境。官方 Leaf node 说明 指出,leaf 可以跨越安全域,流量受 leaf connection 身份的 publish/subscribe 权限约束。
Leaf 不是“便宜的跨地域集群”。断网期间 Core 消息不会自动排队;需要离线保存时要在边缘部署 JetStream,并明确它是独立 domain、扩展中心 domain 还是通过 source/mirror 汇聚。错误地让边缘与中心共享同一存储或模糊 domain,会让恢复和选主不可预测。
Gateway:连接多个集群的 supercluster
Gateway 把多个 cluster 连接成 supercluster,传播跨集群兴趣,但不会把所有 server 做成一个 route 全互联。它适合多个地域各自保持本地集群,再按订阅兴趣转发流量。官方 gateway 文档 强调 cluster 和 gateway 使用不同协议与端口。
跨地域链路故障时,应事先决定业务是本地继续、快速失败还是由 JetStream 独立保存后异步汇聚。不能把 gateway 连通性当成跨地域数据副本保证。
容量、积压与成本
JetStream 容量不是 volume 大小这么简单。至少要估算:
每日原始写入 = 峰值消息速率 × 平均消息字节 × 每日有效秒数
保留容量 = 每日原始写入 × 保留天数 × 副本数
磁盘预算 = 保留容量 + 索引/块开销 + 重投/增长余量还要分别限制 server、account、stream 和 consumer:
server 的 max_file_store、max_memory_store 防止单实例被写满。stream 的 max_bytes、max_age、max_msgs 和 max_msg_size 定义业务保留边界。discard old 会牺牲最老未消费消息,discard new 会把背压暴露给生产者。
consumer 的 max_ack_pending 和 pull batch 限制并发,避免慢消费者把消息都变成 in-flight。max_deliver 控制反复失败的放大,但需要配套失败转存和人工处置。
容量告警应看趋势而非一个固定百分比。可执行的判断包括:stream bytes 按当前增长率是否会在保留窗口前触顶;consumer pending 的增长斜率是否持续大于处理速率;ack pending 是否长期贴近上限;redelivery ratio 是否突然上升;磁盘同步与发布确认延迟是否共同恶化。
故障证据与排查顺序
发布成功,但 stream 没增长
先执行:
nats stream find acme.order.dev.events.created
nats stream info ORDERS_DEV没有 stream 捕获该 subject 时,普通 Core publish 仍可能成功。使用 JetStream publish 可让生产者直接收到“no response from stream”或存储错误。还要检查 account 是否正确,因为相同 subject 在不同 account 中相互隔离。
consumer 不继续投递
nats consumer info ORDERS_DEV ORDER_WORKER_DEV如果 Num Ack Pending 达到 Max Ack Pending,说明应用拿走消息但没有及时确认;如果 pending 很大而 ack pending 很小,通常是消费者吞吐不足或停止拉取;如果 redelivered 快速增长,检查处理耗时、错误重试和 AckWait。
重启后对象或消息消失
docker inspect nats-dev --format '{{json .Mounts}}'
docker volume inspect nats-dev-data
docker logs nats-dev --tail 200常见原因是没有启用 -js、store_dir 指向临时目录、volume 挂载路径与启动参数不一致,或实际启动了一个新 account/domain。不要先重建 stream,这可能掩盖旧数据仍在另一个目录的事实。
slow consumer 与断连
Core subscriber 处理速度赶不上推送时,客户端缓冲区最终会触发 slow consumer 和断连。JetStream pull consumer 更容易通过批次和 max_ack_pending 做背压。排查要同时看 /connz 的 pending bytes、客户端错误回调、GC/事件循环停顿与下游数据库延迟,不能只调大缓冲区。
清理、回滚与恢复演练
删除 stream 或 volume 之前,先证明备份真的能在另一个环境恢复。为便于核对,先向 ORDERS_DEV 发布一条唯一消息但不要消费,并记录 stream 与 durable consumer 的状态:
nats pub acme.order.dev.events.created \
'{"eventId":"evt-backup-001","orderId":"ord-backup-1"}' \
--jetstream
nats stream info ORDERS_DEV
nats consumer info ORDERS_DEV ORDER_WORKER_DEV
mkdir -p ./backups
nats stream backup ORDERS_DEV ./backups/ORDERS_DEV-snapshot备份命令完成时应报告 stream、消息和 consumer 的写出结果。备份目录不是一句“命令成功”就结束:把目录归档后计算校验值,记录 NATS Server/CLI 版本、account、stream 最后序号、消息数、consumer acknowledgement floor 和生成时间,并把归档放入受访问控制的备份存储。stream payload 可能含个人信息、订单数据和内部 subject,备份的权限、加密、保留期与删除审计不能低于生产消息本身。
恢复必须指向隔离的 NATS 环境或隔离 account,并确保目标中不存在同名 stream。不要为了测试恢复而删除生产 stream:
export RESTORE_NATS_URL='nats://127.0.0.1:5222'
nats --server "$RESTORE_NATS_URL" stream restore \
./backups/ORDERS_DEV-snapshot
nats --server "$RESTORE_NATS_URL" stream info ORDERS_DEV
nats --server "$RESTORE_NATS_URL" \
consumer info ORDERS_DEV ORDER_WORKER_DEV
nats --server "$RESTORE_NATS_URL" \
consumer next ORDERS_DEV ORDER_WORKER_DEV --count 1 --no-ack恢复验收要逐项比较 subjects、retention、storage、replicas、消息数、首尾 sequence,以及 durable consumer 的 filter、ack policy、delivered、ack floor 和 pending。最后一条命令应能读到 evt-backup-001;确认内容和元数据后再显式 ack,并检查 acknowledgement floor 前移。若恢复环境只有一个节点,而备份中的 stream 要求三个副本,应先准备三节点恢复环境,或通过受控的配置变更降低副本数并记录它改变了故障容忍能力,不能把副本不足当成数据损坏。
只有 file storage stream 支持这种快照;memory stream 不能依靠 stream backup 承担灾难恢复。stream backup 也不替代 server 配置、operator/account JWT、签名密钥、TLS 私钥、外部失败记录和应用 inbox 数据库的备份。恢复演练必须覆盖这些依赖,否则消息虽然回来,客户端仍可能因为身份、权限或幂等状态缺失而无法安全继续处理。
开发实验的对象按 consumer、stream、服务顺序清理:
nats consumer rm ORDERS_DEV ORDER_WORKER_DEV --force
nats stream rm ORDERS_DEV --force
docker compose down需要彻底删除本地数据时再执行:
docker compose down -vdown -v 不属于普通停止动作,它会删除 JetStream 数据。共享环境删除 consumer 会改变 InterestPolicy 保留,删除 stream 会永久删除消息,因此必须经过 owner 确认并先导出配置、记录 pending 和保留窗口。
升级或变更前至少演练三类回滚:
客户端配置错误时恢复旧 subject/consumer 配置,并确认新旧消费者不会同时产生重复副作用。server 节点滚动失败时回退镜像,确认 stream leader 与 replica 恢复,而不是只看进程 ready。消息模型变更失败时让生产者停止新 schema,消费者兼容读取旧新版本,不能依靠 purge 清掉不兼容消息。
架构选型:什么时候不用 JetStream
Core NATS 适合在线 request/reply、服务通知和可重建状态,优势是延迟低、运维对象少。JetStream 适合离线恢复、工作队列、事件回放和可确认发布,但引入了磁盘、Raft、保留策略、consumer 生命周期和容量治理。
如果需求是超长时间事件归档、重度流式计算生态或跨团队大规模日志平台,Kafka/Pulsar 等日志型平台可能有更成熟的分区与计算生态;如果需求是复杂路由、逐消息 TTL、传统 DLQ 与 AMQP 模型,RabbitMQ 可能更直接。选型不能只比较吞吐数字,要比较消息语义、故障恢复、团队能力、跨地域模型和总维护成本。
团队长期治理
一个可维护的 NATS 平台至少要有以下所有权:
平台团队维护 server、cluster、TLS、operator/account、存储、升级、备份恢复与容量基线。业务团队维护 subject 契约、schema、stream/consumer 声明、幂等键、重试和失败消息处置。安全团队审查账户边界、凭证生命周期、监控端点、跨 account import/export 和跨地域连接。
值班人员能从 publish ack、stream sequence、consumer sequence、ack pending、redelivery 与应用幂等记录恢复一条消息的完整证据链。
上线评审时,应能回答这些问题:消息使用 Core 还是 JetStream,为什么;stream 达到容量限制时丢旧还是拒新;消费者宕机多久仍可恢复;重复处理如何证明无副作用;最大积压多久能追平;凭证如何吊销;一个节点、一个可用区或一条跨地域链路失效时,生产和消费分别发生什么。
当这些答案都能被配置、指标和演练结果支撑,NATS 才不再只是“能发一条消息”的轻量工具,而是可治理的事件通信基础设施。
