企业级SpringBoot集成NATS:高性能消息队列实战落地

🤵♂️ 个人主页:Java开发与君同行
✍🏻作者简介:Java学习者
🐋 希望大家多多支持,我们一起进步!😄
如果文章对你有帮助的话,
欢迎评论 💬点赞👍🏻 收藏 📂加关注+
目录
3.1 安装 NATS 服务器(企业推荐 Docker 部署)
一、前言
在微服务、云原生架构盛行的当下,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 基础测试
启动 NATS 服务(Docker部署需映射持久化目录,修改启动命令):
docker run -d -p 4222:4222 -p 8222:8222 -v /host/nats/data:/data/nats/persistence --name nats-server nats启动 SpringBoot 项目
访问普通消息测试接口:
http://127.0.0.1:8080/nats/send?msg=Hello NATS控制台输出:
普通消息发送成功,主题:nats.test.topic,内容:Hello NATS【测试主题接收消息】:Hello NATS
8.2 持久化验证(企业级关键)
-
访问持久化消息测试接口:
http://127.0.0.1:8080/nats/send/persist?msg=Persistent Message Test -
控制台输出:
持久化消息发送成功,主题:nats.persist.test.topic,内容:Persistent Message Test【持久化测试主题接收消息】:Persistent Message Test -
停止 SpringBoot 项目,再次发送一条持久化消息(此时无订阅者消费)
-
重启 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();
}
}));
十、生产环境注意事项
**必须开启用户名密码认证**,避免未授权访问
**集群部署** 避免单点故障,提升可用性
合理设置**重连次数、超时时间**,适配生产环境网络波动
消息**序列化统一**(推荐 JSON),避免跨服务解析异常
监听异常捕获,防止单个消息消费失败导致整个监听器崩溃
对接监控系统(如Prometheus+Grafana),实时观测 NATS 状态、消息吞吐、消费延迟
持久化目录**定期清理**,避免磁盘占满(可通过max-age自动清理)
十一、总结
本文实现了 **SpringBoot + NATS** 企业级集成,包含:
标准企业级项目结构
生产级连接配置(重连、认证、心跳)
核心持久化方案(防止消息丢失,适配生产场景)
发布/订阅、请求/响应、队列组等核心通信模型
接口化设计,代码可直接落地
NATS 凭借**极致性能 + 轻量级 + 云原生**特性,正在成为微服务架构的首选消息中间件,特别适合对**实时性、高并发**有要求的企业系统。
资料获取,更多粉丝福利,关注下方公众号获取

更多推荐
所有评论(0)