系列

第06篇:手写JavaRPC框架之执行层实战

2026-07-314 min read手写RPC框架

作者: 西魏陶渊明

博客: https://blog.springlearn.cn/

天下代码一大抄, 抄来抄去有提高, 看你会抄不会抄!

@[toc]

一、目标

通过前面几篇文章的学习,我们已经完成了通信层的建立,如下我们可以通过下面代码实现服务间的通信。

  • server() 使用Mojito构建一个服务端, 绑定 6666 端口。
  • clientSync() 使用Mojito构建一个客户端,连接 connect("127.0.0.1", 6666), 发送一条任意请求。
java
1 /** 2 * @author liuxin 3 * 个人博客:https://java.springlearn.cn 4 * 公众号:西魏陶渊明 {关注获取学习源码} 5 * 2022/8/11 23:12 6 */ 7 @Test 8 @DisplayName("构建服务端【阻塞方式】") 9 public void server() throws Exception { 10 Mojito.server(RpcRequest.class, RpcResponse.class) 11 // 业务层,读取请求对象,返回结果 12 .businessHandler((channelContext, request) -> new RpcResponse()) 13 .create().start(6666); 14 } 15 16 @Test 17 @DisplayName("构建客户端【同步方式】") 18 public void clientSync() throws Exception{ 19 Client<RpcRequest, RpcResponse> client = Mojito.client(RpcRequest.class, RpcResponse.class) 20 .connect("127.0.0.1", 6666); 21 System.out.println(client.send(new RpcRequest())); 22 }

实现通信只是第一步,下面我们要基于我们的通信层来设计我们的RPC框架, 第05篇:手写JavaRPC框架之执行层思路已经将RPC接口的设计给讲清楚了,那么本篇文章我们即开始着手开始实现吧。

本篇文章主要有以下几个知识点

  1. 执行体: 用于服务端执行远程请求方法, 用于客户端实现服务端接口的实例化。
  2. 容错策略: 用于客户端在发起远程请求失败后的处理逻辑。
  3. 负载均衡: 用于客户端均匀的向不同服务器发送请求,以实现保护服务端的能力。

1.1 执行体

什么是执行体,在java中我们执行一个方法。如何抽象成统一的执行模型呢? 这句话怎么理解呢? 如下代码示例。

java
1User user = new User(); 2user.getName();// jay 3user.getAge();// 18

抽象成统一的执行体可能就是下面这样

java
1Invoker invoker = new Invoker(User.class); 2invoker.invoker(user,"getName") // jay 3invoker.invoker(user,"getAge"); // 18

那么我们为什么要这么做呢? 我们思考一个场景,客户端向服务端发送了一条指令,比如说是获取user的名字。那么我们服务端怎么来执行这个动作呢?

  1. 首先拿到User.class的实例对象。
  2. 然后通过反射,执行getName方法。

当有了一个统一的执行模型后,就可以不在使用反射。(因为反射已经被封装成一个统一的执行体了。) 这就是统一的执行体的好处,即屏蔽了反射的细节。提供了方便的简洁的执行API。

如下图我们设计了一个统一的接口 Invoker

  • AbstractClusterInvoker 是客户端的抽象类
  • AbstractInvoker 是服务端的抽象类 具体的代码细节见下文。

1.2 容错策略

什么容错策略, 客户端在进行远程调用的时候,当请求失败后, 如何处理? 不在处理? 还是重试? 还是换一个服务器继续执行? 这些统称为容错策略。我们这里也做了容错的设计,容错策略跟dubbo是一致的。提供以下四种能力。

1.3 Balance 负载均衡

什么是负载均衡这里就不赘述了, 可以参考小编这篇文章。负载均衡

这里我们直接写代码实现。

java
1/** 2 * 负载均衡 3 * 4 * @author liuxin 5 * 2022/8/28 20:07 6 */ 7public interface LoadBalance { 8 9 /** 10 * 从众多的执行体中选择一个执行 11 * 12 * @param invokers 执行体集合 13 * @param invocation 执行参数 14 * @param <T> 泛型 15 * @return Invoker 16 */ 17 <T> Invoker<T> routeBalance(List<Invoker<T>> invokers, Invocation invocation); 18}

1.3.1 RandomLoadBalance(随机负载)

java
1public class RandomLoadBalance implements LoadBalance { 2 3 private final static Random random = new Random(); 4 5 @Override 6 public <T> Invoker<T> routeBalance(List<Invoker<T>> invokers, Invocation invocation) { 7 return invokers.get(random.nextInt(invokers.size())); 8 } 9}

1.3.2 AverageLoadBalance(平均负载)

java
1public class AverageLoadBalance implements LoadBalance { 2 3 4 private static final AtomicInteger atomicInteger = new AtomicInteger(); 5 6 7 @Override 8 public <T> Invoker<T> routeBalance(List<Invoker<T>> invokers, Invocation invocation) { 9 int count = atomicInteger.get(); 10 if (count == Integer.MAX_VALUE) { 11 atomicInteger.set(0); 12 } 13 atomicInteger.incrementAndGet(); 14 return invokers.get(count % invokers.size()); 15 } 16}

1.3.3 平滑加权算法

主要解决上面那种不平滑的方案。这种方案是由nginx (opens new window)提出来的。这种算法是最好的。这里并没有提供实现。但是核心方法在前面文章中已经提供过了。感兴趣的可以看下前面文章。

二、执行体设计

前面我们讲了执行体的好处,这里开始实战。

首先定义接口

java
1public interface Invoker<T> { 2 3 /** 4 * 服务接口 5 * 6 * @return Class<T> 7 */ 8 Class<T> getInterface(); 9 10 /** 11 * 远程接口调用 12 * 13 * @param invocation 接口请求参数 14 * @return Result 统一的执行结果 15 * @throws RpcException Rpc异常 16 */ 17 Result invoke(Invocation invocation) throws RpcException; 18 19 /** 20 * 执行体封装的原始对象 21 * @return 原始对象 22 */ 23 default T getSource() { 24 return null; 25 } 26}

2.1 服务端设计

服务端要将本地的对象实例,转换成执行体。 对于服务端来说,只是将本地的java对象转换成 Invoker,所以它的实现类名叫 LocalInvoker

java
1/** 2 * 本地对象Invoker封装 3 * @author liuxin 4 * 个人博客:https://java.springlearn.cn 5 * 公众号:西魏陶渊明 {关注获取学习源码} 6 * 2022/8/28 23:25 7 */ 8public class LocalInvoker<T> extends AbstractInvoker<T> { 9 10 /** 11 * 原始对象 12 */ 13 private Object source; 14 15 private LocalInvoker() { 16 } 17 18 public LocalInvoker(Object source) { 19 this.source = source; 20 } 21 22 @Override 23 public T getSource() { 24 return (T)source; 25 } 26} 27

可以看到实现是非常简单的,只需要通过构造将对象包装,即可转换成一个 Invoker, 之所以这样是因为主要的逻辑都在抽象类 AbstractInvoker 中,下面我们看抽象类的设计。

java
1public abstract class AbstractInvoker<T> implements Invoker<T> { 2 3 /** 4 * 方法缓存 5 */ 6 private final Map<Class<?>, MethodReflectorCache> reflectorCacheMap = new ConcurrentHashMap<>(); 7 8 @Override 9 @SuppressWarnings("unchecked") 10 public Class<T> getInterface() { 11 return (Class<T>) getSource().getClass(); 12 } 13 14 @Override 15 public Result invoke(Invocation invocation) throws RpcException { 16 MethodReflectorCache reflectorCache = reflectorCacheMap.get(invocation.getInterface()); 17 if (Objects.isNull(reflectorCache)) { 18 reflectorCache = new MethodReflectorCache(invocation.getInterface()); 19 reflectorCacheMap.put(invocation.getInterface(), reflectorCache); 20 } 21 String methodName = invocation.getMethodName(); 22 Object[] arguments = invocation.getArguments(); 23 try { 24 JvmInvoker method = reflectorCache.findMethod(methodName, arguments); 25 Object invoke = method.invoke(getSource(), arguments); 26 return new RpcResult(invoke); 27 } catch (NoSuchMethodException | InvocationTargetException | IllegalAccessException e) { 28 return new RpcResult(e); 29 } 30 } 31 32 /** 33 * 具体执行类 34 * 35 * @return Object 36 */ 37 public abstract T getSource(); 38} 39
  • reflectorCacheMap 用来保存不同对象的方法缓存
  • MethodReflectorCache 方法反射对象缓存
  • Invocation 方法的参数对象

然后我们看invoke方法

  1. 首先从 reflectorCacheMap 看这个对象是否有方法的反射缓存,如果有就直接用,没有就生成。 可以看到生成方法的缓存对象也是比较简单的。new MethodReflectorCache(invocation.getInterface())
  2. 从参数执行体Invocation拿到要执行的方法名,和参数对象数组。
  3. 然后从根据名称和参数找到 JvmInvoker 执行方法。

2.2 客户端设计

前面说将了支持4种

  • FailFastClusterInvoker 快速失败
  • FailBackClusterInvoker 失败返回
  • FailRetryClusterInvoker失败重试
  • FailoverClusterInvoker 失败转移

2.2.1 FailFastClusterInvoker 快速失败

请求错误后,直接就返回失败,剩下的交给业务去处理。

java
1 public Result doDirectInvoker(Invoker<T> invoker, Invocation invocation) { 2 Result invoke; 3 try { 4 invoke = doInvoke(invoker,invocation); 5 } catch (Throwable t) { 6 // 失败后直接抛异常 7 invoke = new RpcResult(t); 8 } 9 return invoke; 10 }

2.2.2 FailBackClusterInvoker 失败返回

请求失败后,直接返回,与前者不同的是,会异步进行一次重试。

java
1 public Result doDirectInvoker(Invoker<T> invoker, Invocation invocation) { 2 Result invoke; 3 try { 4 invoke = doInvoke(invoker, invocation); 5 } catch (Throwable t) { 6 // 添加到队列中 7 addRetryTask(invoker, invocation); 8 // 如果失败,构建一个FailBackResult 9 invoke = new FailBackResult(); 10 } 11 return invoke; 12 }

2.2.3 FailRetryClusterInvoker 失败重试

请求失败后,会继续针对这个集群的节点进行重试。第一次重试1s,第二次重试2s,第三次重试3s。一共重试3次。

java
1 public Result doDirectInvoker(Invoker<T> invoker, Invocation invocation) { 2 Result invoke; 3 try { 4 // 基于重试策略 5 invoke = invokeRetry.call(() -> new RetryInvoker(invoker).invoke(invocation)); 6 } catch (Throwable t) { 7 invoke = new RpcResult(t); 8 } 9 return invoke; 10 }

2.2.4 FailoverClusterInvoker 失败转移

请求失败后,会寻找下一个集群的节点进行处理,当所有节点都重试失败后,返回异常。

java
1 @Override 2 public Result doDirectInvoker(Invoker<T> invoker, Invocation invocation) { 3 Result invoke; 4 try { 5 // 基于重试策略 6 invoke = doInvoke(invoker, invocation); 7 } catch (Throwable t) { 8 invoke = failover(invoker, invocation); 9 } 10 return invoke; 11 } 12 13 /** 14 * 失败转移 15 * 16 * @param invoker 当前执行器 17 * @param invocation 执行参数 18 * @return Result 执行结果 19 */ 20 public Result failover(Invoker<T> invoker, Invocation invocation) { 21 List<Invoker<T>> invokers = getInvokers(invocation); 22 Throwable lastThrowable = null; 23 for (int i = 0; i < invokers.size(); i++) { 24 Invoker<T> overInvoker = invokers.get(i); 25 if (!overInvoker.equals(invoker)) { 26 try { 27 // 成功了直接返回,失败了,继续轮训 28 return overInvoker.invoke(invocation); 29 } catch (Throwable t) { 30 // 最后一个异常 31 lastThrowable = t; 32 } 33 } 34 } 35 return new RpcResult(lastThrowable); 36 }

那么你准备好跟我一起 Coding 了吗?