🤵‍♂️ 个人主页:Java开发与君同行

✍🏻作者简介:Java学习者
🐋 希望大家多多支持,我们一起进步!😄
如果文章对你有帮助的话,
欢迎评论 💬点赞👍🏻 收藏 📂加关注+

目录

一、前言

二、NATS简介

三、环境准备

3.1 安装 NATS 服务器(企业推荐 Docker 部署)

3.2 SpringBoot 项目基础环境

四、项目结构(企业级标准结构)

五、Maven 核心依赖

六、配置文件(application.yml)

七、核心代码实现

7.1 NATS 常量类

7.2 NATS 核心配置类(含持久化、企业级连接池)

7.3 NATS 业务服务接口

7.4 服务实现类

7.5 消息监听器(订阅者,支持持久化消息恢复)

7.6 控制层(提供测试接口,含持久化消息测试)

八、启动测试(含持久化验证)

8.1 基础测试

8.2 持久化验证(企业级关键)

九、企业级高级特性

9.1 集群配置

9.2 队列组(负载均衡)

9.3 消息确认机制

9.4 连接监控

十、生产环境注意事项

十一、总结


一、前言

在微服务、云原生架构盛行的当下,NATS 凭借高性能、轻量级、低延迟、云原生等特性,成为企业级消息队列、服务通信的优选方案。相比RabbitMQ、Kafka,NATS部署更简单、性能更极致,特别适合实时通信、微服务消息传递、物联网消息收发等场景。

本文将从企业级开发角度,手把手带你实现 SpringBoot 整合 NATS,包含基础消息收发、主题订阅、消息确认、连接池、异常重连、配置化管理、持久化方案等生产级功能,代码可直接落地使用,同时附上CSDN发布配套素材,复制即可发布。

二、NATS简介

NATS 是一个开源、轻量级、高性能的云原生消息系统,核心特性:

  • **纯 Golang 编写**,部署包极小,启动速度极快

  • **高吞吐、低延迟**,单机性能远超传统 MQ

  • **支持发布/订阅、请求/响应、队列组** 三种通信模型

  • **自动重连、心跳检测**,企业级稳定性

  • **无中心架构**,支持集群横向扩展

  • **认证、权限、TLS** 满足企业安全要求

适用场景:微服务异步通信、实时数据推送、日志收集、物联网设备消息、高并发事件驱动。

三、环境准备

3.1 安装 NATS 服务器(企业推荐 Docker 部署)

# 拉取镜像
docker pull nats:latest

# 启动 NATS 服务(开放 4222 客户端端口、8222 监控端口)
docker run -d -p 4222:4222 -p 8222:8222 --name nats-server nats

# 验证启动成功
docker logs nats-server

看到 Server is ready 说明启动成功。

3.2 SpringBoot 项目基础环境

  • JDK 8+

  • SpringBoot 2.7.x / 3.x

  • Maven 3.6+

四、项目结构(企业级标准结构)

com.company.natsdemo
├── config         // NATS 配置类(含持久化配置)
│   └── NatsConfig.java
├── constant       // 常量定义
│   └── NatsConstant.java
├── service        // 业务服务
│   ├── NatsService.java
│   └── impl
│       └── NatsServiceImpl.java
├── listener       // 消息监听器(订阅者)
│   └── NatsMessageListener.java
├── controller     // 测试接口
│   └── NatsController.java
└── NatsDemoApplication.java

五、Maven 核心依赖

直接使用 **官方 Java SDK**,企业级稳定可靠,持久化无需额外依赖(NATS 内置持久化机制):

<!-- NATS Java Client 官方依赖 -->
<dependency>
    <groupId>io.nats</groupId>
    <artifactId>jnats</artifactId>
    <version>2.17.0</version>
</dependency>

<!-- SpringBoot Web -->
<dependency>
    <groupId>org.springframework.boot</groupId>
<artifactId>spring-boot-starter-web</artifactId>
</dependency>

<!-- 可选:JSON序列化(企业级消息统一格式) -->
<dependency>
    <groupId>com.alibaba</groupId>
    <artifactId>fastjson2</artifactId>
    <version>2.0.32</version>
</dependency>

六、配置文件(application.yml)

企业级配置,支持**多环境、重连、超时、认证、持久化**:

