调用接口时,RPC 背后发生了什么?探究 RPC 底层原理
前言
假设项目中存在这样一个接口:
public interface UserService {
User findById(long id);
}
正常情况下,我们需要创建它的实现类:
UserService userService = new UserServiceImpl();
User user = userService.findById(1001L);
方法调用发生在当前 JVM 中:
调用方
↓
UserServiceImpl.findById()
↓
返回 User
但现在,真正的 UserServiceImpl 并不在当前进程中。
它可能运行在另一台服务器:
客户端 JVM 服务端 JVM
UserService UserServiceImpl
客户端仍然希望这样调用:
User user = userService.findById(1001L);
这时问题就出现了:
接口没有本地实现类,
userService对象从哪里来?
调用
findById()时,怎样把方法名和参数发送到另一台机器?
TCP 只传输字节,服务端怎样知道这些字节代表哪个方法?
多个线程共用一条连接时,多个请求粘在一起怎么办?
服务端响应顺序变化后,客户端怎样知道哪个响应属于哪个线程?
writeAndFlush()成功,是否代表远程方法已经执行成功?
当前线程要等待远程结果,应该使用
CountDownLatch还是 Future?
一条连接能不能同时承载多个请求?
连接池中的连接应该借出后独占,还是允许并发复用?
服务端收到接口名和方法名后,怎样找到真正的实现类?
这些问题共同构成了 RPC 的核心。
RPC 看上去像一次普通的 Java 方法调用:
User user = userService.findById(1001L);
底层实际经历的是:
方法调用
↓
动态代理
↓
请求对象
↓
序列化
↓
协议编码
↓
网络发送
↓
服务端解码
↓
反射调用
↓
响应编码
↓
请求匹配
↓
返回结果
RPC 的本质不是“远程调用和本地调用完全一样”。
而是:
在接口调用的外表之下,隐藏一次完整的网络通信。
一、RPC 是什么?
RPC 是 Remote Procedure Call,也就是远程过程调用。
从使用者角度看:
User user = userService.findById(1001L);
从底层看:
客户端 服务端
findById(1001)
│
│ 请求:
│ 服务名、方法名、参数
├──────────────────────────────────→
│
│ 查找实现类
│ 反射调用方法
│ 得到 User
│
│ 响应:
│ User 或异常
│←──────────────────────────────────
│
返回 User
一个最小 RPC 系统至少需要解决六件事:
1. 调用拦截
2. 消息表达
3. 序列化
4. 网络传输
5. 请求匹配
6. 服务执行
继续展开:
调用拦截
动态代理
消息表达
请求对象和响应对象
序列化
对象与字节互相转换
网络传输
TCP、Netty、连接管理
请求匹配
requestId、Future
服务执行
服务注册、方法查找、反射调用
二、代理从哪里来?
客户端只有接口:
public interface UserService {
User findById(long id);
}
接口不能直接实例化:
new UserService(); // 编译错误
但 RPC 客户端并不需要真正的业务实现类。
它需要的是一个代理对象。
1. 创建代理
JDK 动态代理可以在运行时创建一个实现指定接口的对象:
@SuppressWarnings("unchecked")
public static <T> T createProxy(
Class<T> serviceInterface,
RpcClient rpcClient
) {
Object proxy = Proxy.newProxyInstance(
serviceInterface.getClassLoader(),
new Class<?>[]{serviceInterface},
(proxyObject, method, arguments) -> {
return rpcClient.invoke(
serviceInterface,
method,
arguments
);
}
);
return (T) proxy;
}
使用方式:
UserService userService =
createProxy(UserService.class, rpcClient);
虽然没有 UserServiceImpl,但返回的对象确实实现了:
UserService
JDK 的 Proxy.newProxyInstance() 会创建一个实现指定接口的代理实例。调用代理对象的方法时,这次调用会被编码并分派给关联的 InvocationHandler.invoke(),其中可以取得代理对象、Method 和参数数组。
2. 调用被拦截
执行:
userService.findById(1001L);
不会进入本地业务实现。
它会进入:
(proxyObject, method, arguments) -> {
// 构造远程请求
}
此时可以取得:
接口:UserService
方法:findById
参数:[1001]
返回类型:User
参数类型:[long]
然后把它们转换成一条 RPC 请求。
3. Object 方法
动态代理还有一个容易忽略的问题。
下面这些方法也可能进入 InvocationHandler:
proxy.toString();
proxy.hashCode();
proxy.equals(other);
JDK 代理对 Object 中的 hashCode()、equals() 和 toString() 也会进行分派,因此代理处理器通常需要单独处理它们,不能把每次 toString() 都发送到远程服务器。
例如:
private static Object invokeObjectMethod(
Object proxy,
Method method,
Object[] args
) {
return switch (method.getName()) {
case "toString" ->
"RpcProxy@" + System.identityHashCode(proxy);
case "hashCode" ->
System.identityHashCode(proxy);
case "equals" ->
proxy == args[0];
default ->
throw new UnsupportedOperationException(
method.toString()
);
};
}
完整判断:
if (method.getDeclaringClass() == Object.class) {
return invokeObjectMethod(
proxyObject,
method,
arguments
);
}
三、请求装什么?
调用信息不能只发送方法名。
假设存在重载:
User findById(long id);
User findById(String username);
如果请求中只有:
methodName = findById
服务端无法确定应该调用哪一个方法。
一个 RPC 请求至少需要:
requestId
serviceName
methodName
parameterTypes
arguments
可以定义为:
public record RpcRequest(
long requestId,
String serviceName,
String methodName,
String[] parameterTypeNames,
Object[] arguments
) {
}
构造请求:
private RpcRequest createRequest(
long requestId,
Class<?> serviceInterface,
Method method,
Object[] arguments
) {
String[] parameterTypeNames =
Arrays.stream(method.getParameterTypes())
.map(Class::getName)
.toArray(String[]::new);
return new RpcRequest(
requestId,
serviceInterface.getName(),
method.getName(),
parameterTypeNames,
arguments == null
? new Object[0]
: arguments
);
}
1. 为什么要有服务名?
客户端和服务端不共享对象实例。
服务端需要根据:
com.example.UserService
找到:
UserServiceImpl
因此需要维护服务注册表:
Map<String, Object> services =
new ConcurrentHashMap<>();
注册:
services.put(
UserService.class.getName(),
new UserServiceImpl()
);
查找:
Object service =
services.get(request.serviceName());
真实 RPC 框架中的服务标识通常还可能包含:
接口名
版本
分组
命名空间
例如:
com.example.UserService:1.0:production
否则两个版本的同名服务无法区分。
2. 为什么要有参数类型?
参数对象本身不一定足以定位方法。
例如传入:
null
时,无法从参数对象获取类型。
参数还可能是:
接口类型
父类类型
基本类型
包装类型
所以请求中应该明确携带方法签名。
3. 为什么要有 requestId?
假设十个线程共用一条连接:
线程 A → 请求 101
线程 B → 请求 102
线程 C → 请求 103
服务端可能按照下面的顺序完成:
响应 102
响应 103
响应 101
如果没有请求 ID,客户端无法判断每个响应应该交给哪个线程。
因此,每条请求都需要一个关联标识:
requestId
四、响应装什么?
响应不能只有返回值。
远程调用可能出现:
正常返回
业务异常
服务不存在
方法不存在
反序列化失败
服务端超时
系统异常
可以定义:
public record RpcResponse(
long requestId,
boolean success,
Object result,
String errorType,
String errorMessage
) {
public static RpcResponse success(
long requestId,
Object result
) {
return new RpcResponse(
requestId,
true,
result,
null,
null
);
}
public static RpcResponse failure(
long requestId,
Throwable throwable
) {
return new RpcResponse(
requestId,
false,
null,
throwable.getClass().getName(),
throwable.getMessage()
);
}
}
请求和响应必须携带同一个 ID:
请求:
requestId = 101
响应:
requestId = 101
这样客户端才能完成关联。
五、对象如何变成字节?
TCP 不能直接发送:
RpcRequest
它只能发送字节。
因此需要序列化:
RpcRequest
↓
byte[]
接收端再反序列化:
byte[]
↓
RpcRequest
1. Java 原生序列化
学习版本中,可以使用:
public static byte[] serialize(Object value)
throws IOException {
ByteArrayOutputStream output =
new ByteArrayOutputStream();
try (ObjectOutputStream objectOutput =
new ObjectOutputStream(output)) {
objectOutput.writeObject(value);
}
return output.toByteArray();
}
反序列化:
public static Object deserialize(byte[] bytes)
throws IOException, ClassNotFoundException {
try (ObjectInputStream input =
new ObjectInputStream(
new ByteArrayInputStream(bytes)
)) {
return input.readObject();
}
}
ObjectOutputStream 可以把实现了 Serializable 的对象图写入 OutputStream,ObjectInputStream 可以在另一进程中重建对象。序列化内容不仅包含字段值,还可能包含类描述、对象引用关系和流格式信息。
2. 头部不能这样写
笔记中有一种思路:
把 Header 对象序列化
然后假设 Header 长度固定为 110 字节
这是不可靠的。
例如:
byte[] headerBytes = serialize(header);
序列化后的长度可能受到以下因素影响:
类描述信息
字段变化
类名变化
序列化版本
对象引用
字符串内容
流头信息
ObjectOutputStream 构造时会写入流魔数和版本信息,写对象时还会维护对象句柄及类描述。因此,Java 对象序列化结果不适合作为“天然固定长度的协议头”。
协议头应该使用确定的二进制字段编码。
3. 安全问题
Java 原生反序列化不应该直接处理不可信网络数据。
Oracle 的序列化过滤文档指出,反序列化过滤器可以限制允许反序列化的类、数组长度、对象图深度、引用数量和数据量;过滤机制默认不会自动生效,需要显式配置。
学习项目可以使用 Java 序列化理解流程。
生产系统更适合选择具有明确数据模型的序列化方案,并对:
类型白名单
消息大小
嵌套深度
字段范围
进行限制。
六、TCP 没有消息边界
假设客户端连续发送两个请求:
请求 A:100 字节
请求 B:200 字节
TCP 不保证服务端按照两次读取返回:
第一次 read:100 字节
第二次 read:200 字节
服务端可能读到:
情况一:
第一次读 300 字节
情况二:
第一次读 40 字节
第二次读 260 字节
情况三:
第一次读 150 字节
第二次读 150 字节
这就是经常说的:
半包
粘包
它们并不是 TCP 出错。
TCP 本来就是一个连续字节流。
应用协议必须自己定义消息边界。
七、协议怎么设计?
一个简单的 RPC 协议可以分成:
Header
Body
结构如下:
┌─────────┬─────────┬─────────┬───────────┬────────────┬────────────┐
│ Magic │ Version │ Flags │ Serializer│ Request ID │ Body Length│
│ 4 bytes │ 1 byte │ 1 byte │ 1 byte │ 8 bytes │ 4 bytes │
└─────────┴─────────┴─────────┴───────────┴────────────┴────────────┘
┌───────────────────────────────────────────────────────────────────┐
│ Body │
│ N bytes │
└───────────────────────────────────────────────────────────────────┘
为了对齐,还可以增加一个保留字节:
Magic 4
Version 1
Flags 1
Serializer 1
Reserved 1
Request ID 8
Body Length 4
----------------
Header 20 bytes
1. Magic
魔数用于快速判断:
这是不是我们的协议?
例如:
private static final int MAGIC = 0x52504331;
可以把它理解为:
RPC1
收到错误魔数时,应立即关闭连接或丢弃数据。
2. Version
协议迟早会升级:
版本 1
版本 2
版本 3
没有版本字段,服务端难以兼容新旧客户端。
3. Flags
Flags 可以表达:
请求还是响应
普通消息还是心跳
是否压缩
是否异常
例如:
0000 0001:请求
0000 0010:响应
0000 0100:心跳
0000 1000:异常
4. Serializer
用于声明 Body 使用哪种序列化方式:
1:Java
2:JSON
3:自定义二进制
5. Request ID
用于关联请求与响应。
6. Body Length
用于解决 TCP 消息边界问题。
接收端先读取固定的 20 字节 Header:
得到 bodyLength = 300
然后继续等待:
至少再收到 300 字节
一条完整消息才组装完成。
八、如何编码?
协议头不应通过 ObjectOutputStream 序列化,而应逐字段写入。
使用 Netty ByteBuf:
public final class RpcProtocol {
public static final int MAGIC = 0x52504331;
public static final byte VERSION = 1;
public static final int HEADER_LENGTH = 20;
private RpcProtocol() {
}
}
编码器:
public final class RpcMessageEncoder
extends MessageToByteEncoder<RpcMessage> {
private final RpcSerializer serializer;
public RpcMessageEncoder(
RpcSerializer serializer
) {
this.serializer = serializer;
}
@Override
protected void encode(
ChannelHandlerContext context,
RpcMessage message,
ByteBuf output
) throws Exception {
byte[] body =
serializer.serialize(message.body());
output.writeInt(RpcProtocol.MAGIC);
output.writeByte(RpcProtocol.VERSION);
output.writeByte(message.flags());
output.writeByte(serializer.id());
output.writeByte(0);
output.writeLong(message.requestId());
output.writeInt(body.length);
output.writeBytes(body);
}
}
这里的 Header 长度始终是:
20 字节
不会因为 Java 类字段或类名发生变化。
九、如何拆包?
Netty 已经提供了按长度字段拆分消息的:
LengthFieldBasedFrameDecoder
它会读取消息中的长度字段,并等待完整帧到达后再向后续 Handler 输出。该解码器专门用于包含长度字段的二进制协议,并允许配置长度字段位置、长度、调整值和需要跳过的头部字节。
我们的协议中:
Header 长度:20
bodyLength 的偏移:
4 + 1 + 1 + 1 + 1 + 8
= 16
bodyLength 自身长度:
4
所以可以配置:
new LengthFieldBasedFrameDecoder(
8 * 1024 * 1024,
16,
4,
0,
0
);
参数含义:
最大帧长度:8 MiB
长度字段偏移:16
长度字段宽度:4
长度调整:0
保留整个 Header
然后再解析协议:
public final class RpcMessageDecoder
extends ByteToMessageDecoder {
private final SerializerRegistry serializers;
public RpcMessageDecoder(
SerializerRegistry serializers
) {
this.serializers = serializers;
}
@Override
protected void decode(
ChannelHandlerContext context,
ByteBuf input,
List<Object> output
) throws Exception {
int magic = input.readInt();
if (magic != RpcProtocol.MAGIC) {
throw new CorruptedFrameException(
"Invalid RPC magic: " + magic
);
}
byte version = input.readByte();
byte flags = input.readByte();
byte serializerId = input.readByte();
input.skipBytes(1);
long requestId = input.readLong();
int bodyLength = input.readInt();
if (bodyLength < 0
|| bodyLength > 8 * 1024 * 1024) {
throw new CorruptedFrameException(
"Invalid body length: " + bodyLength
);
}
byte[] body = new byte[bodyLength];
input.readBytes(body);
RpcSerializer serializer =
serializers.get(serializerId);
Object content =
serializer.deserialize(body);
output.add(
new RpcMessage(
version,
flags,
requestId,
content
)
);
}
}
Pipeline:
channel.pipeline().addLast(
new LengthFieldBasedFrameDecoder(
8 * 1024 * 1024,
16,
4,
0,
0
),
new RpcMessageDecoder(serializers),
new RpcMessageEncoder(serializer),
new RpcResponseHandler(pendingRequests)
);
十、响应如何找回线程?
现在假设三个线程共用一条 Channel:
线程 A 发送 requestId=101
线程 B 发送 requestId=102
线程 C 发送 requestId=103
客户端必须保存三份“未来结果”:
101 → Future A
102 → Future B
103 → Future C
可以使用:
ConcurrentHashMap<
Long,
CompletableFuture<RpcResponse>
> pendingRequests;
CompletableFuture 既是一个 Future,也可以由其他线程显式调用 complete() 或 completeExceptionally() 完成,并支持在完成后触发后续处理。
1. 发送前注册
发送请求时:
public CompletableFuture<RpcResponse> send(
Channel channel,
RpcRequest request
) {
CompletableFuture<RpcResponse> future =
new CompletableFuture<>();
pendingRequests.put(
request.requestId(),
future
);
channel.writeAndFlush(request)
.addListener(writeFuture -> {
if (!writeFuture.isSuccess()) {
CompletableFuture<RpcResponse> removed =
pendingRequests.remove(
request.requestId()
);
if (removed != null) {
removed.completeExceptionally(
writeFuture.cause()
);
}
}
});
return future;
}
必须先:
放入 pendingRequests
再发送。
不能反过来:
先发送
再保存 Future
因为网络响应可能非常快:
请求发送
↓
服务端立即响应
↓
响应 Handler 查找 requestId
↓
Future 还没有放入 Map
于是响应无法匹配。
2. 收到响应
public final class RpcResponseHandler
extends SimpleChannelInboundHandler<RpcResponse> {
private final ConcurrentMap<
Long,
CompletableFuture<RpcResponse>
> pendingRequests;
public RpcResponseHandler(
ConcurrentMap<
Long,
CompletableFuture<RpcResponse>
> pendingRequests
) {
this.pendingRequests = pendingRequests;
}
@Override
protected void channelRead0(
ChannelHandlerContext context,
RpcResponse response
) {
CompletableFuture<RpcResponse> future =
pendingRequests.remove(
response.requestId()
);
if (future == null) {
return;
}
if (response.success()) {
future.complete(response);
} else {
future.completeExceptionally(
new RpcRemoteException(
response.errorType(),
response.errorMessage()
)
);
}
}
}
处理流程:
收到响应 102
↓
pendingRequests.remove(102)
↓
得到 Future B
↓
future.complete(response)
↓
线程 B 继续执行
十一、为什么不只用 Latch?
学习阶段可以这样设计:
CountDownLatch latch =
new CountDownLatch(1);
请求线程:
latch.await();
响应线程:
latch.countDown();
CountDownLatch 初始化时具有固定计数;await() 会等待计数降到 0,countDown() 使计数递减。它是一次性同步工具,计数归零后不能重新设置。
但 RPC 结果不仅需要表达:
完成
还需要表达:
返回值
远程异常
发送失败
超时
取消
如果使用 CountDownLatch,还需要额外维护:
class RpcCallback {
private final CountDownLatch latch =
new CountDownLatch(1);
private volatile Object result;
private volatile Throwable error;
}
这其实是在手工实现一个 Future。
所以更自然的模型是:
CompletableFuture<RpcResponse>
它同时承载:
完成状态
返回值
异常
回调
组合操作
十二、两个 Future 不一样
RPC 中通常同时存在两类 Future。
1. ChannelFuture
ChannelFuture writeFuture =
channel.writeAndFlush(request);
它表示:
这次 Channel I/O 操作是否完成
Netty 的所有 Channel I/O 操作都是异步的,方法会立即返回 ChannelFuture。这个 Future 表示写入、连接、绑定或关闭等 I/O 操作的完成状态。
它不表示:
远程方法已经执行完毕
即使:
channel.writeAndFlush(request)
.sync();
执行成功,也只能说明对应的写操作完成。
不能说明:
服务端已经反序列化请求
服务端已经调用业务方法
数据库已经提交
响应已经返回客户端
2. RPC Future
CompletableFuture<RpcResponse> rpcFuture;
它表示:
这个 requestId 对应的远程调用结果
两者关系是:
ChannelFuture
管理网络写操作
RpcFuture
管理远程方法结果
完整流程:
writeAndFlush(request)
↓
ChannelFuture 成功
↓
请求成功交给网络层
↓
等待服务端处理
↓
收到 RpcResponse
↓
RpcFuture 完成
这是手写 RPC 中非常重要的边界。
十三、代理怎样返回结果?
代理层拿到 RPC Future 后,可以选择同步等待:
RpcResponse response =
rpcFuture.get(
3,
TimeUnit.SECONDS
);
然后返回:
return response.result();
完整代理:
public final class RpcProxyFactory {
private final RpcClient rpcClient;
private final AtomicLong requestIds =
new AtomicLong();
public RpcProxyFactory(
RpcClient rpcClient
) {
this.rpcClient = rpcClient;
}
public <T> T create(
Class<T> serviceInterface
) {
Object proxy = Proxy.newProxyInstance(
serviceInterface.getClassLoader(),
new Class<?>[]{serviceInterface},
(proxyObject, method, arguments) -> {
if (method.getDeclaringClass()
== Object.class) {
return invokeObjectMethod(
proxyObject,
method,
arguments
);
}
long requestId =
requestIds.incrementAndGet();
RpcRequest request =
createRequest(
requestId,
serviceInterface,
method,
arguments
);
CompletableFuture<RpcResponse> future =
rpcClient.send(request);
RpcResponse response;
try {
response = future.get(
3,
TimeUnit.SECONDS
);
} catch (TimeoutException exception) {
rpcClient.cancel(requestId);
throw new RpcTimeoutException(
"RPC timeout: "
+ requestId,
exception
);
}
if (!response.success()) {
throw new RpcRemoteException(
response.errorType(),
response.errorMessage()
);
}
return response.result();
}
);
return serviceInterface.cast(proxy);
}
}
这样调用者看到的仍然是:
User user =
userService.findById(1001L);
但当前线程实际上经历了:
代理拦截
↓
构造请求
↓
创建 Future
↓
发送网络消息
↓
等待 Future
↓
响应完成 Future
↓
返回结果
十四、同步还是异步?
接口可以设计为同步:
User findById(long id);
代理内部需要等待:
future.get(timeout, unit);
也可以设计为异步:
CompletableFuture<User> findById(long id);
代理直接返回 Future:
CompletableFuture<RpcResponse> rpcFuture =
rpcClient.send(request);
return rpcFuture.thenApply(response -> {
if (!response.success()) {
throw new RpcRemoteException(
response.errorType(),
response.errorMessage()
);
}
return (User) response.result();
});
调用者:
CompletableFuture<User> future =
userService.findById(1001L);
future.thenAccept(user -> {
System.out.println(user.name());
});
同步接口更接近普通 Java 方法:
容易使用
调用线程会等待
异步接口则:
不阻塞调用线程
需要处理 Future
调用链更复杂
RPC 框架可以同时支持两种接口,但不能把异步接口再次阻塞等待,否则就失去了它的意义。
十五、一条连接够吗?
假设客户端有十个调用线程。
一种设计是:
十个线程
↓
共用一条 TCP 连接
只要具备:
消息边界
requestId
并发安全的发送
响应匹配
流量控制
一条连接可以同时承载多个未完成请求:
Channel
├── requestId 101
├── requestId 102
├── requestId 103
└── requestId 104
这叫请求多路复用。
1. 为什么还要多条连接?
单连接可能出现:
发送缓冲区拥塞
单个 EventLoop 压力较高
单连接故障影响全部请求
带宽或拥塞窗口受限
所以客户端可能为一个服务节点建立多条连接:
Provider A
├── Channel 0
├── Channel 1
├── Channel 2
└── Channel 3
但连接数量不是越多越好。
每条连接都会消耗:
文件描述符
内核 Socket 状态
发送和接收缓冲区
心跳任务
连接维护成本
更合理的策略是从少量连接开始,根据:
请求吞吐
单连接积压
事件循环负载
延迟
故障隔离需求
进行压测调整。
十六、连接池是什么?
传统数据库连接池经常使用:
借出连接
↓
独占使用
↓
归还连接
但支持 requestId 多路复用的 RPC Channel 不一定需要被某个线程独占。
更接近:
从多个可用 Channel 中选择一个
↓
在该 Channel 上发送请求
↓
请求通过 requestId 独立等待响应
连接池可以这样表示:
public final class RpcChannelPool {
private final List<Channel> channels;
private final AtomicInteger index =
new AtomicInteger();
public RpcChannelPool(
List<Channel> channels
) {
this.channels =
List.copyOf(channels);
}
public Channel select() {
int size = channels.size();
for (int attempt = 0;
attempt < size;
attempt++) {
int position = Math.floorMod(
index.getAndIncrement(),
size
);
Channel channel =
channels.get(position);
if (channel.isActive()
&& channel.isWritable()) {
return channel;
}
}
throw new RpcUnavailableException(
"No available RPC channel"
);
}
}
1. 按服务节点管理
可以维护:
ConcurrentHashMap<
InetSocketAddress,
RpcChannelPool
> pools;
结构:
10.0.0.1:8080 → ChannelPool A
10.0.0.2:8080 → ChannelPool B
10.0.0.3:8080 → ChannelPool C
2. 不要每条连接创建 EventLoopGroup
错误做法:
private Channel createConnection() {
NioEventLoopGroup group =
new NioEventLoopGroup();
// 创建一条连接
}
这样每建立一条连接都创建一组线程和 Selector。
更合理的是:
一个 RpcClient
↓
一个共享 EventLoopGroup
↓
管理多条客户端 Channel
NioEventLoopGroup 是面向 NIO Selector Channel 的多线程 EventLoopGroup;一个 EventLoop 在 Channel 注册后负责其 I/O,而且通常会管理多条 Channel。
客户端工厂:
public final class RpcClientFactory
implements AutoCloseable {
private final NioEventLoopGroup eventLoopGroup =
new NioEventLoopGroup();
private final Bootstrap bootstrap;
public RpcClientFactory(
ChannelInitializer<SocketChannel> initializer
) {
this.bootstrap = new Bootstrap()
.group(eventLoopGroup)
.channel(NioSocketChannel.class)
.handler(initializer);
}
public CompletableFuture<Channel> connect(
InetSocketAddress address
) {
CompletableFuture<Channel> result =
new CompletableFuture<>();
bootstrap.connect(address)
.addListener(future -> {
if (future.isSuccess()) {
result.complete(
((ChannelFuture) future)
.channel()
);
} else {
result.completeExceptionally(
future.cause()
);
}
});
return result;
}
@Override
public void close() {
eventLoopGroup.shutdownGracefully();
}
}
Netty 的 Bootstrap 是用于初始化客户端 Channel 的辅助类,TCP 客户端通过 connect() 发起连接。
十七、服务端怎样执行?
服务端收到请求后,需要完成:
解码请求
↓
查找服务
↓
查找方法
↓
调用方法
↓
构造响应
↓
发送响应
1. 服务注册表
public final class ServiceRegistry {
private final ConcurrentMap<String, Object> services =
new ConcurrentHashMap<>();
public void register(
Class<?> serviceInterface,
Object implementation
) {
if (!serviceInterface.isInstance(
implementation
)) {
throw new IllegalArgumentException(
"Implementation does not implement "
+ serviceInterface.getName()
);
}
services.put(
serviceInterface.getName(),
implementation
);
}
public Object get(String serviceName) {
Object service = services.get(serviceName);
if (service == null) {
throw new RpcServiceNotFoundException(
serviceName
);
}
return service;
}
}
2. 查找方法
private Method resolveMethod(
Object service,
RpcRequest request
) throws ClassNotFoundException,
NoSuchMethodException {
String[] typeNames =
request.parameterTypeNames();
Class<?>[] parameterTypes =
new Class<?>[typeNames.length];
for (int i = 0;
i < typeNames.length;
i++) {
parameterTypes[i] =
resolveType(typeNames[i]);
}
return service.getClass().getMethod(
request.methodName(),
parameterTypes
);
}
基本类型需要单独处理:
private Class<?> resolveType(
String typeName
) throws ClassNotFoundException {
return switch (typeName) {
case "boolean" -> boolean.class;
case "byte" -> byte.class;
case "short" -> short.class;
case "int" -> int.class;
case "long" -> long.class;
case "float" -> float.class;
case "double" -> double.class;
case "char" -> char.class;
case "void" -> void.class;
default -> Class.forName(typeName);
};
}
3. 执行调用
public RpcResponse invoke(
RpcRequest request
) {
try {
Object service =
registry.get(
request.serviceName()
);
Method method =
resolveMethod(
service,
request
);
Object result =
method.invoke(
service,
request.arguments()
);
return RpcResponse.success(
request.requestId(),
result
);
} catch (InvocationTargetException exception) {
Throwable target =
exception.getTargetException();
return RpcResponse.failure(
request.requestId(),
target
);
} catch (Throwable throwable) {
return RpcResponse.failure(
request.requestId(),
throwable
);
}
}
4. 缓存 Method
反复执行:
getMethod()
会产生额外查找开销。
可以用方法签名作为 Key:
服务名
+
方法名
+
参数类型列表
缓存:
ConcurrentHashMap<
MethodKey,
Method
> methodCache;
十八、业务不能堵住 I/O
服务端 Handler 可能运行在 Netty EventLoop 上。
一个 EventLoop 通常会处理多条 Channel 的 I/O。如果直接在 Handler 中执行慢数据库查询或长时间业务计算,同一 EventLoop 管理的其他连接也会延迟处理。
错误模式:
@Override
protected void channelRead0(
ChannelHandlerContext context,
RpcRequest request
) {
// 可能阻塞 5 秒
RpcResponse response =
serviceInvoker.invoke(request);
context.writeAndFlush(response);
}
可以使用业务线程池:
public final class RpcRequestHandler
extends SimpleChannelInboundHandler<RpcRequest> {
private final ExecutorService businessExecutor;
private final ServiceInvoker serviceInvoker;
public RpcRequestHandler(
ExecutorService businessExecutor,
ServiceInvoker serviceInvoker
) {
this.businessExecutor =
businessExecutor;
this.serviceInvoker =
serviceInvoker;
}
@Override
protected void channelRead0(
ChannelHandlerContext context,
RpcRequest request
) {
businessExecutor.execute(() -> {
RpcResponse response =
serviceInvoker.invoke(request);
context.writeAndFlush(response);
});
}
}
不过把任务放入业务线程池以后,还需要考虑:
线程池大小
队列容量
拒绝策略
请求超时
服务关闭
上下文传播
过载保护
无界队列会把压力从 EventLoop 转移到内存,并不是真正的治理。
十九、超时如何处理?
请求发送后,不能无限等待。
发送时注册超时任务:
public CompletableFuture<RpcResponse> send(
Channel channel,
RpcRequest request,
Duration timeout
) {
CompletableFuture<RpcResponse> result =
new CompletableFuture<>();
long requestId =
request.requestId();
pendingRequests.put(
requestId,
result
);
ScheduledFuture<?> timeoutTask =
channel.eventLoop().schedule(
() -> {
CompletableFuture<RpcResponse> removed =
pendingRequests.remove(
requestId
);
if (removed != null) {
removed.completeExceptionally(
new RpcTimeoutException(
"Request timed out: "
+ requestId
)
);
}
},
timeout.toMillis(),
TimeUnit.MILLISECONDS
);
result.whenComplete(
(response, throwable) ->
timeoutTask.cancel(false)
);
channel.writeAndFlush(request)
.addListener(writeFuture -> {
if (!writeFuture.isSuccess()) {
CompletableFuture<RpcResponse> removed =
pendingRequests.remove(
requestId
);
if (removed != null) {
removed.completeExceptionally(
writeFuture.cause()
);
}
}
});
return result;
}
这里必须覆盖几个结束路径:
正常响应
写入失败
请求超时
连接关闭
客户端取消
无论哪条路径结束,都必须从:
pendingRequests
中删除。
否则会发生内存泄漏:
请求已经永远不会返回
但 Future 仍保存在 Map 中
二十、连接断开怎么办?
当 Channel 断开时,这条连接上的所有未完成请求都不可能再收到响应。
因此应完成全部 Future:
@Override
public void channelInactive(
ChannelHandlerContext context
) {
RpcConnectionClosedException exception =
new RpcConnectionClosedException(
context.channel().remoteAddress()
);
pendingRequests.forEach(
(requestId, future) -> {
if (pendingRequests.remove(
requestId,
future
)) {
future.completeExceptionally(
exception
);
}
}
);
context.fireChannelInactive();
}
更完整的设计中,pendingRequests 最好按 Channel 隔离:
Channel A
└── PendingMap A
Channel B
└── PendingMap B
这样 Channel A 断开时,只失败 A 上的请求,不会错误影响 B。
二十一、请求 ID 怎么生成?
笔记中使用随机 long 作为 UID。
随机数不能保证绝对唯一。
更简单的单进程实现是:
AtomicLong requestIds =
new AtomicLong();
生成:
long requestId =
requestIds.incrementAndGet();
只要保证:
同一客户端实例内
尚未完成的请求 ID 不重复
就能完成关联。
分布式追踪或跨进程全局标识可以使用更大的结构:
节点标识
时间戳
进程启动序列
本地递增序号
不过 RPC 请求关联 ID 和全局业务 ID 是两个不同概念。
requestId 主要服务于:
一端连接上的请求响应匹配
二十二、重试安全吗?
RPC 超时只表示:
客户端在规定时间内没有拿到响应
它不一定表示:
服务端没有执行方法
可能发生:
客户端发送成功
↓
服务端修改数据库成功
↓
响应在网络中丢失
↓
客户端超时
如果客户端自动重试:
同一个扣款请求再次执行
可能造成重复操作。
因此重试需要考虑:
幂等性
业务请求 ID
去重表
状态查询
重试次数
退避策略
RPC 框架不能仅根据:
TimeoutException
就假设远端方法没有执行。
二十三、完整调用链
现在重新看这一行代码:
User user =
userService.findById(1001L);
完整过程如下。
客户端
1. 调用代理对象 findById()
2. InvocationHandler 拦截调用
3. 生成 requestId
4. 封装 RpcRequest
- serviceName
- methodName
- parameterTypes
- arguments
5. 创建 CompletableFuture
6. requestId → Future 放入 PendingMap
7. 序列化 RpcRequest
8. 编码 Header + Body
9. 选择可用 Channel
10. writeAndFlush()
网络
11. TCP 传输字节流
12. 服务端 LengthFieldBasedFrameDecoder
根据 bodyLength 组装完整消息
13. 协议解码器解析 Header 和 Body
服务端
14. 根据 serviceName 查找实现对象
15. 根据方法名和参数类型查找 Method
16. 在业务线程池中调用 Method
17. 得到结果或异常
18. 构造 RpcResponse
19. 使用相同 requestId
20. 编码并发送响应
客户端返回
21. 客户端解码响应
22. 根据 requestId 查找 Future
23. 从 PendingMap 删除 Future
24. complete 或 completeExceptionally
25. 代理线程恢复执行
26. 返回 User 或抛出异常
最终,调用者看到的仍然只是:
User user =
userService.findById(1001L);
二十四、常见误区
1. 动态代理会执行远程方法
错误。
动态代理只负责拦截本地接口调用。
真正的远程执行由:
消息封装
网络发送
服务端分发
共同完成。
2. 只发送方法名就够了
错误。
方法可能重载,还需要参数类型。
服务也可能存在多个版本,还需要完整服务标识。
3. TCP 一次 write 对应一次 read
错误。
TCP 是字节流,没有应用消息边界。
必须通过:
长度字段
分隔符
固定帧长
定义协议帧。
4. 粘包是 Netty 的问题
错误。
粘包和半包来自 TCP 字节流语义,Netty 只是提供了方便的拆帧工具。
5. Header 对象序列化后长度固定
错误。
Java 对象序列化包含流头、类描述和对象图信息,不能把某次测得的长度写死为协议头长度。
6. requestId 使用随机数就绝不会重复
错误。
随机只能降低碰撞概率,不能从定义上保证唯一。
7. 多线程共用连接必须按发送顺序返回
错误。
只要协议中有 requestId,响应可以乱序完成。
8. 一条连接一次只能有一个请求
错误。
具有消息边界和 requestId 后,同一连接可以存在多个并发未完成请求。
9. 连接池一定要借出和归还
不一定。
支持多路复用的 RPC Channel 可以同时承载多个请求,不一定需要由一个调用线程独占。
10. writeAndFlush().sync() 等于等待 RPC 返回
错误。
它等待的是 Channel 写操作,不是远程业务结果。
11. ChannelFuture 就是 RPC Future
错误。
ChannelFuture
表示 I/O 操作结果
RPC Future
表示远程方法结果
12. CountDownLatch 是最合适的 RPC 回调
不一定。
它只能表达一次性计数归零,还要额外维护结果和异常;CompletableFuture 更符合“未来结果”的语义。
13. Future 完成后不用删除
错误。
不删除会让 PendingMap 不断增长。
14. 超时后 Future 可以继续留着
错误。
超时后应从 PendingMap 原子删除;迟到响应只能忽略、记录或交给专门策略处理。
15. 连接断开只影响下一次请求
错误。
这条连接上的全部未完成请求都需要立即失败。
16. 每个连接创建一个 NioEventLoopGroup
错误。
EventLoopGroup 应被多个 Channel 共享,否则会创建大量线程和 Selector。
17. 服务端可以直接在 EventLoop 查询数据库
不推荐。
一个 EventLoop 通常管理多条 Channel,长时间阻塞会延迟同一 EventLoop 上的其他连接。
18. Java 序列化可以直接接收任意网络对象
危险。
对不可信数据进行反序列化前,应限制允许类型、对象图复杂度和消息大小,并显式配置过滤策略。
19. RPC 超时说明服务端没有执行
错误。
服务端可能已经执行成功,只是响应没有及时返回。
20. 自动重试一定安全
错误。
非幂等方法可能被重复执行。
总结
RPC 的目标是:
让调用远程服务
看起来像调用本地接口
但这个“像”只存在于调用形式上。
本地调用:
方法参数
↓
线程栈
↓
实现类
↓
返回值
远程调用:
方法参数
↓
动态代理
↓
RpcRequest
↓
序列化
↓
协议编码
↓
TCP
↓
协议解码
↓
服务注册表
↓
反射调用
↓
RpcResponse
↓
requestId
↓
Future
↓
返回值
动态代理负责:
把方法调用转换成请求
协议负责:
让通信双方理解字节含义
长度字段负责:
从 TCP 字节流中切出完整消息
requestId 负责:
在并发和乱序响应中匹配请求
CompletableFuture 负责:
保存未来的返回值或异常
连接池负责:
管理和复用多条网络连接
服务注册表负责:
从服务名找到实现对象
反射负责:
根据方法签名执行真正的业务方法
需要特别区分两种 Future:
ChannelFuture
表示网络 I/O 操作结果
RpcFuture
表示远程方法调用结果
也需要区分两种“成功”:
消息写出成功
≠
服务端执行成功
≠
远端业务提交成功
最后,可以用一句话概括整个 RPC 调用过程:
RPC 用动态代理隐藏调用入口,用协议和序列化跨越网络,用 requestId 和 Future 找回结果,最终把一次异步、乱序的网络通信,重新包装成一次看似普通的方法调用。
评论区