myrpc 学习笔记-001 实现简易版 rpc
项目地址 欢迎访问
笔记总览 myrpc 学习笔记-001 实现简易版 rpc myrpc 学习笔记-002 配置加载 myrpc 学习笔记-003 Mock 服务代理 myrpc 学习笔记-004 序列化实现和 SPI 机制 myrpc 学习笔记-005 注册中心 myrpc 学习笔记-006 自定义协议 myrpc 学习笔记-007 负载均衡 myrpc 学习笔记-008 重试机制 myrpc 学习笔记-009 容错机制 myrpc 学习笔记-0010 启动机制和注解驱动
了解rpc
RPC(Remote Procedure Call)远程过程调用,简单说就是:让你像调用本地函数一样,调用另一台机器上的函数,不用关心网络、通信、序列化这些底层细节。
简易实现
架构图 v1.0.0
- rpc 的简易版实现
- 网络通信是基于 vertx 实现的
- 序列化器使用的是 JDK 原生序列化器
- 注册中心使用的是本地注册 存的是服务名成和对应的实现类

一、模块划分(4个核心模块)
核心原则:模块职责单一,降低耦合,模拟真实微服务架构中“服务通信-接口定义-服务实现-服务调用”的完整链路,所有模块协同完成 RPC 远程调用的全流程。
1. 公共模块(common)
核心定位:统一接口定义中心,存放所有需要远程调用的接口,供服务提供者和服务消费者共同依赖,避免接口重复定义,保证双方接口一致性。
核心内容:
- 定义远程调用的业务接口(如 UserService、OrderService),仅包含方法声明,不包含实现(类似“契约”)。
- (可选)存放公共工具类、常量(如序列化方式、端口号、服务名称),供其他模块复用。
注意:该模块仅提供接口,不依赖其他任何模块,是整个 RPC 实现的“基础依赖”。
2. RPC 核心模块(rpc-core)
核心定位:RPC 通信核心,封装远程调用的底层细节(网络通信、序列化/反序列化、请求/响应封装),提供通用的远程调用能力,供服务提供者和消费者调用。
核心内容(必实现):
- 网络通信:基于 TCP 实现客户端(消费者侧)和服务端(提供者侧)的通信,负责请求发送、响应接收(可使用 Socket 编程)。
- 序列化/反序列化:将 Java 对象(请求参数、响应结果)转换为字节流(网络传输需要),反之亦然(可选:JDK 序列化、JSON、Protobuf 等)。
- 请求/响应封装:定义请求对象(包含接口名、方法名、参数类型、参数值)和响应对象(包含返回结果、异常信息),规范通信数据格式。
▼java复制代码package com.longlong.myrpc.model; import lombok.AllArgsConstructor; import lombok.Builder; import lombok.Data; import lombok.NoArgsConstructor; import java.io.Serializable; import static com.longlong.myrpc.constant.RpcConstant.DEFAULT_SERVICE_VERSION; /** * RPC 请求(客户端 -> 服务端) */ @Data @Builder @AllArgsConstructor @NoArgsConstructor public class RpcRequest implements Serializable { /** * 服务名称(接口全限定名) */ private String serviceName; /** * 方法名称 */ private String methodName; /** * 服务版本 */ private String serviceVersion = DEFAULT_SERVICE_VERSION; /** * 参数类型列表 */ private Class<?>[] parameterTypes; /** * 参数列表 */ private Object[] parameters; }
▼java复制代码package com.longlong.myrpc.model; import lombok.AllArgsConstructor; import lombok.Builder; import lombok.Data; import lombok.NoArgsConstructor; import java.io.Serializable; /** * RPC 响应(服务端 -> 客户端) */ @Data @Builder @AllArgsConstructor @NoArgsConstructor public class RpcResponse implements Serializable { /** * 响应数据(方法返回值) */ private Object data; /** * 响应数据类型 */ private Class<?> dataType; /** * 错误信息(调用失败时使用) */ private String message; /** * 异常信息 */ private Exception exception; }
- 服务注册与发现(简易版):这里先实现本地注册中心,使用map存储服务接口名和服务实现类对象。后续可扩展其他注册中心。
注意:该模块是 RPC 实现的核心,不依赖业务逻辑(不依赖公共模块的具体接口),仅提供通用通信能力。
3. 服务提供者模块(provider)
核心定位:接口实现与服务暴露,依赖公共模块的接口,实现具体业务逻辑,同时依赖 RPC 核心模块,将自身服务暴露出去,供消费者远程调用。
核心内容:
- 依赖 common 模块,实现公共模块中定义的远程接口(如 UserServiceImpl 实现 UserService),编写具体业务逻辑。
- 依赖 rpc-core 模块,启动服务端(监听指定端口),接收消费者的 RPC 请求。
- 将自身服务注册到注册中心(如通过 rpc-core 提供的注册方法,注册服务名称+自身地址)。
- 接收请求后,解析请求参数(反序列化),调用本地实现的接口方法,将结果封装为响应对象(序列化),返回给消费者。
4. 服务消费者模块(consumer)
核心定位:远程接口调用,依赖公共模块的接口,依赖 RPC 核心模块,通过代理方式(静态/动态)调用远程服务提供者的接口,无需关心底层通信。
核心内容:
- 依赖 common 模块,通过接口声明,调用远程方法(形式上和调用本地方法一致)。
- 依赖 rpc-core 模块,通过注册中心获取服务提供者的地址,启动客户端,发送 RPC 请求。
- 通过代理(静态/动态)封装 RPC 调用细节:调用接口方法时,自动触发 RPC 核心模块的请求发送、响应接收逻辑,最终返回远程方法的执行结果。
二、前期准备
这里在网络中传输的是 RpcRequest RpcResponse 对象,需要先进行序列化操作
JDK 原生序列化实现
▼java复制代码package com.longlong.myrpc.serializer; import java.io.IOException; /** * 序列化器接口 * 定义序列化、反序列化的标准方法 */ public interface Serializer { /** * 序列化 * 将 Java 对象转换为字节数组(用于网络传输) * * @param object 要序列化的对象 * @param <T> 对象类型 * @return 字节数组 * @throws IOException 序列化异常 */ <T> byte[] serialize(T object) throws IOException; /** * 反序列化 * 将字节数组转换为 Java 对象(用于接收数据后还原对象) * * @param bytes 字节数组 * @param type 目标对象类型 * @param <T> 目标类型 * @return 反序列化后的对象 * @throws IOException 反序列化异常 */ <T> T deserialize(byte[] bytes, Class<T> type) throws IOException; }
▼java复制代码package com.longlong.myrpc.serializer; import java.io.ByteArrayInputStream; import java.io.ByteArrayOutputStream; import java.io.IOException; import java.io.ObjectInputStream; import java.io.ObjectOutputStream; /** * JDK 原生序列化器实现 */ public class JdkSerializer implements Serializer { /** * 序列化:将对象转为字节数组 * @param object 要序列化的对象 * @return 序列化后的字节数组 * @throws IOException IO异常 */ @Override public <T> byte[] serialize(T object) throws IOException { // 创建字节数组输出流 ByteArrayOutputStream outputStream = new ByteArrayOutputStream(); // 创建对象输出流,用于写入对象 ObjectOutputStream objectOutputStream = new ObjectOutputStream(outputStream); // 写入对象 objectOutputStream.writeObject(object); // 关闭流 objectOutputStream.close(); // 返回字节数组 return outputStream.toByteArray(); } /** * 反序列化:将字节数组转为对象 * @param bytes 字节数组 * @param type 目标对象类型 * @return 反序列化后的对象 * @throws IOException IO异常 */ @Override public <T> T deserialize(byte[] bytes, Class<T> type) throws IOException { ByteArrayInputStream inputStream = new ByteArrayInputStream(bytes); ObjectInputStream objectInputStream = new ObjectInputStream(inputStream); try { // 读取对象并强转为指定类型 return (T) objectInputStream.readObject(); } catch (ClassNotFoundException e) { // 捕获类找不到异常,包装为运行时异常抛出 throw new IOException("反序列化失败,未找到目标类", e); } finally { // 关闭流 objectInputStream.close(); } } }
实现网络通信
实现基于 Vert.x 的 RPC 服务端 + 客户端网络通信
服务端
▼java复制代码/** * HTTP 服务器接口 * 定义统一的启动服务方法 */ public interface HttpServer { /** * 启动服务器 * @param port 端口号 */ void doStart(int port); }
▼java复制代码/** * Vert.x HTTP 服务器实现 * 用于 RPC 框架接收客户端请求、处理远程调用 */ public class VertxHttpServer implements HttpServer { /** * 启动 HTTP 服务 * @param port 监听端口 */ @Override public void doStart(int port) { // 1. 获取 Vertx 实例(Vert.x 核心对象) Vertx vertx = Vertx.vertx(); // 2. 创建 HTTP 服务器 HttpServer server = vertx.createHttpServer(); // 3. 设置请求处理器(所有客户端请求都会进入这里) server.requestHandler(new HttpServerHandler()); // 4. 监听端口,启动服务 server.listen(port, result -> { if (result.succeeded()) { System.out.println("RPC 服务端启动成功,端口:" + port); } else { System.err.println("RPC 服务端启动失败:" + result.cause()); } }); } // 测试启动服务端 public static void main(String[] args) { new VertxHttpServer().doStart(8080); } }
▼java复制代码/** * HTTP 请求处理器 * 接收 RPC 请求,调用对应服务,并返回响应 */ public class HttpServerHandler implements Handler<HttpServerRequest> { /** * 处理请求 * @param request Vert.x HTTP 请求 */ public void handle(HttpServerRequest request) { // 1. 创建序列化器 Serializer serializer = new JdkSerializer(); // 异步处理请求体 request.bodyHandler(body -> { RpcResponse rpcResponse = new RpcResponse(); try { // 2. 反序列化:字节数组 -> RpcRequest byte[] bytes = body.getBytes(); RpcRequest rpcRequest = serializer.deserialize(bytes, RpcRequest.class); // 3. 获取服务实现类 Class<?> implClass = LocalRegistry.getService(rpcRequest.getServiceName()); if (implClass == null) { throw new RuntimeException("服务未注册: " + rpcRequest.getServiceName()); } // 4. 反射调用方法 Method method = implClass.getMethod( rpcRequest.getMethodName(), rpcRequest.getParameterTypes() ); Object result = method.invoke(implClass.newInstance(), rpcRequest.getParameters()); // 5. 封装响应 rpcResponse.setData(result); rpcResponse.setDataType(result.getClass()); } catch (Exception e) { // 异常处理 rpcResponse.setMessage(e.getMessage()); rpcResponse.setException(e); e.printStackTrace(); } // 6. 序列化响应并返回 try { byte[] responseBytes = serializer.serialize(rpcResponse); HttpServerResponse response = request.response(); response.putHeader("Content-Type", "application/octet-stream"); response.end(io.vertx.core.buffer.Buffer.buffer(responseBytes)); } catch (IOException e) { e.printStackTrace(); request.response().end(); } }); } }
客户端
▼java复制代码public class VertxHttpClient { /** * 定义 RPC 服务端地址 */ private static final String SERVICE_ADDRESS = "http://localhost:8080"; public RpcResponse sendRequest(RpcRequest rpcRequest) { // 1. 创建序列化器 Serializer serializer = new JdkSerializer(); try { // 3. 序列化请求对象 byte[] requestBytes = serializer.serialize(rpcRequest); // 4. 发送 HTTP 请求给服务端 try (HttpResponse httpResponse = HttpRequest.post(SERVICE_ADDRESS) .body(requestBytes) .execute()) { byte[] responseBytes = httpResponse.bodyBytes(); // 5. 反序列化响应 RpcResponse res = serializer.deserialize(responseBytes, RpcResponse.class); return res; } } catch (Exception e) { e.printStackTrace(); } return null; } }
服务注册与发现
本地注册中心
▼java复制代码package com.longlong.myrpc.registry; import java.util.Map; import java.util.concurrent.ConcurrentHashMap; /** * 本地注册中心 * 存储:服务接口名 -> 服务实现类对象 */ public class LocalRegistry { /** * 注册信息存储 * 使用线程安全的 ConcurrentHashMap,保证高并发场景下安全 * key:服务名称(接口全限定名,如:com.longlong.myrpc.service.UserService) * value:服务实现类的 Class 对象(用于反射创建实例) */ private static final Map<String, Class<?>> SERVICE_MAP = new ConcurrentHashMap<>(); /** * 注册服务 * * @param serviceName 服务名称(接口名) * @param implClass 服务实现类 */ public static void register(String serviceName, Class<?> implClass) { SERVICE_MAP.put(serviceName, implClass); } /** * 获取服务 * * @param serviceName 服务名称 * @return 服务实现类 */ public static Class<?> getService(String serviceName) { return SERVICE_MAP.get(serviceName); } /** * 删除服务 * * @param serviceName 服务名称 */ public static void removeService(String serviceName) { SERVICE_MAP.remove(serviceName); } }
服务提供者注册服务
▼java复制代码public class EasyProviderExample { public static void main(String[] args) { //注册服务到本地 LocalRegistry.register(UserService.class.getName(), UserServiceImpl.class); // 启动 web 服务 HttpServer httpServer = new VertxHttpServer(); httpServer.doStart(RpcApplication.getRpcConfig().getServerPort()); } }
三、实现简易 rpc :两种代理方式(核心重点)
代理的核心目的:屏蔽 RPC 底层通信细节,让消费者调用远程接口时,完全像调用本地方法一样,无需手动处理网络请求、序列化等操作。代理类负责将“接口调用”转换为“RPC 远程调用”。
1. 静态代理实现
核心原理
针对每个远程接口,手动编写一个代理类(如 UserServiceProxy),代理类实现该接口,在接口方法中封装 RPC 调用逻辑(请求发送、响应接收),消费者直接调用代理类的方法,间接完成远程调用。
实现步骤
- 在 common 模块定义 UserService 接口
▼java复制代码public interface UserService { /** * 获取用户信息 */ User getUser(User user); }
- 在 consumer 模块编写 UserServiceProxy 类,实现 UserService 接口。
- 在代理类的 getUserById 方法中,调用 rpc-core 模块的客户端能力:
- 封装请求对象(接口名:UserService,方法名:getUserById,参数:id)。
- 序列化请求对象,通过 Socket 发送到服务提供者地址。
- 接收服务提供者返回的响应对象,反序列化后,返回结果。
▼java复制代码/** * UserService 静态代理 * 用于客户端:封装网络请求,调用远程服务 */ public class UserServiceProxy implements UserService { /** * 远程服务地址 */ private static final String URL = "http://localhost:8080"; @Override public User getUser(User user) { // 1. 创建序列化器 Serializer serializer = new JdkSerializer(); // 2. 构造 RPC 请求 RpcRequest rpcRequest = RpcRequest.builder() .serviceName(UserService.class.getName()) .methodName("getUser") .parameterTypes(new Class[]{User.class}) .parameters(new Object[]{user}) .build(); try { // 3. 将 RPC 请求 序列化 byte[] bodyBytes = serializer.serialize(rpcRequest); // 4. 发送 HTTP 请求给服务端 byte[] resultBytes = HttpRequest.post(URL) .body(bodyBytes) .execute() .bodyBytes(); // 5. 将收到的响应信息反序列化 得到 RPC 响应 RpcResponse rpcResponse = serializer.deserialize(resultBytes, RpcResponse.class); // 6. 从 RPC 响应拿到方法返回值 return (User) rpcResponse.getData(); } catch (Exception e) { e.printStackTrace(); } return null; } }
- 消费者使用时,直接创建 UserServiceProxy 实例,调用其方法即可
▼java复制代码public class EasyConsumerExample { public static void main(String[] args) { // 静态代理 UserService userService = new UserServiceProxy(); User user = new User(); user.setName("zhangsan123"); User newUser = userService.getUser(user); } }
优点与缺点
- 优点:实现简单,无需依赖额外框架,代码直观,适合入门理解 RPC 原理。
- 缺点:冗余度高,每个远程接口都需要编写对应的代理类;维护成本高,接口修改后,代理类也需要同步修改。
2. 动态代理实现
核心原理
不手动编写代理类,而是在程序运行时,通过动态代理框架(如 JDK 动态代理、CGLIB)动态生成代理对象。代理对象会拦截接口的所有方法调用,统一封装 RPC 调用逻辑,适用于所有远程接口,无需重复编写代码。
常用方式:JDK 动态代理(基于接口,推荐,因为我们的远程接口都在 common 模块定义,刚好符合 JDK 动态代理的要求)。
实现步骤
- 在 consumer 模块编写动态代理工厂类(RpcProxyFactory),提供一个静态方法,用于生成远程接口的代理对象。
- 代理工厂类中,实现 InvocationHandler 接口,重写 invoke 方法(核心逻辑):
- invoke 方法会被自动调用(当消费者调用代理对象的接口方法时)。
- 在 invoke 方法中,获取当前调用的接口名、方法名、参数,封装为请求对象。
- 调用 rpc-core 模块的客户端能力,发送请求、接收响应,反序列化后返回结果。
▼java复制代码public class ServiceProxyFactory { /** * 私有构造方法 * <p> * 1. 禁止外部通过 new 关键字实例化当前工厂类 * 2. 工厂类只提供静态方法,无需创建对象,符合工具类设计规范 * </p> */ private ServiceProxyFactory() { } /** * 获取远程服务的动态代理对象 * <p> * 核心方法:根据传入的服务接口类,生成一个实现了该接口的代理实例 * 所有对接口方法的调用,都会被转发到 ServiceProxy 中的 invoke 方法处理 * </p> * * @param <T> 泛型类型,代表服务接口类型(如:UserService、OrderService) * @param serviceClass 服务接口的 Class 对象(必须是接口,不能是普通类) * @return 生成的动态代理对象,实现了指定的服务接口 */ private static <T> T getProxy(Class<T> serviceClass) { // 调用 JDK 原生动态代理,创建代理对象 return (T) Proxy.newProxyInstance( // 类加载器:使用服务接口自身的类加载器加载代理类 serviceClass.getClassLoader(), // 要代理的接口数组:JDK 动态代理必须基于接口,这里只代理当前传入的服务接口 new Class[]{serviceClass}, // 调用处理器:所有接口方法的调用都会被这个处理器拦截处理 // ServiceProxy 中实现了远程调用、序列化、网络请求等核心逻辑 new ServiceProxy() ); } }
▼java复制代码/** * 服务代理(JDK 动态代理) * 作用:统一为所有服务接口生成代理对象,发送 RPC 请求 */ public class ServiceProxy implements InvocationHandler { /** * 定义 RPC 服务端地址 */ private static final String SERVICE_ADDRESS = "http://localhost:8080"; /** * 动态代理核心方法:拦截所有接口方法调用 */ @Override public Object invoke(Object proxy, Method method, Object[] args) throws Throwable { // 1. 创建序列化器 Serializer serializer = new JdkSerializer(); // 2. 构造 RPC 请求 RpcRequest rpcRequest = RpcRequest.builder() .serviceName(method.getDeclaringClass().getName()) .methodName(method.getName()) .parameterTypes(method.getParameterTypes()) .parameters(args) .build(); try { // 3. 序列化请求对象 byte[] requestBytes = serializer.serialize(rpcRequest); // 4. 发送 HTTP 请求给服务端 try (HttpResponse httpResponse = HttpRequest.post(SERVICE_ADDRESS) .body(requestBytes) .execute()) { byte[] responseBytes = httpResponse.bodyBytes(); // 5. 反序列化响应 RpcResponse rpcResponse = serializer.deserialize(responseBytes, RpcResponse.class); // 6. 返回结果 return rpcResponse.getData(); } } catch (Exception e) { e.printStackTrace(); } return null; } }
- 消费者使用时,通过代理工厂获取代理对象(如 UserService proxy = RpcProxyFactory.getProxy(UserService.class)),直接调用接口方法即可。
▼java复制代码public class EasyConsumerExample { public static void main(String[] args) { // 动态代理 UserService userService = ServiceProxyFactory.getProxy(UserService.class); User user = new User(); user.setName("zhangsan123"); User newUser = userService.getUser(user); } }
优点与缺点
- 优点:无冗余代码,一个代理工厂可生成所有远程接口的代理对象;维护成本低,接口修改后,无需修改代理逻辑。
- 缺点:实现难度略高于静态代理,需要理解 JDK 动态代理的原理;仅支持基于接口的代理(若需代理类,可使用 CGLIB)。
四、流程分析
- 模块依赖关系(重要):
- provider 依赖 → common + rpc-core
- consumer 依赖 → common + rpc-core
- rpc-core 依赖 → 无(独立核心)
- common 依赖 → 无(基础接口)
- 执行流程:
- 启动 provider(注册服务、监听端口)→ 启动 consumer(获取代理对象)→ 消费者调用代理方法 → 代理类触发 RPC 调用 → provider 接收请求并执行 → 结果返回给消费者。
