今天学习了掘金小册里面的《Java开发者的RPC实战课》,写了一些笔记给大家分享一下(还没看完所以笔记还不完全✨)


Java开发者的RPC



远程调用


  1. HTTP :HTTP协议
  2. RPC : TCP/IP协议

性能:RPC>HTTP


PRC常用框架


Dobbo、Grpc、Thritf

合格的rpc功能包含:服务发现机、负载均衡策略、容错降级、调用重试、异步队列、事件解耦


认识PRC


rpc的核心组件:

  1. Server :服务端
  2. Client:客户端
  3. Server Stub:服务端收到Client发送的数据进行消息解包,调用本地方法
  4. Client Stub: 将客户端的请求参数、服务名称、服务地址进行统一打包,发送给Server方


整体结构分析

调用流程分析


首先本地客户端通知到了本地存根(sub),接着本地存根需要对数据格式进行包装、网络请求进行封装,然后按照一定的规则将数据包发送的目标机器上
服务端收到传过来的数据包然后对其按照事先约定好的规则进行解码,从而识别到数据包内部信息。将请求转发到本地函数中进行处理,处理完的数据进行返回给客户端
客户端存根接收到服务的时候需要对其进行解码。最后得到了最终结果。

简单的一个设计流程,在此基础上对功能进行不断地扩充,最终落实成一个成熟的rpc框架


代理层的设计


从使用角度来说,希望使用这款框架期间,在远程调用的时候能隐藏其内部细节,让其像调用本地方法一样方便

通过代码理解

public class Client {
   public static void main(String[] args) {
       //调用一次远程服务
       Server server = new Server("127.0.0.1",9999);
       server.doConnect();
       Object sendResponse = server.doRef("sendSms","这是一条短信信息",10001);
       System.out.println(sendResponse);
  }
   
}

我们希望在本地应用进行服务的连接,调用短信发送功能,进行远程调用,然后得到数据的返回,这段代码看似是简单的调用,实则隐藏了许多内部的细节


server.doRef 内部就是一个代理的手段,内部可以设计一个统一的代理组件,辅助开发者发送远程服务的调用,最后获取数据的返回

在此可以联想到代理模式,需要给某对象提供一个访问的代理,访问对象不适合或不能直接访问目标对象,就可以将代理对象作为访问对象跟目标对象之间的媒介


代码模式的优点:

  1. 代理对象能作为客户端与目标对象的媒介,能保护目标对象;
  2. 代理对象能扩展目标对象的功能;
  3. 代理对象能将客户端与目标对象分离,起到了解耦的作用,增加了程序的可扩展性


面对客户端发送请求我们可以设计一个代理层,处理一些内部细节,将内部细节隐蔽起来,让调用者无感知。以下是模拟一张请求调用的流程图


路由层的设计


问题的引出:当目标服务众多,客户端需要怎么确认最终的请求服务

设计思路:引入一个路由的角色,由他选择最终的请求服务

客户端请求都会经过一个路由层,由路由层内部规则去匹配对应的provider服务

路由层设计的时候需要考虑的点:

  1. 如何获取到provider服务
  2. 如何从集群服务中做筛选
  3. 如何设计能较好地兼容后期路由的扩展功能


协议层的设计


客户端在使用RPC框架进行远程调用的时候,需要对数据进行统一的包装和组织,最终才发送到目标机器进行接收解析。因此对数据的各种序列化、反序列化、协议的组装可以统一封装在协议层中。

router模块会负责提供最终需要调用服务提供者的具体信息,然后将对应的地址信息、请求参数传入到Protocol中,最终由协议层对数据封装成对应的协议体,然后序列化处理后发送给服务端


链路层的设计


从本地请求,到protocol发送数据,整个链路中要考虑一个可扩展性,比如一些自定义的过滤、服务分组等,可通过新增一个链路模块,类似于责任链的模式


注册中心层的设计


当服务提供者呈现集群模式的时候,客户端需获取到Provider的多个信息,这时我们需引入一个注册中心的角色

服务提供者会将自己的地址、接口、分组等详细信息都上报到注册中心模块,并且当服务上线、下线的时候通知到注册中心,服务调用方只需要订阅注册中心即可

市面上常见的注册中心:Nacos、Zookeeper、etcd、Eureka、Redis等

需要重点关注的点:

  1. 如何与注册中心进行基本的连接访问
  2. 如何监听服务数据在注册中心的实时变化
  3. 如果注册中心出现了异常,需要有哪些安全手段


容错层的设计


