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 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.csdn.net/qq_47848696/article/details/117746807
  • https://zhuanlan.zhihu.com/p/628371358 用户+密码连接

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

如若内容造成侵权/违法违规/事实不符,请联系我们进行投诉反馈qq邮箱809451989@qq.com,一经查实,立即删除!

相关文章

编程入门(七)【虚拟机VMware安装Linux系统Ubuntu】

读者大大们好呀&#xff01;&#xff01;!☀️☀️☀️ &#x1f525; 欢迎来到我的博客 &#x1f440;期待大大的关注哦❗️❗️❗️ &#x1f680;欢迎收看我的主页文章➡️寻至善的主页 文章目录 &#x1f525;前言&#x1f680;Ubuntu知多少&#x1f680;安装的前期准备&am…

放开了去的 ulimit

放开了去的 ulimit 放开了去的 ulimitulimit简介临时修改打开文件数目永久修改系统总打开句柄限制更多信息 放开了去的 ulimit ulimit简介 对于高并发或者频繁读写文件的应用程序而言&#xff0c;有时可能需要修改系统能够打开的最多文件句柄数&#xff0c;否则就可能会出现t…

3389,为了保障3389端口的安全,我们可以采取的措施

3389端口&#xff0c;作为远程桌面协议&#xff08;RDP&#xff09;的默认端口&#xff0c;广泛应用于Windows操作系统中&#xff0c;以实现远程管理和控制功能。然而&#xff0c;正因为其广泛使用&#xff0c;3389端口也成为许多潜在安全威胁的入口。因此&#xff0c;确保3389…

使用C#实现VS窗体应用——画图板

✅作者简介&#xff1a;大家好&#xff0c;我是 Meteors., 向往着更加简洁高效的代码写法与编程方式&#xff0c;持续分享Java技术内容。&#x1f34e;个人主页&#xff1a;Meteors.的博客&#x1f49e;当前专栏&#xff1a;小项目✨特色专栏&#xff1a; 知识分享&#x1f96d…

ARM32开发——总线与时钟

&#x1f3ac; 秋野酱&#xff1a;《个人主页》 &#x1f525; 个人专栏:《Java专栏》《Python专栏》 ⛺️心若有所向往,何惧道阻且长 文章目录 APB总线时钟树时钟树 外部晶振内部晶振 在这个例子中&#xff0c;这条大街和巴士构成了一套系统&#xff0c;我们称之为AHB总线。 …

【数据结构】详解二叉树

文章目录 1.树的结构及概念1.1树的概念1.2树的相关结构概念1.3树的表示1.4树在实际中的应用 2.二叉树的结构及概念2.1二叉树的概念2.2特殊的二叉树2.2.1满二叉树2.2.2完全二叉树 2.3 二叉树的性质2.4二叉树的存储结构2.4.1顺序结构2.4.2链表结构 1.树的结构及概念 1.1树的概念…

乐观锁 or 悲观锁 你怎么选?

你有没有听过这样一句话&#xff1a;悲观者正确&#xff0c;乐观者成功​。那么今天我来分享下什么是乐观锁​和悲观锁。 乐观锁和悲观锁有什么区别&#xff0c;它们什么场景会用 乐观锁 乐观锁基于这样的假设&#xff1a;多个事务在同一时间对同一数据对象进行操作的可能性很…

三体中的冯诺依曼

你叫冯诺依曼&#xff0c;是一位科学家。你无法形容眼前的现态&#xff0c;你不知道下一次自己葬身火海会是多久&#xff0c;你也不知道会不会下一秒就会被冰封&#xff0c;你唯一知道的&#xff0c;就是自己那寥寥无几的科学知识&#xff0c;你可能会抱着他们终身&#xff0c;…

基于Android Studio记事本系统

目录 项目介绍 图片展示 运行环境 获取方式 项目介绍 具有登录&#xff0c;注册&#xff0c;记住密码&#xff0c;自动登录的功能&#xff1b; 可以新增记事本&#xff0c;编辑&#xff0c;删除记事本信息&#xff0c;同时可以设置主标题&#xff0c;内容&#xff0c;以及…

SpringBoot【1】集成 Druid