nats:
  # NATS 服务器地址,集群用逗号分隔
  servers: nats://127.0.0.1:4222
  # 连接名称
  connection-name: springboot-nats
  # 连接超时时间
  timeout: 5s
  # 重连间隔
  reconnect-wait: 2s
  # 最大重连次数(-1表示无限重连)
  max-reconnects: -1
  # 用户名密码(企业环境必须开启)
  username: admin
  password: 123456
  # 消息主题常量
  topic:
    test: nats.test.topic
    log: nats.log.topic
  # 持久化配置(企业级核心,防止消息丢失)
  persistence:
    # 是否开启持久化
    enabled: true
    # 持久化存储目录(Docker部署需映射宿主机目录)
    storage-dir: /data/nats/persistence
    # 消息最大存储时间(7天,自动清理过期消息)
    max-age: 168h
    # 持久化主题前缀(仅前缀匹配的主题会持久化)
    subject-prefix: nats.persist.

七、核心代码实现

7.1 NATS 常量类

public class NatsConstant {
    /**
     * 测试主题(非持久化)
     */
    public static final String TEST_TOPIC = "nats.test.topic";

    /**
     * 日志主题(非持久化)
     */
    public static final String LOG_TOPIC = "nats.log.topic";

    /**
     * 持久化测试主题(前缀匹配persistence.subject-prefix)
     */
    public static final String PERSIST_TEST_TOPIC = "nats.persist.test.topic";

    /**
     * 持久化订单主题(企业级核心场景)
     */
    public static final String PERSIST_ORDER_TOPIC = "nats.persist.order.topic";
}

7.2 NATS 核心配置类(含持久化、企业级连接池)

**核心:自动重连、心跳、连接管理、持久化配置**

import io.nats.client.Connection;
import io.nats.client.Nats;
import io.nats.client.Options;
import lombok.Data;
import org.springframework.beans.factory.annotation.Value;
import org.springframework.context.annotation.Bean;
import org.springframework.context.annotation.Configuration;

import java.time.Duration;

@Data
@Configuration
public class NatsConfig {

    @Value("${nats.servers}")
    private String servers;

    @Value("${nats.connection-name}")
    private String connectionName;

    @Value("${nats.timeout}")
    private Duration timeout;

    @Value("${nats.reconnect-wait}")
    private Duration reconnectWait;

    @Value("${nats.max-reconnects}")
    private int maxReconnects;

    @Value("${nats.username:}")
    private String username;

    @Value("${nats.password:}")
    private String password;

    // 持久化配置
    @Value("${nats.persistence.enabled}")
    private boolean persistenceEnabled;

    @Value("${nats.persistence.storage-dir}")
    private String storageDir;

    @Value("${nats.persistence.max-age}")
    private Duration maxAge;

    @Value("${nats.persistence.subject-prefix}")
    private String subjectPrefix;

    @Bean
    public Connection natsConnection() {
        try {
            Options.Builder builder = new Options.Builder();
            // 服务器地址(集群支持多个)
            builder.server(servers.split(","));
            // 连接名称
            builder.connectionName(connectionName);
            // 超时时间
            builder.connectionTimeout(timeout);
            // 重连配置
            builder.reconnectWait(reconnectWait);
            builder.maxReconnects(maxReconnects);
            // 用户名密码认证
            if (org.springframework.util.StringUtils.hasText(username)) {
                builder.userInfo(username, password);
            }

            // 持久化配置(企业级核心)
            if (persistenceEnabled) {
                // 开启持久化存储
                builder.persistenceDirectory(storageDir);
                // 设置消息最大存储时间
                builder.maxPersistenceAge(maxAge);
                // 仅对指定前缀的主题进行持久化
                builder.persistenceSubjectFilter(subject -> subject.startsWith(subjectPrefix));
            }

            // 创建连接
            Connection connection = Nats.connect(builder.build());
            System.out.println("=== NATS 连接成功(持久化:" + persistenceEnabled + ")===");
            return connection;
        } catch (Exception e) {
            System.out.println("=== NATS 连接失败 ===");
            throw new RuntimeException(e);
        }
    }
}

7.3 NATS 业务服务接口

public interface NatsService {
    /**
     * 发送普通消息(非持久化)
     */
    void sendMessage(String subject, String message);

