NATS-研究学习

NATS-研究学习


文章目录

    • NATS-研究学习
    • @[toc]
      • 介绍说明
      • 提供的服务内容
      • 各模式介绍测试使用
        • 发布订阅(Publish Subscribe)
        • 请求响应(Request Reply)
        • 队列订阅&分享工作(Queue Subscribers & Sharing Work)
        • 小杭写的Demo
      • 简单安装使用与测试
      • JetStream 简单使用Demo
      • Spring 项目整合
      • Nkey 认证连接
      • 参考资料

介绍说明

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 somethingSystem.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 {//subscribeConnection 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 somethingSystem.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 streamNatsJsUtils.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 messageString data = exArgs.msgCount < 2 ? exArgs.message : exArgs.message + "-" + x;// create a typical NATS messageMessage 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 modelString response = new String(msg.getData(), StandardCharsets.UTF_8);//do somethingSystem.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 {@Autowiredprivate 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.csdn.net/qq_47848696/article/details/117746807
  • https://zhuanlan.zhihu.com/p/628371358 用户+密码连接

本文来自互联网用户投稿,该文观点仅代表作者本人,不代表本站立场。本站仅提供信息存储空间服务,不拥有所有权,不承担相关法律责任。如若转载,请注明出处:http://www.mzph.cn/web/20502.shtml

如若内容造成侵权/违法违规/事实不符,请联系多彩编程网进行投诉反馈email:809451989@qq.com,一经查实,立即删除!

相关文章

运放的自激振荡问题

运放的自激振荡指的是当运算放大器加电后&#xff0c;在没有外部信号输入的情况下&#xff0c;输出端会出现高频类似于正弦波的波形。 运算放大器产生自激的原因以及解决办法-CSDN博客 a)当振荡由分布电容、电感等引起时&#xff0c;可通过反馈端并联电容&#xff0c;抵消影响…

【开源】课程管理平台 JAVA+Vue.js+SpringBoot+MySQL

目录 一、项目介绍 课程管理模块 作业题目模块 考试阅卷模块 教师评价模块 部门角色菜单模块 二、项目截图 三、核心代码 一、项目介绍 Vue.jsSpringBoot前后端分离新手入门项目《课程管理平台》&#xff0c;包括课程管理模块、作业题目模块、考试阅卷模块、教师评价模…

spoon工具的安装与配置

spoon对应的jdk包下载资源地址 spoon软件下载资源地址 首先需要安装jdk&#xff0c;配置java环境&#xff0c;安装好后&#xff0c;cmd一下&#xff0c;查看java -version&#xff0c;看看是否成功安装&#xff0c;如果失败&#xff0c;查看系统环境变量&#xff0c;去配置jdk…

Java | Leetcode Java题解之第122题买卖股票的最佳时机II

题目&#xff1a; 题解&#xff1a; class Solution {public int maxProfit(int[] prices) {int ans 0;int n prices.length;for (int i 1; i < n; i) {ans Math.max(0, prices[i] - prices[i - 1]);}return ans;} }

Python保存为json中文Unicode乱码解决json.dump()

保存为json中文Unicode乱码&#xff1a; 可以看到&#xff0c;中文字符没有乱码&#xff0c;只是出现了反斜杠&#xff0c;此时解决方法应考虑是否进行了二次序列化。 一、原因1 在dump时加入ensure_asciiFalse 即可解决&#xff0c;即json.dump(json_data, f, indent4, en…

阿里云布置net core 项目

一、 创建镜像 给镜像添加触发器&#xff0c;编译的时候会触发k8s集群里的taget链接&#xff0c;从而更新项目 二&#xff0c;创建k8s集群 使用镜像创建 添加基本信息 镜像名称&#xff1a;镜像仓库》基本信息公网地址镜像Tag:创建镜像时的镜像版本镜像配置为&#xff1a;总…

opencv笔记(13)—— 停车场车位识别

一、所需数据介绍 car1.h5 是训练后保存的模型 class_directionary 是0&#xff0c;1的分类 二、图像数据预处理 对输入图片进行过滤&#xff1a; def select_rgb_white_yellow(self,image): #过滤掉背景lower np.uint8([120, 120, 120])upper np.uint8([255, 255, 255])#…

C#WPF数字大屏项目实战04--设备运行状态

1、引入Livecharts包 项目中&#xff0c;设备运行状态是用饼状图展示的&#xff0c;因此需要使用livechart控件&#xff0c;该控件提供丰富多彩的图形控件显示效果 窗体使用控件 2、设置饼状图的显示图例 通过<lvc:PieChart.Series>设置环状区域 3、设置饼状图资源样…

【TB作品】MSP430G2553单片机,智能储物柜

智能储物柜将实现的功能&#xff1a; 1在超市或者机场场景下&#xff0c;用户需要进行物品暂存时。按下储物柜键盘的需求按键&#xff0c;智能储物柜将会随机为用户分配一个还没使用的柜子&#xff0c;屏幕提示用户选择密码存储方式或者身份证存储方式&#xff1b; 2 用户选择密…

禁止Windows Defender任务计划程序

开始键->搜索“任务计划程序”->“任务计划程序库”->“Microsoft”->"Windows"->"Windows Defender"->右边四项

43-5 waf绕过 - 安全狗简介及安装

一、安全狗安装 安装安全狗需要开启 Apache 系统服务。如果 Apache 系统服务未开启,安装过程中可能会出现无法填入服务名称的问题,导致无法继续安装。为避免此问题,可以先在虚拟机中安装 PHPStudy。 安装PHPStudy 下载、安装phpstudy并启动(安装过程可以一路下一步,也…

安装VS2017后,离线安装Debugging Tools for Windows(QT5.9.2使用MSVC2017 64bit编译器)

1、背景 安装VS2017后&#xff0c;Windows Software Development Kit - Windows 10.0.17763.132的Debugging Tools for Windows默认不会安装&#xff0c;如下图。这时在QT5.9.2无法使用MSVC2017 64bit编译器。 2、在线安装 如果在线安装参考之前的文章&#xff1a; Qt5.9.2初…

Windows操作系统提权之系统服务漏洞提权Always Install Elevated

Always Install Elevated 1.形成原因 任意用户以NT AUTHORITY\SYSTEM权限安装 i。 AlwaysInstallElevated是一个策略设置&#xff0c;当在系统中使用Windows Installer安装任何程序时&#xff0c;该参数允许非 特权用户以system权限运行MSI文件。如果目标系统上启用了这一设…

leetcode - 20.有效的括号(LinkedHashMap)

leetcode题目有效的括号&#xff0c;分类是easy&#xff0c;但是博主前前后后提交了几十次才通过&#xff0c;现在记录一下使用Java语言的写法。 题目链接: 20.有效的括号 题目描述&#xff1a; 给定一个只包括 (&#xff0c;)&#xff0c;{&#xff0c;}&#xff0c;[&…

【漏洞复现】WordPress Country State City Dropdown CF7插件 SQL注入漏洞(CVE-2024-3495)

0x01 产品简介 Country State City Dropdown CF7插件是一个功能强大、易于使用的 WordPress 插件&#xff0c;它为用户在联系表单中提供国家.州/省和城市的三级下拉菜单功能&#xff0c;帮助用户更准确地填写地区信息。同时&#xff0c;插件的团队和支持也非常出色&#xff0c…

香橙派 Kunpeng Pro使用教程:从零开始打造个人私密博客

一、引言 在这个日益互联的世界中&#xff0c;单板计算机已经成为创新和个性化解决方案的重要载体。而在单板计算机领域&#xff0c;香橙派 Kunpeng Pro凭借其强大的性能和灵活的应用潜力&#xff0c;正逐渐吸引着全球开发者和技术爱好者的目光。 作为一款集成了华为的鲲鹏处…

单调栈原理+练习

首先用一道题引出单调栈 码蹄集 (matiji.net) 首先画一个图演示山的情况&#xff1a; 最暴力的做法自然是O(n方)的双循环遍历&#xff0c;这么做的思想是求出当前山右侧有多少座比它小的山&#xff0c;遇见第一个高度大于等于它的就停止。 但是对于我们所求的答案数&#xff…

codefun的蓝桥杯国赛之旅

前言 好久没有刷算法了&#xff0c;今天完成了我的蓝桥杯国赛之旅&#xff01; 总的来说&#xff0c;比赛的过程不是很顺利&#xff0c;只能ac两道题目&#xff0c;好多题都是有思路&#xff0c;但是要么是写不出来&#xff0c;要么是debug不出来&#xff0c;多重背包&#xf…

Docker(Centos7+)

先确定是否 Centos 7 及以上的版本 查看是否 ping 通外网 linux centos7运行下面的代码&#xff0c;基本上都可以正常安装 # 删除之前的docker残留 yum -y remove docker*yum install -y yum-utilsyum-config-manager --add-repo http://mirrors.aliyun.com/docker-ce/linux/…

Controller类明明写了@CrossOrigin跨域注解,但还是有跨域问题

可能是写的过滤器干扰到了跨域处理。如&#xff1a; 此时&#xff0c;先注释掉过滤器注解&#xff0c;让其不生效&#xff0c;就可以避免干扰跨域处理了 不过&#xff0c;这只能暂时解决该问题&#xff0c;毕竟过滤器还是要用的&#xff0c;后续我再探索一下。。。。。。。