在远程调用的时候可能会出现一些异常情况,市面上常见的一些RPC框架都会提供一些容错方面的处理手段:

  1. 超时重试
  2. 快速失败
  3. 无线重试
  4. 出现异常后回调指定方法
  5. 无视失败

面对这些异常,我们可以抽象出一层 容错层 进行处理


服务提供者的线程池设计


当请求发送到了服务提供者的时候,服务提供方需要对其进行解码,然后再本地进行核心处理,在解码这一步骤中我们统一交给专门的线程进行处理

涉及的相关技术点:

  1. io线程与worker线程的拆分
  2. 调用结果和客户端请求的唯一匹配
  3. 客户端请求后的同步转为异步处理
  4. 单一请求队列和多请求队列的设计差异性


接入层的设计


整套RPC组件基本设计实现后,要考虑如何将其接入到实际开发项目当中,团队主要技术使用Spring,所以这款RPC框架也应该介入到starter组件中,让使用Spring的技术团队更好的接入使用,最后是整体的架构图


小结


对RPC框架整体进行基本的分层:

  1. 代理层:负责底层调用细节的封装
  2. 路由层:通过调用策略来选择目标服务
  3. 协议层:服务数据的转码封装、序列化、反序列化等
  4. 链路层:负责执行一些自定义的过滤链路,可供后期的二次扩展
  5. 注册中心层:关注服务的上下线,以及一些权重,配置动态调整等功能
  6. 容错层:当出现服务调用失败,应该执行的一些策略、兜底等功能
  7. 接入层:如何将框架接入常用框架Spring中
  8. 公共层:主要存放一些通用配置、工具类、缓存等信息


本地调用和RPC调用的区别


形象的比喻:
事:人要吃苹果,尝试苹果的甜度
本地调用:我在果园,苹果摘下来我就可以吃了,当场发表意见,苹果甜不甜。
PRC调用:人不在果园,果农将苹果摘下,打包,装箱,运输,最后送达接收人的手中,接收人拆箱,尝试苹果,通过微信发消息,告诉我苹果的甜度。


理解网络通信模型的核心

阻塞IO技术


基于BIO实现的阻塞IO服务端程序代码

/**
*
* @author cong
* @date 2023/07/25
*/
public class BioServer {
   private static ExecutorService executors = Executors.newFixedThreadPool(10);
   public static void main(String[] args) throws IOException {
       ServerSocket serverSocket = new ServerSocket();
       serverSocket.bind(new InetSocketAddress(1009));
       try {
           while (true) {
              //堵塞状态点--1
               Socket socket = serverSocket.accept();
               System.out.println("获取新连接");
               executors.execute(new Runnable() {
                   @Override
                   public void run() {
                       while (true){
                           InputStream inputStream = null;
                           try {
                               //堵塞的状态点--2
                               inputStream = socket.getInputStream();
                               byte[] result = new byte[1024];
                               int len = inputStream.read(result);
                               if(len!=-1){
                                   System.out.println("[response] "+new String(result,0,len));
                                   OutputStream outputStream = socket.getOutputStream();
                                   outputStream.write("response data".getBytes());
                                   outputStream.flush();
                              }
                          } catch (IOException e) {
                               e.printStackTrace();
                               break;
                          }
                      }
                  }
              });
          }
      }catch (Exception e){
           e.printStackTrace();
      }
  }
}

当前代码主要关注两个函数 accept、read 在传统的BIO技术中会在这两个模块发生阻塞

服务端创建了Socket之后会堵塞住等待accept函数的连接

当客户端连接上了服务端后,accept的堵塞状态就会放开进入到read环节(读取客户端传过来的网络数据)

客户端如果一直没有发送数据过来,那么服务端的read调用方法就会一直处于堵塞状态,如果数据通过网络抵达了网卡缓冲区,此时则会将数据从内核态拷贝至用户态,然后返回给read调用方。

出现问题:如果客户端只是连上了服务端,但是没有进行发送数据,那么服务端就会一直处于阻塞状态,此时引出一种技术叫 非阻塞IO技术。


非阻塞IO技术


对阻塞IO进行性能改造,可以尝试让其在accept函数获取到客户端连接后,专门创建一个线程来处理read函数,整体流程如下图:

是否疑惑?为什么来一个请求再创建一个线程,这部分工作交给线程池去完成不好吗?


没错,线程池能提强帮助我们创建好一点数量的线程,当请求量大的时候还能有队列缓冲,以及增加worker线程的效果,从而降低这块实现的难度。但是这样的设计依然存在不足点,因为在用户态层面调用的read函数依然是阻塞的。就现阶段而言,这种技术方案还不能称之为非阻塞IO技术。