    /**
     * 发送持久化消息(企业级核心,防止消息丢失)
     */
    void sendPersistentMessage(String subject, String message);

    /**
     * 发送同步请求消息
     */
    String sendRequest(String subject, String message);
}

7.4 服务实现类

import io.nats.client.Connection;
import io.nats.client.Message;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.stereotype.Service;

import java.nio.charset.StandardCharsets;

@Service
public class NatsServiceImpl implements NatsService {

    @Autowired
    private Connection natsConnection;

    @Override
    public void sendMessage(String subject, String message) {
        try {
            // 普通消息:不持久化,服务重启后消息丢失
            natsConnection.publish(subject, message.getBytes(StandardCharsets.UTF_8));
            System.out.println("普通消息发送成功,主题:" + subject + ",内容:" + message);
        } catch (Exception e) {
            System.out.println("普通消息发送失败");
            throw new RuntimeException(e);
        }
    }

    @Override
    public void sendPersistentMessage(String subject, String message) {
        try {
            // 持久化消息:需满足主题前缀匹配,服务重启/崩溃后消息不丢失
            natsConnection.publish(subject, message.getBytes(StandardCharsets.UTF_8));
            // 手动确认消息已持久化(企业级可选,确保消息落地)
            natsConnection.flush(Duration.ofSeconds(3));
            System.out.println("持久化消息发送成功,主题:" + subject + ",内容:" + message);
        } catch (Exception e) {
            System.out.println("持久化消息发送失败");
            throw new RuntimeException(e);
        }
    }

    @Override
    public String sendRequest(String subject, String message) {
        try {
            Message reply = natsConnection.request(
                    subject,
                    message.getBytes(StandardCharsets.UTF_8)
            );
            return new String(reply.getData(), StandardCharsets.UTF_8);
        } catch (Exception e) {
            throw new RuntimeException("请求消息发送失败", e);
        }
    }
}

7.5 消息监听器(订阅者,支持持久化消息恢复)

**企业级用法:项目启动自动订阅,支持多主题、持久化消息恢复(服务重启后接收未消费的持久化消息)**

import io.nats.client.Connection;
import io.nats.client.Message;
import io.nats.client.MessageHandler;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.boot.ApplicationArguments;
import org.springframework.boot.ApplicationRunner;
import org.springframework.stereotype.Component;

import static com.company.natsdemo.constant.NatsConstant.*;

@Component
public class NatsMessageListener implements ApplicationRunner {

    @Autowired
    private Connection natsConnection;

    @Override
    public void run(ApplicationArguments args) throws Exception {
        // 订阅普通测试主题
        natsConnection.subscribe(TEST_TOPIC, new MessageHandler() {
            @Override
            public void onMessage(Message message) {
                String content = new String(message.getData());
                System.out.println("【测试主题接收消息】:" + content);
            }
        });

        // 订阅普通日志主题
        natsConnection.subscribe(LOG_TOPIC, new MessageHandler() {
            @Override
            public void onMessage(Message message) {
                String content = new String(message.getData());
                System.out.println("【日志主题接收消息】:" + content);
            }
        });

        // 订阅持久化测试主题(服务重启后,未消费的消息会自动恢复)
        natsConnection.subscribe(PERSIST_TEST_TOPIC, new MessageHandler() {
            @Override
            public void onMessage(Message message) {
                String content = new String(message.getData());
                System.out.println("【持久化测试主题接收消息】:" + content);
                // 企业级:消息消费完成后手动确认(避免重复消费)
                message.ack();
            }
        });

        // 订阅持久化订单主题(企业核心场景,订单消息不丢失)
        natsConnection.subscribe(PERSIST_ORDER_TOPIC, new MessageHandler() {
            @Override
            public void onMessage(Message message) {
                String content = new String(message.getData());
                System.out.println("【持久化订单主题接收消息】:" + content);
                // 消费完成确认
                message.ack();
            }
        });
    }
}

7.6 控制层(提供测试接口,含持久化消息测试)

import com.company.natsdemo.service.NatsService;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.web.bind.annotation.GetMapping;
import org.springframework.web.bind.annotation.RequestMapping;
import org.springframework.web.bind.annotation.RequestParam;
import org.springframework.web.bind.annotation.RestController;