SpringBoot 集成 Druid 前言创建项目修改 pom.xml 文件添加配置文件开发 java 代码启动类 - DruidApplication配置文件-propertiesDruidConfigPropertyDruidMonitorProperty 配置文件-configDruidConfig 控制层DruidController 运行验证Druid 的监控应用程序 前言 JDK版本&…

【HarmonyOS - ArkTS - 状态管理】

概述 本文主要是从页面和应用级两个方面介绍了ArkTS中的状态管理的一些概念以及如何使用。可能本文比较粗略&#xff0c;细节化请前往官网(【状态管理】)学习&#xff0c;若有理解错误的地方&#xff0c;欢迎评论指正。 装饰器总览 由于下面会频繁提及到装饰器&#xff0c;所…

【CH32V305FBP6】调试入坑指南

1. 无法烧录程序 现象 MounRiver Studio WXH-LinkUtility 解决方法 前提&#xff1a;连接复位引脚 或者 2. 无法调试 main.c 与调试口冲突&#xff0c;注释后调试 // USART_Printf_Init(115200);

orin部署tensorrt、cuda、cudnn、pytorch、onnx

绝大部分参考https://blog.csdn.net/qq_41336087/article/details/129661850 非orin可以参考https://blog.csdn.net/JineD/article/details/131201121 报错显卡驱动安装535没法安装、原始是和l4t-cuda的部分文件冲突 Options marked [*] produce a lot of output - pipe it t…

三方语言中调用, Go Energy GUI编译的dll动态链接库CEF

如何在其它编程语言中调用energy编译的dll动态链接库&#xff0c;以使用CEF 或 LCL库 Energy是Go语言基于LCL CEF开发的跨平台GUI框架, 具有很容易使用CEF 和 LCL控件库 interface 便利 示例链接 正文 为方便起见使用 python 调用 go energy 编译的dll 准备 系统&#x…

使用compile_commands.json配置includePath环境,解决vscode中引入头文件处有波浪线的问题

通过编译时生成的 compile_commands.json 文件自动完成对 vscode 中头文件路径的配置&#xff0c;实现 vscode 中的代码的自动跳转。完成头文件路径配置后&#xff0c;可以避免代码头部导入头文件部分出现波浪线&#xff0c;警告说无法正确找到头文件。 步骤 需要在 vscode 中…

Linux--进程间通信(1)(匿名管道)

目录 1.了解进程通信 1.1进程为什么要通信 1.2 进程如何通信 1.3进程间通信的方式 2.管道 2.1管道的初步理解 2.2站在文件描述符的角度-进一步理解管道 2.3 管道的系统调用接口&#xff08;匿名管道&#xff09; 2.3.1介绍接口函数&#xff1a; 2.3.2编写一个管道的代…

windows系统配置dns加快访问github 实用教程一(图文保姆级教程)

第一步、打开网页 https://tool.lu/ip IP地址查询 - 在线工具 输入www.github.com 或者github.com 点击网页查询按钮, 获取对应github网站对应的ip 完整操作步骤如上图所示,可以很清晰的看到github网站的ip显示地区是美国也就是说该网站服务器是在国外, 这也就是为什么我们在…

IDEA 中导入脚手架后该如何处理?

MySQL数据库创建啥的&#xff0c;没啥要说的&#xff01;自行配置即可&#xff01; 1.pom.xml文件&#xff0c;右键&#xff0c;add Maven Project …………&#xff08;将其添加为Maven&#xff09;【下述截图没有add Maven Project 是因为目前已经是Maven了&#xff01;&…

redis 高可用及哨兵模式 @by_TWJ

目录 1. 高可用2. redis 哨兵模式3. 图文的方式让我们读懂这几个算法3.1. Raft算法 - 图文3.2. Paxos算法 - 图文3.3. 区别&#xff1a; 1. 高可用 在 Redis 中&#xff0c;实现 高可用 的技术主要包括 持久化、复制、哨兵 和 集群&#xff0c;下面简单说明它们的作用&#xf…

今日分享丨按场景定制界面

遇到问题 我们在写文档或者代码时&#xff0c;会遇到需要书写重复或者类似内容的情况。快捷的做法是&#xff1a;先复制粘贴此相似内容&#xff0c;再修改差异。那么开发人员在设计界面的时候&#xff0c;也会遇到同类型的界面有重复的特性&#xff0c;比如报销类型的单据&…