- NATS-研究学习
- @[toc]
- 介绍说明
- 提供的服务内容
- 各模式介绍测试使用
- 发布订阅(Publish Subscribe)
- 请求响应(Request Reply)
- 队列订阅&分享工作(Queue Subscribers & Sharing Work)
- 小杭写的Demo
- 简单安装使用与测试
- JetStream 简单使用Demo
- Spring 项目整合
- Nkey 认证连接
- 参考资料
NATS-研究学习- 文章目录
介绍说明
NATS是一个go语言开发的开源的、轻量、高性能的原生消息系统。消息由主题处理,不依赖于网络位置。它提供了应用程序或服务与底层物理网络之间的抽象层。数据被编码并作为消息,由发布者发送。消息由一个或多个订阅者接收、解码和处理。
NATS使程序可以很容易地跨不同的环境、语言、云提供商和内部系统进行通信。客户机通常通过单个URL连接到NATS系统,然后向主题订阅或发布消息。通过这种简单的设计,NATS允许程序共享通用的消息处理代码,隔离资源和相互依赖。
NATS核心提供最多一次的服务质量。
默认情况下,NATS是一种即发即弃的消息传递系统。
如果订户没有收听主题(没有主题匹配),或者在发送消息时未激活,则不会收到消息。
如果需要高级的东东,可以试用NATS Streaming 进行,属于NATS的一个服务模块了。
**优点:**使用简单,配置简单。速度极快,性能良好。
多语言支持,不依赖于网络位置,client端只需知道nats的节点和约定好的subject名称即可。
**缺点:**对服务器稳定性要求较高,机房出现故障,导致nats server端需要重连。可能需要重启nats-server。
在消息timeout后,需要在reconnection里要重新初始化连接,不方便。
提供的服务内容
NATS支持各种消息传递模型,包括:
发布订阅(Publish Subscribe)
请求回复(Request Reply)
队列订阅(Queue Subscribers )
提供的功能:
纯粹的发布订阅模型(Pure pub-sub)
服务器集群(Cluster mode server)
自动精简订阅者(Auto-pruning of subscribers)
基于文本协议(Text-based protocol)
多服务质量保证(Multiple qualities of service - QoS)
各模式介绍测试使用
<dependency> <groupId>io.nats</groupId> <artifactId>jnats</artifactId> <version>2.16.13</version> </dependency>
发布订阅(Publish Subscribe)
NATS将publish/subscribe消息分发模型实现为一对多通信,发布者在 subject 上发送消息,并且监听该Subject在任何活动的订阅者都会收到该消息。
Demo:【测试可用】
//publish
Connection nc = Nats.connect("nats://127.0.0.1:4222");
nc.publish("subject", "hello world".getBytes(StandardCharsets.UTF_8));
//subscribe [这个时间内就之后收到一个,就结束了]
Subscription sub = nc.subscribe("subject");
Message msg = sub.nextMessage(Duration.ofMillis(500));
String response = new String(msg.getData(), StandardCharsets.UTF_8);
//或者是基于回调的subscribe [这个程序可以保持,持续接收信息]
//subscribe
Dispatcher d = nc.createDispatcher(msg ->{
String response = new String(msg.getData(), StandardCharsets.UTF_8);
//do something
})
d.subscribe("subject");
请求响应(Request Reply)
Request-Reply是现代分布式系统中的常见模式。发布者(crm)发送一个请求,应用程序(ybind,fpga-agent)要么在响应时等待一定的超时,要么异步接收响应。Request()是一个简单方便的API,它提供了一个伪同步的方式,使用了超时timeout设置。它创建了一个收件箱(收件箱是一种subject类型,对请求者唯一),订阅subject,然后发布你的请求消息(消息带reply地址)设置为收件箱的subject,然后等待响应,或者超时取消。
Demo:【测试可用】
// publish
Connection nc = Nats.connect("nats://127.0.0.1:4222");
String reply = "replyMsg"; // 这个相当于回到的主题
//请求回应方法回调
Dispatcher d = nc.createDispatcher(msg -> {
System.out.println("reply: " + JSON .toJSONString(msg));
}) ;
d.unsubscribe(reply , 1);
//订阅请求
d.subscribe(reply);
//发布请求
nc.publish("requestSub", reply, "request".getBytes(StandardCharsets.UTF_8));
//subscribe
Connection nc = Nats.connect("nats://127.0.0.1:4222");
//注册订阅
Dispatcher dispatcher = nc.createDispatcher(msg -> {
System.out.println(JSON.toJSONString(msg));
nc.publish(msg.getReplyTo(), "this is reply".getBytes(StandardCharsets.UTF_8));
});
dispatcher.subscribe("requestSub");
队列订阅&分享工作(Queue Subscribers & Sharing Work)
NATS提供称为队列订阅的负载均衡功能。
主要功能是将具有相同queue名字的subject进行负载均衡。
要创建一个消息队列,订阅者需注册一个队列名。所有的订阅者用同一个队列名,形成一个队列组。当消息发送到主题后,队列组会自动选择一个成员接收消息。尽管队列组有多个订阅者,但每条消息只能被组中的一个订阅者接收。
Demo:【测试可用】
// Subscribe
Connection nc = Nats.connect();
Dispatcher d = nc.createDispatcher(msg -> {
//do something
System.out.println("msg: " + new String(msg.getData(),StandardCharsets.UTF_8));
});
d.subscribe("subject", "queName"); //差别就是这个了
小杭写的Demo
/**
* 发布Demo
*/
public class NatsPublish {
public static void main(String[] args) throws IOException, InterruptedException {
// publishSubscribe();
requestReply();
}
/**
* test 请求响应(Request Reply) 模式
*/
public static void requestReply() throws IOException, InterruptedException {
Connection nc = Nats.connect("nats://192.168.137.xxx:4222");
// 这个相当于回到的主题
String reply = "replyMsg-qingqiuxinxi";
//请求回应方法回调
Dispatcher d = nc.createDispatcher(msg ->{
System.out.println("=========收到返回的信息============");
System.out.println("reply:get retuen: " + JSON.toJSONString(msg));
System.out.println( JSON.parseObject(JSON.toJSONString(msg)).get("data") );
String data = (String) JSON.parseObject(JSON.toJSONString(msg)).get("data");
System.out.println( new String(Base64.decode( data )) );
});
d.unsubscribe(reply , 1);
//订阅请求
d.subscribe(reply);
//发布请求
System.out.println( "订阅信息:"+reply );
nc.publish("requestSub", reply, "请求参数,巴拉巴拉1".getBytes(StandardCharsets.UTF_8));
// 下面这些用来负载测试的
nc.publish("requestSub", reply, "请求参数,巴拉巴拉2".getBytes(StandardCharsets.UTF_8));
nc.publish("requestSub", reply, "请求参数,巴拉巴拉3".getBytes(StandardCharsets.UTF_8));
nc.publish("requestSub", reply, "请求参数,巴拉巴拉4".getBytes(StandardCharsets.UTF_8));
nc.publish("requestSub", reply, "请求参数,巴拉巴拉5".getBytes(StandardCharsets.UTF_8));
nc.publish("requestSub", reply, "请求参数,巴拉巴拉6".getBytes(StandardCharsets.UTF_8));
}
/**
* test 发布订阅(Publish Subscribe) 模式
*/
public static void publishSubscribe() throws IOException, InterruptedException {
Connection nc = Nats.connect("nats://192.168.137.xxx:4222");
nc.publish("subject", "hello world1111122211111111".getBytes(StandardCharsets.UTF_8));
}
}
/**
* 订阅Demo
*/
public class NatsSubscribe {
public static void main(String[] args) throws IOException, InterruptedException {
// publishSubscribe();
requestReply();
}
/**
* test 请求响应(Request Reply) 模式
*/
public static void requestReply() throws IOException, InterruptedException {
//subscribe
Connection nc = Nats.connect("nats://192.168.137.xxx:4222");
//注册订阅
Dispatcher dispatcher = nc.createDispatcher(msg -> {
System.out.println("=======收到请求信息===========");
System.out.println(JSON.toJSONString(msg));
String data = (String) JSON.parseObject(JSON.toJSONString(msg)).get("data");
System.out.println( new String(Base64.decode( data )) );
nc.publish(msg.getReplyTo(), "这个是返回的数据,啦啦啦啦啦".getBytes(StandardCharsets.UTF_8));
});
dispatcher.subscribe("requestSub");
// 队列订阅就换成下面这个,负载测试,都启动几个服务,就可以看到接受效果了
// dispatcher.subscribe("requestSub", "queName");
}
/**
* test 发布订阅(Publish Subscribe) 模式
*/
public static void publishSubscribe() throws IOException, InterruptedException {
Connection nc = Nats.connect("nats://192.168.137.xxx:4222");
// //subscribe [这个时间内就之后收到一个,就结束了]
// Subscription sub = nc.subscribe("subject");
// Message msg = sub.nextMessage(Duration.ofMillis(50000));
// String response = new String(msg.getData(), StandardCharsets.UTF_8);
// System.out.println(response);
//subscribe [这个程序可以保持,持续接收信息]
Dispatcher d = nc.createDispatcher(msg ->{
String response = new String(msg.getData(), StandardCharsets.UTF_8);
//do something
System.out.println(response);
});
d.subscribe("subject");
}
}
简单安装使用与测试
# 官方安装NATS[单台]
docker pull nats
docker network create nats
docker run --name nats --network nats -p 4222:4222 -p 8222:8222 nats --http_port 8222 -js
# 192.168.137.xxx : 4222
# 然后用,上文中小杭的Demo试试,基础的功能就可以了解了。
JetStream 简单使用Demo
目前这个的Demo使用的是官方的封装例子方法。
结果是,创建流之后,发送数据。消费端接入会获取全部数据。除非消息被删除,否则每次都是全部获取。
当然,正常获取的时候,由于持久化,只要没有删除,消费端都可以请求再次获取的。
// 创建发送流 和 数据
public static void main(String[] args) throws Exception {
jetStream(args);
}
public static void jetStream(String[] args) throws Exception {
ExampleArgs exArgs = ExampleArgs.builder("Publish", args, "")
.defaultStream("example-stream")
.defaultSubject("example-subject")
.defaultMessage("hello")
.defaultMsgCount(10)
.defaultServer("nats://192.168.137.xxx:4222")
.build();
String hdrNote = exArgs.hasHeaders() ? ", with " + exArgs.headers.size() + " header(s)" : "";
System.out.printf("\nPublishing to %s%s. Server is %s\n\n", exArgs.subject, hdrNote, exArgs.server);
try (Connection nc = Nats.connect(ExampleUtils.createExampleOptions(exArgs.server))) {
JetStream js = nc.jetStream();
// Create the stream
NatsJsUtils.createStreamOrUpdateSubjects(nc, exArgs.stream, exArgs.subject);
int stop = exArgs.msgCount < 2 ? 2 : exArgs.msgCount + 1;
for (int x = 1; x < stop; x++) {
// make unique message data if you want more than 1 message
String data = exArgs.msgCount < 2 ? exArgs.message : exArgs.message + "-" + x;
// create a typical NATS message
Message msg = NatsMessage.builder()
.subject(exArgs.subject)
.headers(exArgs.headers)
.data(data, StandardCharsets.UTF_8)
.build();
PublishAck pa = js.publish(msg);
System.out.printf("Published message %s on subject %s, stream %s, seqno %d.\n",
data, exArgs.subject, pa.getStream(), pa.getSeqno());
}
}
catch (Exception e) {
e.printStackTrace();
}
}
// 消费,并删除流中已处理的数据
public static void main(String[] args) throws Exception {
jetStream();
}
public static void jetStream() throws IOException, InterruptedException, JetStreamApiException {
Connection nc = Nats.connect("nats://192.168.137.xxx:4222");
Dispatcher disp = nc.createDispatcher(msg -> {
System.out.println("ddddddd"+msg);
});
JetStream js = nc.jetStream();
MessageHandler handler = (msg) -> {
// Process the message.
// Ack the message depending on the ack model
String response = new String(msg.getData(), StandardCharsets.UTF_8);
//do something
System.out.println(response);
System.out.println(msg);
System.out.println("处理一下数据,然后要删除掉!!");
System.out.println(msg.metaData());
try {
// 处理完数据,要把数据删除掉的,否则会一直在持久队列中。
JetStreamManagement jsm = nc.jetStreamManagement();
jsm.deleteMessage(msg.metaData().getStream(),msg.metaData().streamSequence());
} catch (IOException | JetStreamApiException e) {
e.printStackTrace();
}
};
boolean autoAck = true;
js.subscribe("example-subject", disp, handler, autoAck);
}
一些复杂的功能,还是有需要的时候,再研究一下官方的Demo程序,比文档好理解多了。。。。
Spring 项目整合
参考开源项目:wanlinus/nats-streaming-spring
代码直接打包了;这里就记录一下使用。
// pom
<dependency>
<groupId>io.nats</groupId>
<artifactId>jnats</artifactId>
<version>2.16.13</version>
<scope>compile</scope>
</dependency>
// 配置
spring:
nats:
natsUrls: nats://192.168.137.xxx:4222
// 启动类
@EnableNats
@SpringBootApplication
public class AppApplication {
// 测试类
@Component
@RestController
@RequestMapping("/test")
public class TestController extends BaseController {
@Autowired
private Connection cconnection;
@GetMapping("/test")
public String test(HttpServletRequest request){
String msg = "send msg " + DateUtil.now();
// 测试发送普通消息
cconnection.publish("xixi", msg.getBytes(StandardCharsets.UTF_8));
return "test-success";
}
/**
* 接收 JetStream 的消息
* @param message
*/
@Subscribe(value="haha",type = "JetStream")
public void message1(Message message) {
System.out.println("接收 JetStream 的消息,进行处理。。。。。。");
System.out.println(message);
System.out.println(message.getSubject() + " : " + new String(message.getData()));
}
/**
* 接收普通消息
* @param message
*/
@Subscribe(value="xixi")
public void message2(Message message) {
System.out.println("接收普通消息,进行处理。。。。。。");
System.out.println(message);
System.out.println(message.getSubject() + " : " + new String(message.getData()));
}
}
其他类型的封装 和 发送操作,就真实需要的时候再继续完善一下了。
Nkey 认证连接
AuthHandler authHandler = Nats.staticCredentials("UCVU4OEHWAxxxxxxxxxxxxDDIxxxxxBMYxxxxxxxxxxxxxxxxxx".toCharArray(),"SUAMMIOB6xxxxxxxxxxxxxxxxxSHYxxxx7MUxxxxxxxxxxx5FCI".toCharArray());
Options.Builder builder = new Options.Builder()
// 配置 nats 服务器地址
.servers(new String[]{"nats://xxxx.xxxx.xxx:4222"})
.authHandler(authHandler);
Connection nc = Nats.connect(builder.build());
参考资料
- 简单看看:https://www.jianshu.com/p/341082dadd3e
- 详细点说明:http://www.guoxiaolong.cn/blog/?id=10376
- JetStream:https://docs.nats.io/nats-concepts/jetstream
- Developing With NATS:https://docs.nats.io/using-nats/developer
- https://docs.nats.io/running-a-nats-service/nats_docker/nats-docker-tutorial
- 官方javademo:https://github.com/nats-io/nats.java 参考这个 【这个重点的样子】
- 发布订阅:https://blog.****.net/qq_47848696/article/details/117746807
- https://zhuanlan.zhihu.com/p/628371358 用户+密码连接