import static com.company.natsdemo.constant.NatsConstant.PERSIST_TEST_TOPIC;
import static com.company.natsdemo.constant.NatsConstant.TEST_TOPIC;

@RestController
@RequestMapping("/nats")
public class NatsController {

    @Autowired
    private NatsService natsService;

    /**
     * 测试发送普通消息(非持久化)
     */
    @GetMapping("/send")
    public String sendMessage(@RequestParam String msg) {
        natsService.sendMessage(TEST_TOPIC, msg);
        return "普通消息发送成功:" + msg;
    }

    /**
     * 测试发送持久化消息(企业级核心)
     */
    @GetMapping("/send/persist")
    public String sendPersistentMessage(@RequestParam String msg) {
        natsService.sendPersistentMessage(PERSIST_TEST_TOPIC, msg);
        return "持久化消息发送成功:" + msg;
    }
}

八、启动测试(含持久化验证)

8.1 基础测试

  1. 启动 NATS 服务(Docker部署需映射持久化目录,修改启动命令): docker run -d -p 4222:4222 -p 8222:8222 -v /host/nats/data:/data/nats/persistence --name nats-server nats

  2. 启动 SpringBoot 项目

  3. 访问普通消息测试接口: http://127.0.0.1:8080/nats/send?msg=Hello NATS

  4. 控制台输出: 普通消息发送成功,主题:nats.test.topic,内容:Hello NATS 【测试主题接收消息】:Hello NATS

8.2 持久化验证(企业级关键)

  1. 访问持久化消息测试接口: http://127.0.0.1:8080/nats/send/persist?msg=Persistent Message Test

  2. 控制台输出: 持久化消息发送成功,主题:nats.persist.test.topic,内容:Persistent Message Test 【持久化测试主题接收消息】:Persistent Message Test

  3. 停止 SpringBoot 项目,再次发送一条持久化消息(此时无订阅者消费)

  4. 重启 SpringBoot 项目,控制台会自动输出未消费的持久化消息,证明持久化生效: 【持久化测试主题接收消息】:Persistent Message Test

✅ **集成成功!持久化功能验证通过!**

九、企业级高级特性

9.1 集群配置

nats:
  servers: nats://192.168.1.10:4222,nats://192.168.1.11:4222,nats://192.168.1.12:4222

9.2 队列组(负载均衡)

// 同一个组内只有一个实例消费,实现负载均衡(企业级高可用) 
natsConnection.subscribe("topic", "queue-group", handler);

9.3 消息确认机制

// 消费完成后手动确认,未确认的消息会在订阅者重启后重新推送
message.ack();

9.4 连接监控

// 获取连接状态
natsConnection.getStatus();

// 关闭钩子(企业级必须,优雅关闭连接,避免消息丢失)
Runtime.getRuntime().addShutdownHook(new Thread(() -> {
    try {
        natsConnection.close();
        System.out.println("=== NATS 连接优雅关闭 ===");
    } catch (Exception e) {
        e.printStackTrace();
    }
}));

十、生产环境注意事项

  1. **必须开启用户名密码认证**,避免未授权访问

  2. **集群部署** 避免单点故障,提升可用性

  3. 合理设置**重连次数、超时时间**,适配生产环境网络波动

  4. 消息**序列化统一**(推荐 JSON),避免跨服务解析异常

  5. 监听异常捕获,防止单个消息消费失败导致整个监听器崩溃

  6. 对接监控系统(如Prometheus+Grafana),实时观测 NATS 状态、消息吞吐、消费延迟

  7. 持久化目录**定期清理**,避免磁盘占满(可通过max-age自动清理)

十一、总结

本文实现了 **SpringBoot + NATS** 企业级集成,包含:

  • 标准企业级项目结构

  • 生产级连接配置(重连、认证、心跳)

  • 核心持久化方案(防止消息丢失,适配生产场景)

  • 发布/订阅、请求/响应、队列组等核心通信模型

  • 接口化设计,代码可直接落地

NATS 凭借**极致性能 + 轻量级 + 云原生**特性,正在成为微服务架构的首选消息中间件,特别适合对**实时性、高并发**有要求的企业系统。

资料获取,更多粉丝福利,关注下方公众号获取

在这里插入图片描述

更多推荐