# go-nats **Repository Path**: golang_common/go-nats ## Basic Information - **Project Name**: go-nats - **Description**: nats队列中间件 - **Primary Language**: Unknown - **License**: Not specified - **Default Branch**: master - **Homepage**: None - **GVP Project**: No ## Statistics - **Stars**: 0 - **Forks**: 0 - **Created**: 2026-09-09 - **Last Updated**: 2026-09-11 ## Categories & Tags **Categories**: Uncategorized **Tags**: None ## README # go-nats 一个基于官方 `github.com/nats-io/nats.go v1.53.1` 封装的中文生产级 NATS 客户端库。 ```bash go get gitee.com/golang_common/go-nats@latest ``` ```go import "gitee.com/golang_common/go-nats" ``` ## 一、先看这张表,选择你需要的用法 | 你的业务场景 | 应该选哪种用法 | 用法示例文件 | 消息会落盘? | 消费者离线能补收? | 可靠性 | | --- | --- | --- | --- | --- | --- | | 微服务同步调用:A 调 B,等 B 返回结果 | 请求回复 Request/Reply | `examples/core_request_reply/main.go` | 不会 | 不会 | 最多一次 | | 实时广播通知:多个服务都要立刻收到 | 发布订阅 Publish/Subscribe | `examples/core_pubsub/main.go` | 不会 | 不会 | 最多一次 | | 多实例分摊任务:每个任务只处理一次 | 队列组 QueueSubscribe | `examples/core_queue/main.go` | 不会 | 不会 | 最多一次 | | 订单/支付等关键异步事件,不能丢、要补收 | JetStream 常驻消费 Consume | `examples/jetstream_persistent/main.go` | 会 | 会 | 至少一次 | | 定时批处理:每批拉若干条,逐条处理成功后逐条 Ack | JetStream 手动 Fetch | `examples/jetstream_fetch/main.go` | 会 | 会 | 至少一次 | | 订单超时、预约任务、指定未来时间触发 | JetStream 延时队列 | `examples/jetstream_delay_dlq/main.go` | 会 | 会 | 至少一次 | | 重试耗尽后需要保留、告警和人工重放 | JetStream 死信队列 DLQ | `examples/jetstream_delay_dlq/main.go` | 会 | 会 | 至少一次 | 一句话选型: - 要“等结果”的 RPC,用 **请求回复**。 - 要“发出去就完事”的在线广播,用 **发布订阅**。 - 要“多个 Worker 分摊处理”,用 **队列组**。 - 要“离线也不丢、还能重试”,用 **JetStream**。 - 要“未来某个时间再触发”,用 **JetStream 延时队列**。 - 要“失败重试耗尽后保留原消息”,用 **JetStream 死信队列**。 --- ## 二、安装依赖 ```bash go get gitee.com/golang_common/go-nats@latest ``` 本地用 docker-compose 启动单节点 NATS + JetStream 持久化: ```bash cp .env.example deploy/.env docker compose -f deploy/docker-compose.yml up -d --wait ``` 配置自己的连接信息: ```bash export NATS_URL=nats://127.0.0.1:4222 export NATS_USER=nats export NATS_PASSWORD=你的密码 ``` ## 三、所有用法共用的连接代码 下面这段是每种用法都会用到的连接入口,复制一次即可。 ```go package main import ( "context" "time" "gitee.com/golang_common/go-nats" ) // newClient 从环境变量读取 NATS_URL/NATS_USER/NATS_PASSWORD 建立连接。 func newClient() (*natslib.Client, error) { cfg := natslib.ConfigFromEnv() if cfg.URL == "" { cfg.URL = "nats://127.0.0.1:4222" } ctx, cancel := context.WithTimeout(context.Background(), 10*time.Second) defer cancel() client, err := natslib.Connect(ctx, cfg) if err != nil { return nil, err } return client, nil } ``` --- ## 四、用法 1:发布订阅 Publish/Subscribe > 示例文件:`examples/core_pubsub/main.go` > 下面的每种用法默认你已经把第三节的 `newClient` 放进同一个包; > 如果不想拼代码,也可以直接运行本仓库 `examples/*/main.go` 里的完整程序。 ### 适合什么场景 - 实时推送、广播通知。 - 一个服务发布事件,多个在线服务各自收到一份。 - 例如:价格变更通知、登录状态变更、公共数据刷新。 ### 有什么弊端 - Core NATS 不持久化,不返回服务端 ack。 - 发布时订阅者不在线,这条消息直接丢弃。 - 订阅者后来再启动,收不到离线期间的消息。 - 服务重启不会补发,可靠性是 at-most-once。 ### 具体意思 - 比如 A 往 "order" 里面发布了一条消息,那么 B,C,D都各自订阅了 "order",那么他们都会收到这条消息。 ### 直接复制使用 ```go package main import ( "context" "fmt" "log" "time" "gitee.com/golang_common/go-nats" ) func main() { ctx, cancel := context.WithTimeout(context.Background(), 15*time.Second) defer cancel() client, err := newClient() if err != nil { log.Fatal(err) } defer client.Close(context.Background()) // 1. 先订阅 sub, err := client.Subscribe("weather.updates", func(msg *natslib.NatsMsg) { fmt.Printf("收到天气: %s\n", string(msg.Data)) }) if err != nil { log.Fatal(err) } defer client.Unsubscribe(sub, true) // 2. 等订阅注册到服务器 _ = client.Ping(ctx) // 3. 再发布 _ = client.PublishJSON("weather.updates", map[string]any{ "city": "上海", "temperature": 26.5, }) } ``` ### 什么时候不要用 当发布方和订阅方可能不在线、消息还要补收时,不要用这种用法,改用下面的 JetStream。 --- ## 五、用法 2:队列组 QueueSubscribe > 示例文件:`examples/core_queue/main.go` ### 适合什么场景 - 同一类任务在多个实例之间分摊,每条消息只交给一个 Worker。 - 例如:后台图片处理、批量通知发送、爬虫任务队列。 ### 有什么弊端 - 队列组成员都必须在线。 - Worker 离线期间发布的消息不会保留。 - Worker 处理失败没有内置重试机制,进程退出等于消息丢失。 ### 具体意思 - 比如 A 往 "order" 里面发布了一条消息,那么 B,C,D都各自订阅了 "order",那么 ABC 就只会随机一个人处理,不会 ABC 三个都收到这条消息。 ### 直接复制使用 ```go package main import ( "fmt" "time" "gitee.com/golang_common/go-nats" ) func main() { client, err := newClient() if err != nil { panic(err) } defer client.Close(context.Background()) handler := func(msg *natslib.NatsMsg) { fmt.Printf("Worker 处理: %s\n", string(msg.Data)) } // 三个实例都用同一个 queue 名 "workers" _, _ = client.QueueSubscribe("jobs.execute", "workers", handler) _, _ = client.QueueSubscribe("jobs.execute", "workers", handler) _, _ = client.QueueSubscribe("jobs.execute", "workers", handler) for i := 1; i <= 10; i++ { _ = client.Publish("jobs.execute", []byte(fmt.Sprintf("job-%d", i))) } time.Sleep(2 * time.Second) } ``` ### 什么时候不要用 当任务失败必须重试、Worker 重启期间不能丢任务时,用下面的 JetStream 持久化消费代替。 --- ## 六、用法 3:请求回复 Request/Reply > 示例文件:`examples/core_request_reply/main.go` ### 适合什么场景 - 微服务同步 RPC:A 调用 B,等 B 返回结果。 - 例如:用户查询、余额校验、下单接口、翻译服务。 ### 有什么弊端 - 请求有超时,服务端不在线会返回无响应错误。 - 消息不落盘,服务端重启期间请求不会排队。 - 如果服务端有多个实例,不能全部用普通 `Reply` 订阅同一个 subject,否则每个实例都会重复处理。 - 多实例时服务端要使用同一个 queue 名订阅,让 NATS 只挑一个实例处理。 ### 直接复制使用 服务端: ```go package main import ( "encoding/json" "gitee.com/golang_common/go-nats" ) type addRequest struct { A int `json:"a"` B int `json:"b"` } func main() { client, err := newClient() if err != nil { panic(err) } defer client.Close(context.Background()) // 单实例服务端直接 ReplyJSON。 _, _ = client.ReplyJSON("calc.add", func(msg *natslib.NatsMsg) (any, error) { var req addRequest if err := json.Unmarshal(msg.Data, &req); err != nil { return nil, err } return map[string]int{"sum": req.A + req.B}, nil }) select {} } ``` 客户端: ```go package main import ( "context" "fmt" "log" "time" "gitee.com/golang_common/go-nats" ) func main() { client, err := newClient() if err != nil { log.Fatal(err) } defer client.Close(context.Background()) ctx, cancel := context.WithTimeout(context.Background(), 3*time.Second) defer cancel() var resp struct { Sum int `json:"sum"` } err = client.RequestJSON(ctx, "calc.add", map[string]int{"a": 5, "b": 3}, &resp) if err != nil { log.Fatal(err) } fmt.Println("结果:", resp.Sum) } ``` ### 服务端多实例部署时 服务端不要直接用 `ReplyJSON`,改成队列组订阅 + `msg.Respond`: ```go _, _ = client.QueueSubscribe("calc.add", "calc-service-group", func(msg *natslib.NatsMsg) { _ = msg.Respond(calcResult) }) ``` 这样每个请求只会被队列组中的一个实例处理,客户端依然用 `RequestJSON` 调用。 --- ## 七、用法 4:JetStream 常驻持久化消费(推荐给不能丢消息的业务) > 示例文件:`examples/jetstream_persistent/main.go` ### 适合什么场景 - 订单、支付、积分、通知等关键异步事件。 - 发布时消费者不在线,消息需要先落盘,以后补消费。 - 消费失败需要自动重试。 ### 有什么弊端 - 使用复杂度高于 Core。 - 需要先创建 Stream 和 Durable Consumer。 - 业务处理成功后才 Ack,因此是至少一次投递,业务代码必须幂等。 - 同样的消息可能因为网络重试或进程崩溃重复投递。 ### 直接复制使用 ```go package main import ( "context" "fmt" "log" "time" "gitee.com/golang_common/go-nats" ) func main() { client, err := newClient() if err != nil { log.Fatal(err) } defer client.Close(context.Background()) ctx := context.Background() js := client.JetStream() // 1. 幂等创建持久化 Stream streamCfg := natslib.NewLimitsStreamConfig("ORDERS", []string{"orders.>"}) streamCfg.Replicas = 1 _, err = js.EnsureStream(ctx, streamCfg) if err != nil { log.Fatal(err) } dlqCfg := natslib.NewDeadLetterStreamConfig("ORDER_DLQ", []string{"order_dlq.>"}) dlqCfg.Replicas = 1 _, err = js.EnsureStream(ctx, dlqCfg) if err != nil { log.Fatal(err) } // 2. 同步发布并等待服务端 ack,msgID 用于重复窗口内去重 orderID := "order-1001" ack, err := js.PublishWithID(ctx, "orders.created", []byte(`{"order_id":"order-1001"}`), orderID) if err != nil { log.Fatal(err) } fmt.Println("已确认存储:", ack.Stream, ack.Sequence) // 3. 常驻消费:handler 返回 nil 才 Ack,返回 error 自动 Nak 重投。 // 默认重试 3 次;配置 WithDLQ 后,重试耗尽会进入死信队列。 // 真实应用中请放在 goroutine 里运行。 err = js.Consume( ctx, "ORDERS", "order-processor", func(_ context.Context, msg natslib.JetStreamMsg) error { fmt.Printf("处理订单: %s\n", string(msg.Data())) return doBusiness(msg.Data()) }, natslib.WithFilterSubject("orders.created"), natslib.WithMaxAckPending(100), natslib.WithMaxRetries(3), natslib.WithRetryDelay(2*time.Second), natslib.WithDLQ("order_dlq.orders.created"), ) if err != nil { log.Fatal(err) } } func doBusiness(data []byte) error { // 在这里执行你的真实业务,例如写入数据库。 // 只有这里成功返回 nil,消息才会被 Ack。 return nil } ``` ### 什么时候不要用 - 消息丢了也无所谓,只想快速广播。 - 处理逻辑没有做幂等,重复投递会造成重复下单/重复扣款。 --- ## 八、用法 5:JetStream 手动 Fetch > 示例文件:`examples/jetstream_fetch/main.go` ### 适合什么场景 - 定时任务、离线补单、日终批处理。 - 每次最多拉一批消息,然后对每一条逐条处理、逐条 Ack。 ### 有什么弊端 - 不是实时推送,取决于任务调度频率。 - 消息拉下来后不会自动 Ack,开发者必须自己控制成功/失败。 - 如果忘记 Ack,消息会按 AckWait 超时后重新投递,可能造成重复处理。 ### 直接复制使用 注意:`Fetch` 只是“一次最多向服务端要 batch 条”,返回后仍然是逐条处理、逐条 Ack。 ```go package main import ( "context" "log" "time" "gitee.com/golang_common/go-nats" ) func main() { client, err := newClient() if err != nil { log.Fatal(err) } defer client.Close(context.Background()) ctx := context.Background() js := client.JetStream() // 1. 创建持久化 Stream 和 Durable Consumer _, _ = js.EnsureStream(ctx, natslib.NewLimitsStreamConfig("NOTIFICATIONS", []string{"notify.email"})) _, _ = js.EnsureConsumer(ctx, "NOTIFICATIONS", natslib.NewConsumerConfig("email-worker", "notify.email")) // 2. 一次最多拉 10 条,最多等 2 秒;可能少于 10 条,也可能一条都没有。 msgs, err := js.Fetch(ctx, "NOTIFICATIONS", "email-worker", 10, 2*time.Second) if err != nil { log.Fatal(err) } // 3. 返回后逐条处理,每处理成功一条就手动 Ack 一条,不存在“整批一起 Ack”。 for _, msg := range msgs { err := sendEmail(msg.Data()) if err != nil { _ = msg.NakWithDelay(5 * time.Second) // 失败:稍后重投 continue } _ = msg.Ack() // 成功:告诉 JetStream 已处理 } } func sendEmail(data []byte) error { return nil } ``` ### 什么时候不要用 - 消息需要一发布就立刻被消费。 - 不想手动写 Ack/Nak 逻辑。 --- ## 九、延时队列与死信队列 > 示例文件:`examples/jetstream_delay_dlq/main.go` 延时队列使用 JetStream 原生 message schedules。定时配置 subject 和最终目标 subject 必须由同一个 Stream 捕获,因此推荐直接使用 `NewDelayStreamConfig`。NATS 会把定时 配置 subject 当作计划标识,所以本库会在传入的前缀后追加唯一 token,避免多条延时 消息互相覆盖;`PublishDelayedWithID` 则使用业务 `msgID` 的稳定 hash: ```go delayCfg := natslib.NewDelayStreamConfig( "ORDER_TIMEOUTS", []string{"order_timeouts.schedule.>", "order_timeouts.due"}, ) _, _ = js.EnsureStream(ctx, delayCfg) _, err := js.PublishDelayedWithID( ctx, "order_timeouts.schedule", "order_timeouts.due", time.Now().Add(30*time.Minute), []byte(`{"order_id":"order-1001"}`), "order-timeout:order-1001", ) ``` 消费端默认在业务失败后重试 3 次,可通过 `WithMaxRetries` 修改。配置 `WithDLQ` 后,第 4 次投递仍失败时会写入死信 subject,然后终止原消息: ```go err = js.Consume( ctx, "ORDER_TIMEOUTS", "order-timeout-worker", handleOrderTimeout, natslib.WithFilterSubject("order_timeouts.due"), natslib.WithMaxRetries(3), natslib.WithRetryDelay(5*time.Second), natslib.WithDLQ("order_dlq.timeouts"), ) ``` 死信 Stream 可以用 `NewDeadLetterStreamConfig` 创建,默认保留 30 天: ```go dlqCfg := natslib.NewDeadLetterStreamConfig("ORDER_DLQ", []string{"order_dlq.>"}) _, _ = js.EnsureStream(ctx, dlqCfg) ``` 死信消息类型为 `natslib.DeadLetterMessage`,包含原始 payload、headers、Stream、 Consumer、消息序号、投递次数和最后一次处理错误。重试次数中的 `3` 表示首次处理之外 再重试 3 次,也就是最多执行 4 次。 ## 十、消息可靠性总表 | 用法 | 是否落盘 | 服务端 ack | 离线补收 | 失败重试 | 语义 | | --- | --- | --- | --- | --- | --- | | Publish/Subscribe | 否 | 无 | 否 | 无 | 最多一次 | | QueueSubscribe | 否 | 无 | 否 | 无 | 最多一次 | | Request/Reply | 否 | 请求超时控制 | 否 | 无 | 最多一次 | | JetStream Consume | 是 | 有 | 是 | 默认重试 3 次,可配死信队列 | 至少一次 | | JetStream Fetch | 是 | 有 | 是 | 手动 Nak | 至少一次 | JetStream 至少一次语义意味着消息可能重复投递,业务处理要幂等: - 使用稳定业务 ID。 - 用 `PublishWithID` 发布。 - 消费处理前先判断是否已经处理过。 ## 十一、服务端落盘强度 当前 [deploy/nats/nats-server.conf](/Users/admin/Desktop/code/middleware/go-nats/deploy/nats/nats-server.conf) 使用: ```conf sync_interval: "2s" ``` 含义: - 每 2 秒后台批量 fsync 一次。 - 正常服务、优雅重启不丢消息。 - 强杀进程或断电时,最后 2 秒内已 ack 的消息可能丢失。 如果要求“服务端 ack = 已经 fsync 落盘”,改成: ```conf sync_interval: always ``` 写入吞吐会下降到大约 2000 条/秒,但突然断电也不丢已 ack 消息。 ## 十二、性能参考 仓库内置压测工具: ```bash go run ./bench/throughput -n 5000 -payload 256 -workers 8 ``` 本机 Docker Desktop 实测: | 测试 | `always` | `"2s"` | | --- | --- | --- | | 256B 同步写入 | 约 2296 条/秒 | 约 38331 条/秒 | | 1KB 同步写入 | 约 1767 条/秒 | 约 37647 条/秒 | | 持久化消费并 Ack | 十万级/秒 | 十万级/秒 | ## 十三、本仓库可运行示例对照 | 用法 | 可运行文件 | | --- | --- | | 发布订阅 | [examples/core_pubsub/main.go](/Users/admin/Desktop/code/middleware/go-nats/examples/core_pubsub/main.go) | | 队列组 | [examples/core_queue/main.go](/Users/admin/Desktop/code/middleware/go-nats/examples/core_queue/main.go) | | 请求回复 | [examples/core_request_reply/main.go](/Users/admin/Desktop/code/middleware/go-nats/examples/core_request_reply/main.go) | | JetStream 常驻持久化消费 | [examples/jetstream_persistent/main.go](/Users/admin/Desktop/code/middleware/go-nats/examples/jetstream_persistent/main.go) | | JetStream 手动拉取 | [examples/jetstream_fetch/main.go](/Users/admin/Desktop/code/middleware/go-nats/examples/jetstream_fetch/main.go) | | JetStream 延时队列 + 死信队列 | [examples/jetstream_delay_dlq/main.go](/Users/admin/Desktop/code/middleware/go-nats/examples/jetstream_delay_dlq/main.go) | | 吞吐压测 | [bench/throughput/main.go](/Users/admin/Desktop/code/middleware/go-nats/bench/throughput/main.go) |