NATS JetStream 持久化消息

FreeGuideOnline 最新 2026-07-10

bash nats stream add ORDERS
--subjects "orders.>"
--storage file
--retention limits
--max-msgs 1000000
--max-age 7d
--max-bytes 10GB
--discard old


解释常用参数:

- `--subjects`:该流负责捕获的主题模式。`"orders.>"` 表示匹配所有以 `orders.` 开头的主题。
- `--storage file`:使用文件持久化,断电不丢失。
- `--retention limits`:根据消息数量、大小和时间三个限制进行保留。
- `--max-msgs 1000000`:最多保留 100 万条消息。
- `--max-age 7d`:消息最长保留 7 天。
- `--max-bytes 10GB`:流的总大小不超过 10GB。
- `--discard old`:当达到任何限制时,丢弃最旧的消息(也可以设置为 `--discard new` 拒绝新消息)。

执行后会确认并创建流,你可以通过 `nats stream info ORDERS` 查看详细信息。

### 使用内存存储(高性能但非持久)

```bash
nats stream add CACHE \
  --subjects "cache.>" \
  --storage memory \
  --max-msgs 10000

内存存储的读写速度更快,但服务重启会丢失所有消息。

发布和消费持久化消息

发布消息到流

只需向流绑定的任一主题发布消息,消息就会自动持久化。

# 发布一条普通消息
nats pub orders.new "order_id: 1001, item: book"

# 发布带 JetStream 确认的消息(推荐)
nats pub orders.new "order_id: 1002" --js

使用 --js 标志会使用 JetStream 发布的确认机制,确保消息已被流成功存储后才返回成功。

创建 Push 消费者并实时消费

Push 消费者会将消息推送到一个可订阅的交付主题。创建时指定交付目标:

nats consumer add ORDERS MONITOR \
  --deliver all \
  --replay instant \
  --ack explicit \
  --deliver-group monitor
  • --deliver all:从流中最早的消息开始递送(可选 lastnewby_start_sequence)。
  • --replay instant:尽快发送所有历史消息,不模拟原始时间间隔。
  • --ack explicit:需要显式确认每条消息。
  • --deliver-group monitor:指定一个交付主题名称(会变成 _INBOX.> 之类的具体主题)。

创建后,你可以订阅由消费者生成的交付主题来接收消息。通常你不需要手动管理交付主题,在客户端 SDK 中,消费者创建时会自动生成推送订阅。

使用 CLI 快速测试:

# 在终端1中,创建一个临时 push 消费者,直接输出消息
nats consumer sub ORDERS MONITOR --raw

然后在另一个终端发布消息,就可以在消费者窗口看到输出了。

创建 Pull 消费者并手动拉取

Pull 消费者不会自动推送,客户端需要主动发出拉取请求。

nats consumer add ORDERS BATCH \
  --deliver all \
  --ack explicit \
  --pull

拉取消息:

nats consumer next ORDERS BATCH

这会请求并返回一条消息。你也可以批量拉取:

nats consumer next ORDERS BATCH --count 10

消息确认与重投策略

JetStream 提供灵活的确认机制,保证消息至少被成功处理一次。

确认类型

  • Explicit(显式确认):消费者必须对每条消息调用 ack,否则消息会被视为未处理。
  • All(确认所有):确认一条消息意味着同时确认所有序列号小于或等于该消息之前的未确认消息。
  • None(无确认):消息投递后立即被认为已确认,失败不重试。

未确认消息的处理

当消息未被及时确认时,JetStream 会根据消费者配置的 ack_wait 时间和 max_deliver 次数进行重投。超过最大投递次数后,消息可以被转发到一个专门处理死信的流(Dead Letter Queue)。

配置示例:

nats consumer add ORDERS RELIABLE \
  --deliver all \
  --ack explicit \
  --ack-wait 30s \
  --max-deliver 5 \
  --sample 100
  • --ack-wait 30s:30 秒内需要确认,否则重投。
  • --max-deliver 5:每条消息最多投递 5 次,之后标记为有问题的消息(或转移到死信)。
  • --sample 100:按 100% 比例收集消费者层面的遥测数据(监控可见)。

编程语言示例(Go)

以下是用 Go 语言连接 JetStream、发布和 Pull 消费的简化示例,展示持久化交互。

package main

import (
    "fmt"
    "log"
    "time"

    "github.com/nats-io/nats.go"
)

func main() {
    // 连接 NATS
    nc, err := nats.Connect("nats://localhost:4222")
    if err != nil {
        log.Fatal(err)
    }
    defer nc.Close()

    // 获取 JetStream 上下文
    js, err := nc.JetStream()
    if err != nil {
        log.Fatal(err)
    }

    // 检查或创建流(通常事先通过配置文件或管理API创建)
    streamName := "ORDERS"
    _, err = js.AddStream(&nats.StreamConfig{
        Name:     streamName,
        Subjects: []string{"orders.>"},
        Storage:  nats.FileStorage,
        MaxAge:   7 * 24 * time.Hour,
        MaxMsgs:  1_000_000,
        MaxBytes: 10 * 1024 * 1024 * 1024, // 10GB
    })
    if err != nil {
        log.Printf("Stream may already exist: %v", err)
    }

    // 发布持久化消息
    ack, err := js.Publish("orders.new", []byte("order_id: 1001"))
    if err != nil {
        log.Fatal(err)
    }
    fmt.Printf("Published msg with sequence: %d\n", ack.Sequence)

    // 创建 Pull 消费者
    consumerName := "BATCH"
    _, err = js.AddConsumer(streamName, &nats.ConsumerConfig{
        Durable:   consumerName,
        AckPolicy: nats.AckExplicitPolicy,
    })
    if err != nil {
        log.Fatal(err)
    }

    // 拉取消息
    sub, err := js.PullSubscribe("orders.>", consumerName)
    if err != nil {
        log.Fatal(err)
    }

    // 一次拉取 5 条消息,等待最多 2 秒
    msgs, err := sub.Fetch(5, nats.MaxWait(2*time.Second))
    if err != nil {
        log.Fatal(err)
    }

    for _, msg := range msgs {
        fmt.Printf("Received: %s\n", string(msg.Data))
        msg.Ack()
    }
}