如果read函数在调用的时候没有数据抵达,他还是处于阻塞状态,如何对其进行设计升级非阻塞效果呢


JDK的NIO模型中就有相关的设计存在,下面是一段简单的NIO服务端代码

public class NioSocketServer extends Thread {
   ServerSocketChannel serverSocketChannel = null;
   Selector selector = null;
   SelectionKey selectionKey = null;
   public void initServer() throws IOException {
       selector = Selector.open();
       serverSocketChannel = ServerSocketChannel.open();
       //设置为非阻塞模式,默认serverSocketChannel是采用了阻塞模式
       serverSocketChannel.configureBlocking(false);
       serverSocketChannel.socket().bind(new InetSocketAddress(8888));
       selectionKey = serverSocketChannel.register(selector, SelectionKey.OP_ACCEPT);
  }
   @Override
   public void run() {
       while (true) {
           try {
               //默认这里会堵塞
               int selectKey = selector.select();
               if (selectKey > 0) {
                   //获取到所有的处于就绪状态的channel,selectionKey中包含了channel的信息
                   Set<SelectionKey> keySet = selector.selectedKeys();
                   Iterator<SelectionKey> iter = keySet.iterator();
                   //对selectionkey进行遍历
                   while (iter.hasNext()) {
                       SelectionKey selectionKey = iter.next();
                       //需要清空,防止下次重复处理
                       iter.remove();
                       //就绪事件,处理连接
                       if (selectionKey.isAcceptable()) {
                           accept(selectionKey);
                      }
                       //读事件,处理数据读取
                       if (selectionKey.isReadable()) {
                           read(selectionKey);
                      }
                       //写事件,处理写数据
                       if (selectionKey.isWritable()) {
                      }
                  }
              }
          } catch (IOException e) {
               e.printStackTrace();
               try {
                   serverSocketChannel.close();
              } catch (IOException e1) {
                   e1.printStackTrace();
              }
          }
      }
  }
   public void accept(SelectionKey key) {
       try {
           ServerSocketChannel serverSocketChannel = (ServerSocketChannel) key.channel();
           SocketChannel socketChannel = serverSocketChannel.accept();
           System.out.println("conn is acceptable");
           socketChannel.configureBlocking(false);
           //将当前的channel交给selector对象监管,并且有selector对象管理它的读事件
           socketChannel.register(selector, SelectionKey.OP_READ);
      } catch (IOException e) {
           e.printStackTrace();
      }
  }
   public void read(SelectionKey selectionKey) {
       try {
           SocketChannel channel = (SocketChannel) selectionKey.channel();
           ByteBuffer byteBuffer = ByteBuffer.allocate(100);
           int len = channel.read(byteBuffer);
           if (len > 0) {
               byteBuffer.flip();
               byte[] byteArray = new byte[byteBuffer.limit()];
               byteBuffer.get(byteArray);
               System.out.println("NioSocketServer receive from client:" + new String(byteArray,0,len));
               selectionKey.interestOps(SelectionKey.OP_READ);
          }
      } catch (Exception e) {
           try {
               serverSocketChannel.close();
               selectionKey.cancel();
          } catch (IOException e1) {
               e1.printStackTrace();
          }
           e.printStackTrace();
      }
  }
   public static void main(String args[]) throws IOException {
       NioSocketServer server = new NioSocketServer();
       server.initServer();
       server.start();
  }
}

以下是上述代码的流程图

思考一下接下来代码会怎么走,当socket的服务端启动之后,会对每一个socket连接的对象都开启一个线程,然后在循环里面去调read函数,此时的read函数调用不会进入阻塞状态额,但是似乎没有解决掉根本性问题:每次请求来的时候都要创建一个线程来监听客户端的请求。如果客户端在建立连接之后长时间都没有传输数据,那对于服务端而言就会造成资源浪费的情况。

每个请求都要建立一个线程,如何进行优化

我们不妨可以将accept和read分成两个模块来处理,当accept函数接收到了新的接连(本质就是一个文件描述符fd)之后,将其放入一个集合中,然后会有一个后台任务统一对集合中的fd进行遍历执行read操作。 流程如下:

看到这里有些疑问,循环调用read方法岂不是循环进行用户态和内核态之间切换?这样不断地进行上下文切换也不是什么好的设计思路啊,能否将这个循环的操作交给内核态处理呢


待续。。。。


0个评论
点击登录,快来和大家讨论吧~
表情
图片
暂无评论
聪
下载 APP