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

作者: 西魏陶渊明
博客: https://blog.springlearn.cn/
天下代码一大抄, 抄来抄去有提高, 看你会抄不会抄!

@[toc]
一、目标
通过前面几篇文章的学习,我们已经完成了通信层的建立,如下我们可以通过下面代码实现服务间的通信。
- server() 使用Mojito构建一个服务端, 绑定 6666 端口。
- clientSync() 使用Mojito构建一个客户端,连接
connect("127.0.0.1", 6666), 发送一条任意请求。
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.1 执行体
什么是执行体,在java中我们执行一个方法。如何抽象成统一的执行模型呢? 这句话怎么理解呢? 如下代码示例。
1User user = new User();
2user.getName();// jay
3user.getAge();// 18抽象成统一的执行体可能就是下面这样
1Invoker invoker = new Invoker(User.class);
2invoker.invoker(user,"getName") // jay
3invoker.invoker(user,"getAge"); // 18那么我们为什么要这么做呢? 我们思考一个场景,客户端向服务端发送了一条指令,比如说是获取user的名字。那么我们服务端怎么来执行这个动作呢?
- 首先拿到User.class的实例对象。
- 然后通过反射,执行getName方法。
当有了一个统一的执行模型后,就可以不在使用反射。(因为反射已经被封装成一个统一的执行体了。) 这就是统一的执行体的好处,即屏蔽了反射的细节。提供了方便的简洁的执行API。
如下图我们设计了一个统一的接口 Invoker
- AbstractClusterInvoker 是客户端的抽象类
- AbstractInvoker 是服务端的抽象类
具体的代码细节见下文。

1.2 容错策略
什么容错策略, 客户端在进行远程调用的时候,当请求失败后, 如何处理? 不在处理? 还是重试? 还是换一个服务器继续执行? 这些统称为容错策略。我们这里也做了容错的设计,容错策略跟dubbo是一致的。提供以下四种能力。
1.3 Balance 负载均衡
什么是负载均衡这里就不赘述了, 可以参考小编这篇文章。负载均衡
这里我们直接写代码实现。
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(随机负载)
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(平均负载)
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)提出来的。这种算法是最好的。这里并没有提供实现。但是核心方法在前面文章中已经提供过了。感兴趣的可以看下前面文章。

二、执行体设计
前面我们讲了执行体的好处,这里开始实战。
首先定义接口
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
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 中,下面我们看抽象类的设计。
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方法
- 首先从
reflectorCacheMap看这个对象是否有方法的反射缓存,如果有就直接用,没有就生成。 可以看到生成方法的缓存对象也是比较简单的。new MethodReflectorCache(invocation.getInterface()) - 从参数执行体
Invocation拿到要执行的方法名,和参数对象数组。 - 然后从根据名称和参数找到
JvmInvoker执行方法。
2.2 客户端设计
前面说将了支持4种
- FailFastClusterInvoker 快速失败
- FailBackClusterInvoker 失败返回
- FailRetryClusterInvoker失败重试
- FailoverClusterInvoker 失败转移
2.2.1 FailFastClusterInvoker 快速失败
请求错误后,直接就返回失败,剩下的交给业务去处理。
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 失败返回
请求失败后,直接返回,与前者不同的是,会异步进行一次重试。
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次。
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 失败转移
请求失败后,会寻找下一个集群的节点进行处理,当所有节点都重试失败后,返回异常。
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 了吗?
