RPC(Remote Procedure Call)—远程过程调用协议,它是一种通过网络从远程计算机程序上请求服务,而不需要了解底层网络技术的协议。RPC协议假定某些传输协议的存在,如TCP或UDP,为通信程序之间携带信息数据。在OSI网络通信模型中,RPC跨越了传输层和应用层。RPC使得开发包括网络分布式多程序在内的应用程序更加容易。
RPC采用客户机/服务器模式。请求程序就是一个客户机,而服务提供程序就是一个服务器。首先,客户机调用进程发送一个有进程参数的调用信息到服务进程,然后等待应答信息。在服务器端,进程保持睡眠状态直到调用信息的到达为止。当一个调用信息到达,服务器获得进程参数,计算结果,发送答复信息,然后等待下一个调用信息,最后,客户端调用进程接收答复信息,获得进程结果,然后调用执行继续进行。
原理
1. Client端获取一个 RPC 代理对象 proxy
2. 调用 proxy 上的方法, 被 InvocationHandler 实现类 Invoker 的 invoke() 方法捕获
3. invoke() 方法内将 RPC 请求封装成 Invocation 实例, 再向 Server 发送 RPC请求
4. Server端循环接收 RPC请求, 对每一个请求都创建一个 Handler线程处理
5. Handler线程从输入流中反序列化出 Invocation实例, 再调用 Server端的实现方法
6. 调用结束, 向 Client端返回调用结果
一. Invoker 类
InvocationHandler 的实现类
1/** 2 * InvocationHandler 接口的实现类 <br> 3 * Client端代理对象的方法调用都会被 Invoker 的 invoke() 方法捕获 4 */ 5public class Invoker implements InvocationHandler { 6 /** RPC协议接口的 Class对象 */ 7 private Class<?> intface; 8 /** Client 端 Socket */ 9 private Socket client; 10 /** 用于向 Server端发送 RPC请求的输出流 */ 11 private ObjectOutputStream oos; 12 /** 用于接收 Server端返回的 RPC请求结果的输入流 */ 13 private ObjectInputStream ois; 14 15 /** 16 * 构造一个 Socket实例 client, 并连接到指定的 Server端地址, 端口 17 * 18 * @param intface 19 * RPC协议接口的 Class对象 20 * @param serverAdd 21 * Server端地址 22 * @param serverPort 23 * Server端监听的端口 24 */ 25 public Invoker(Class<?> intface, String serverAdd, int serverPort) throws UnknownHostException, IOException { 26 this.intface = intface; 27 client = new Socket(serverAdd, serverPort); 28 } 29 30 @Override 31 public Object invoke(Object proxy, Method method, Object[] args) throws Throwable { 32 33 try { 34 // 封装 RPC请求 35 Invocation invocation = new Invocation(intface, method.getName(), method.getParameterTypes(), args); 36 // 打开 client 的输出流 37 oos = new ObjectOutputStream(client.getOutputStream()); 38 // 序列化, 将 RPC请求写入到 client 的输出流中 39 oos.writeObject(invocation); 40 oos.flush(); 41 42 // 等待 Server端返回 RPC请求结果 // 43 44 // 打开 client 的输入流 45 ois = new ObjectInputStream(client.getInputStream()); 46 // 反序列化, 从输入流中读取 RPC请求结果 47 Object res = ois.readObject(); 48 // 向 client 返回 RPC请求结果 49 return res; 50 } finally { // 关闭资源 51 CloseUtil.closeAll(ois, oos); 52 CloseUtil.closeAll(client); 53 } 54 } 55}
二. Invocation 类
Serializable 的实现类, RPC请求的封装
1/** 2 * RPC调用的封装, 包括以下字段: <br> 3 * methodName: 方法名 <br> 4 * parameterTypes: 方法参数列表的 Class 对象数组 <br> 5 * params: 方法参数列表 6 */ 7@SuppressWarnings("rawtypes") 8public class Invocation implements Serializable { 9 private static final long serialVersionUID = -7311316339835834851L; 10 /** RPC协议接口的 Class对象 */ 11 private Class<?> intface; 12 /** 方法名 */ 13 private String methodName; 14 /** 方法参数列表的 Class 对象数组 */ 15 private Class[] parameterTypes; 16 /** 方法的参数列表 */ 17 private Object[] params; 18 19 public Invocation() { 20 } 21 22 /** 23 * 构造一个 RPC请求的封装 24 * 25 * @param intface 26 * RPC协议接口的 Class对象 27 * @param methodName 28 * 方法名 29 * @param parameterTypes 30 * 方法参数列表的 Class 对象数组 31 * @param params 32 * 方法的参数列表 33 */ 34 public Invocation(Class intface, String methodName, Class[] parameterTypes, Object[] params) { 35 this.intface = intface; 36 this.methodName = methodName; 37 this.parameterTypes = parameterTypes; 38 this.params = params; 39 } 40 41 public Class getIntface() { 42 return intface; 43 } 44 45 public String getMethodName() { 46 return methodName; 47 } 48 49 public Class[] getParameterTypes() { 50 return parameterTypes; 51 } 52 53 public Object[] getParams() { 54 return params; 55 } 56}
三.RPC 类
构造 Client端代理对象, Server端实例
1/** 2 * 一个构造 Server 端实例与 Client 端代理对象的类 3 */ 4public class RPC { 5 6 /** 7 * 获取一个 Client 端的代理对象 8 * 9 * @param intface 10 * RPC协议接口, Client 与 Server 端共同遵守 11 * @param serverAdd 12 * Server 端地址 13 * @param serverPort 14 * Server 端监听的端口 15 * @return Client 端的代理对象 16 */ 17 public static <T> Object getProxy(final Class<T> intface, String serverAdd, int serverPort) 18 throws UnknownHostException, IOException { 19 20 Object proxy = Proxy.newProxyInstance(intface.getClassLoader(), new Class[] { intface }, 21 new Invoker(intface, serverAdd, serverPort)); 22 return proxy; 23 } 24 25 /** 26 * 获取 RPC 的 Server 端实例 27 * 28 * @param intface 29 * RPC协议接口 30 * @param intfaceImpl 31 * Server 端 RPC协议接口的实现 32 * @param port 33 * Server 端监听的端口 34 * @return RPCServer 实例 35 */ 36 public static <T> RPCServer getRPCServer(Class<T> intface, T intfaceImpl, int port) throws IOException { 37 return new RPCServer(intface, intfaceImpl, port); 38 } 39 40}
四. RPCServer 类
Server端接收 RPC请求, 处理请求
1/** 2 * RPC 的 Server端 3 */ 4public class RPCServer { 5 /** Server端的 ServerSocket实例 */ 6 private ServerSocket server; 7 /** Server端 RPC协议接口的实现缓存, 一个接口对应一个实现类的实例 */ 8 private static Map<Class<?>, Object> intfaceImpls = new HashMap<Class<?>, Object>(); 9 10 /** 11 * 构造一个 RPC 的 Server端实例 12 * 13 * @param intface 14 * RPC协议接口的 Class对象 15 * @param intfaceImpl 16 * Server端 RPC协议接口的实现 17 * @param port 18 * Server端监听的端口 19 */ 20 public <T> RPCServer(Class<T> intface, T intfaceImpl, int port) throws IOException { 21 server = new ServerSocket(port); 22 RPCServer.intfaceImpls.put(intface, intfaceImpl); 23 } 24 25 /** 26 * 循环监听并接收 Client端连接, 处理 RPC请求, 向 Client端返回结果 27 */ 28 public void start() { 29 try { 30 while (true) { 31 // 接收 Client端连接, 创建一个 Handler线程, 处理 RPC请求 32 new Handler(server.accept()).start(); 33 } 34 } catch (IOException e) { 35 e.printStackTrace(); 36 } finally { // 关闭资源 37 CloseUtil.closeAll(server); 38 } 39 } 40 41 /** 42 * 向 RPC协议接口的实现缓存中添加缓存 43 * 44 * @param intface 45 * RPC协议接口的 Class对象 46 * @param intfaceImpl 47 * Server端 RPC协议接口的实现 48 */ 49 public static <T> void addIntfaceImpl(Class<T> intface, T intfaceImpl) { 50 RPCServer.intfaceImpls.put(intface, intfaceImpl); 51 } 52 53 /** 54 * 处理 RPC请求的线程类 55 */ 56 private static class Handler extends Thread { 57 /** Server端接收到的 Client端连接 */ 58 private Socket client; 59 /** 用于接收 client 的 RPC请求的输入流 */ 60 private ObjectInputStream ois; 61 /** 用于向 client 返回 RPC请求结果的输出流 */ 62 private ObjectOutputStream oos; 63 /** RPC请求的封装 */ 64 private Invocation invocation; 65 66 /** 67 * 用 Client端连接构造 Handler线程 68 * 69 * @param client 70 */ 71 public Handler(Socket client) { 72 this.client = client; 73 } 74 75 @Override 76 public void run() { 77 try { 78 // 打开 client 的输入流 79 ois = new ObjectInputStream(client.getInputStream()); 80 // 反序列化, 从输入流中读取 RPC请求的封装 81 invocation = (Invocation) ois.readObject(); 82 // 从 RPC协议接口的实现缓存中获取实现 83 Object intfaceImpl = intfaceImpls.get(invocation.getIntface()); 84 // 获取 Server端 RPC协议接口的方法实现 85 Method method = intfaceImpl.getClass().getMethod(invocation.getMethodName(), 86 invocation.getParameterTypes()); 87 // 跳过安全检查 88 method.setAccessible(true); 89 // 调用具体的实现方法, 用 res 接收方法返回结果 90 Object res = method.invoke(intfaceImpl, invocation.getParams()); 91 // 打开 client 的输出流 92 oos = new ObjectOutputStream(client.getOutputStream()); 93 // 序列化, 向输出流中写入 RPC请求的结果 94 oos.writeObject(res); 95 oos.flush(); 96 } catch (Exception e) { 97 e.printStackTrace(); 98 } finally { // 关闭资源 99 CloseUtil.closeAll(ois, oos); 100 CloseUtil.closeAll(client); 101 } 102 } 103 } 104}
五. 测试类
Login类, RPC协议接口
1/** 2 * RPC协议接口, Client 与 Server端共同遵守 3 */ 4public interface Login { 5 /** 6 * 抽象方法 login(), 模拟用户登录传入两个String 类型的参数, 返回 String类型的结果 7 * 8 * @param username 9 * 用户名 10 * @param password 11 * 密码 12 * @return 返回登录结果 13 */ 14 public String login(String username, String password); 15}
LoginImpl类, Server 端 RPC协议接口( Login )的实现类
1/** 2 * Server端 RPC协议接口( Login )的实现类 3 */ 4public class LoginImpl implements Login { 5 /** 6 * 实现 login()方法, 模拟用户登录 7 * 8 * @param username 9 * 用户名 10 * @param password 11 * 密码 12 * @return hello 用户名 13 */ 14 @Override 15 public String login(String username, String password) { 16 return "hello " + username; 17 } 18}
ClientTest类, Client端测试类
1/** 2 * Client端测试类 3 */ 4public class ClientTest { 5 public static void main(String[] args) throws UnknownHostException, IOException { 6 // 获取一个 Client端的代理对象 proxy 7 Login proxy = (Login) RPC.getProxy(Login.class, "192.168.8.1", 8888); 8 // 调用 proxy 的 login() 方法, 返回值为 res 9 String res = proxy.login("rpc", "password"); 10 // 输出 res 11 System.out.println(res); 12 } 13}
ServerTest类, Server端测试类
1/** 2 * Server端测试类 3 */ 4public class ServerTest { 5 public static void main(String[] args) throws ClassNotFoundException, IOException { 6 // 获取 RPC 的 Server 端实例 server 7 RPCServer server = RPC.getRPCServer(Login.class, new LoginImpl(), 8888); 8 // 循环监听并接收 Client 端连接, 处理 RPC 请求, 向 Client 端返回结果 9 server.start(); 10 } 11}
运行 ServerTest, 控制台输出:
Starting Socket Handler for port 8888
运行 ClientTest, 控制台输出:
hello rpc
至此, 实现了基于 Proxy, Socket, IO 的简单版 RPC模型,
对于每一个 RPC请求, Server端都开启一个 Handler线程处理该请求,
在高并发情况下, Server端是扛不住的, 改用 NIO应该表现更好