前言
上一篇文章中,我们已经写出了一个最小 RPC 调用链:
动态代理
↓
构造 RpcRequest
↓
协议编码
↓
Netty 发送
↓
服务端反射调用
↓
返回 RpcResponse
↓
requestId 匹配 Future
单线程测试时,一切正常。
客户端发送一个请求:
Header + Body
服务端读取一次,也恰好得到:
Header + Body
于是我们很容易写出这样的解码代码:
@Override
public void channelRead(
ChannelHandlerContext context,
Object message
) throws Exception {
ByteBuf buffer = (ByteBuf) message;
byte[] headerBytes = new byte[HEADER_LENGTH];
buffer.readBytes(headerBytes);
RpcHeader header =
deserializeHeader(headerBytes);
byte[] bodyBytes =
new byte[header.bodyLength()];
buffer.readBytes(bodyBytes);
RpcRequest request =
deserializeBody(bodyBytes);
handle(request);
}
一个请求时,它可能真的能运行。
但当 20 个线程同时调用远程方法,并复用同一条 TCP 连接时,服务端开始出现:
StreamCorruptedException
EOFException
OptionalDataException
IndexOutOfBoundsException
或者:
请求头可以解析
请求体反序列化失败
这时问题就来了:
客户端明明发送了完整消息,服务端为什么只收到一部分?
一次
writeAndFlush(),为什么不能对应一次channelRead()?
一次读取为什么会同时拿到三条请求?
请求头已经读出来了,请求体不够时应该怎么办?
为什么直接
return仍然会丢数据?
readInt()和getInt()在解码器里有什么区别?
未处理完的数据应该由谁保存?
ByteToMessageDecoder到底替我们完成了什么?
为什么
decode()有时会被连续调用多次?
requestId能不能解决粘包?
客户端和服务端是否都要配置解码器?
这些问题的答案都来自同一个事实:
TCP 传输的是连续字节流,而不是一条条 RPC 消息。
一、不是序列化坏了
先看一个常见错误现场。
RPC 协议包含:
固定长度 Header
+
变长 Body
假设:
Header:20 字节
Body:300 字节
完整消息长度为:
320 字节
第一次网络读取只得到:
Header:20 字节
Body:100 字节
此时缓冲区共有:
120 字节
错误代码先把 Header 读走:
RpcHeader header = readHeader(buffer);
读指针已经向后移动 20 字节:
读取前:
readerIndex
↓
┌────────┬─────────────────────┐
│ Header │ Body 前 100 字节 │
└────────┴─────────────────────┘
读取 Header 后:
readerIndex
↓
┌────────┬─────────────────────┐
│ 已读取 │ Body 前 100 字节 │
└────────┴─────────────────────┘
然后发现 Body 不够:
if (buffer.readableBytes()
< header.bodyLength()) {
return;
}
看起来只是等待下一批数据。
但问题是:
Header 已经被消费
下一批数据到来后,缓冲区可能变成:
Body 前 100 字节
+
Body 后 200 字节
+
下一条消息的 Header
解码器却把当前位置再次当成 Header 起点:
Body 字节
↓
被错误地当作 Header 解析
随后再交给 ObjectInputStream,自然会从一个错误位置读取对象流。
所以异常看起来发生在:
反序列化
真正的错误却发生在:
协议边界判断
+
读指针移动
二、TCP 只认字节
TCP 为应用提供可靠、有序的字节流服务。它保证字节顺序,却没有应用层消息边界的概念。
假设客户端连续发送三条 RPC 请求:
消息 A:320 字节
消息 B:280 字节
消息 C:400 字节
客户端看到的是:
write(A)
write(B)
write(C)
TCP 看到的却更接近:
连续的 1000 个字节
至于接收端每次 read() 得到多少字节,取决于:
数据到达速度
内核接收缓冲区
应用读取时机
接收缓冲区容量
网络分段与重传
系统调度
因此,一次读取可能出现三种情况。
1. 恰好一条
┌───────────┐
│ Message A │
└───────────┘
这是测试中最容易遇到的情况,也最容易让错误代码看起来正确。
2. 一次多条
┌───────────┬───────────┬───────────┐
│ Message A │ Message B │ Message C │
└───────────┴───────────┴───────────┘
如果 Handler 只处理第一条:
decodeOneMessage(buffer);
return;
后面的完整消息就没有在当前调用中继续消费。
3. 一条被拆开
第一次:
┌────────┬──────────────┐
│ Header │ Body Part 1 │
└────────┴──────────────┘
第二次:
┌──────────────┐
│ Body Part 2 │
└──────────────┘
这就是通常所说的半包。
所谓“粘包”和“半包”,并不是 TCP 破坏了数据,而是应用把连续字节流误认为了离散消息。
三、循环还不够
发现一次读取可能包含多条消息后,第一反应通常是增加循环:
while (buffer.readableBytes()
>= HEADER_LENGTH) {
RpcHeader header = readHeader(buffer);
if (buffer.readableBytes()
< header.bodyLength()) {
break;
}
RpcRequest request =
readBody(buffer, header);
handle(request);
}
循环确实解决了:
一次读取包含多条完整消息
但它仍然没有解决:
最后一条消息只有 Header,
Body 还没有到齐
因为 readHeader() 已经移动了 readerIndex。
执行 break 只能停止循环,不能让读指针自动回到 Header 开始位置。
所以,正确解码需要同时解决两个问题:
横向:
一次读取中可能存在多条消息
纵向:
一条消息可能跨越多次读取
只写 while,只能解决第一个问题。
四、边界怎么算?
继续使用上一篇 RPC 协议:
┌─────────┬─────────┬─────────┬────────────┬─────────────┐
│ Magic │ Version │ Flags │ Request ID │ Body Length │
│ 4 bytes │ 1 byte │ 3 bytes │ 8 bytes │ 4 bytes │
└─────────┴─────────┴─────────┴────────────┴─────────────┘
┌────────────────────────────────────────────────────────┐
│ Body │
│ N bytes │
└────────────────────────────────────────────────────────┘
为方便对齐,前三个单字节字段可以设计为:
Version 1
Flags 1
Serializer 1
Reserved 1
完整 Header:
Magic 4
Version 1
Flags 1
Serializer 1
Reserved 1
Request ID 8
Body Length 4
----------------
总计 20 字节
一条完整消息长度为:
frameLength
=
HEADER_LENGTH
+
bodyLength
解码器首先必须回答:
当前缓冲区是否至少有 20 字节?
如果连 Header 都不完整:
直接等待更多字节
Header 完整后,读取 bodyLength,再判断:
当前缓冲区是否至少有
20 + bodyLength 字节?
只有完整帧已经到达,才能推进读指针。
五、先看不动指针
ByteBuf 中两类读取方法必须分清。
1. read 方法
int value = buffer.readInt();
会读取数据,并推进 readerIndex。
ByteBuf.readInt() 会从当前 readerIndex 读取 4 个字节,并将读索引向后移动 4 个字节。
2. get 方法
int value = buffer.getInt(index);
读取指定位置的数据,但不移动 readerIndex。
因此,在还不确定完整消息是否到达时,可以先“偷看”长度字段:
int start = buffer.readerIndex();
int bodyLength =
buffer.getInt(
start + BODY_LENGTH_OFFSET
);
注意不能写成:
buffer.getInt(BODY_LENGTH_OFFSET);
因为当前帧不一定从底层缓冲区索引 0 开始。
Netty 官方文档也特别提醒,自定义帧解码器检查数据时应使用 getInt(in.readerIndex() + offset) 一类绝对读取;数据不足时,应在不修改读索引的情况下返回。
六、手写一个解码器
先看一个不依赖内置长度解码器的实现。
public final class RpcFrameDecoder
extends ByteToMessageDecoder {
private static final int MAGIC =
0x52504331;
private static final int HEADER_LENGTH =
20;
private static final int BODY_LENGTH_OFFSET =
16;
private static final int MAX_BODY_LENGTH =
8 * 1024 * 1024;
@Override
protected void decode(
ChannelHandlerContext context,
ByteBuf input,
List<Object> output
) throws Exception {
if (input.readableBytes()
< HEADER_LENGTH) {
return;
}
int frameStart =
input.readerIndex();
int magic =
input.getInt(frameStart);
if (magic != MAGIC) {
throw new CorruptedFrameException(
"Invalid magic: " + magic
);
}
int bodyLength =
input.getInt(
frameStart
+ BODY_LENGTH_OFFSET
);
if (bodyLength < 0
|| bodyLength > MAX_BODY_LENGTH) {
throw new CorruptedFrameException(
"Invalid body length: "
+ bodyLength
);
}
int frameLength =
HEADER_LENGTH + bodyLength;
if (input.readableBytes()
< frameLength) {
return;
}
int actualMagic =
input.readInt();
byte version =
input.readByte();
byte flags =
input.readByte();
byte serializerId =
input.readByte();
input.skipBytes(1);
long requestId =
input.readLong();
int actualBodyLength =
input.readInt();
if (actualMagic != MAGIC
|| actualBodyLength
!= bodyLength) {
throw new CorruptedFrameException(
"RPC frame changed while decoding"
);
}
byte[] body =
new byte[bodyLength];
input.readBytes(body);
output.add(
new RpcFrame(
version,
flags,
serializerId,
requestId,
body
)
);
}
}
这段代码最关键的不是序列化,而是解码顺序:
1. Header 不够
return
2. 使用 getInt 查看 Body 长度
不移动 readerIndex
3. 验证长度是否合法
4. 完整帧不够
return
5. 完整帧已经到达
才开始 read
6. 解析出一条 RpcFrame
放入 output
因此,即使第一次只到达:
20 字节 Header
+
100 字节 Body
而 Body 总长为 300 字节,readerIndex 仍会停留在 Header 开头。
下一批 200 字节到来后,解码器可以从同一个位置重新判断。
七、谁保存剩余数据?
如果自己继承:
ChannelInboundHandlerAdapter
那么通常需要自己准备一个连接级缓冲区:
上一次剩余字节
+
本次新到字节
自己还要处理:
扩容
复制
读写索引
生命周期
释放
连接关闭
异常清理
ByteToMessageDecoder 已经把这部分逻辑封装了。
它内部维护累积缓冲区,并提供两类 Cumulator:
MERGE_CUMULATOR
把新数据合并到一个 ByteBuf
可能产生内存复制
COMPOSITE_CUMULATOR
使用 CompositeByteBuf 组合数据
尽量避免复制
Netty 官方文档明确说明,ByteToMessageDecoder 会以流式方式把输入字节解码成消息,并通过内部累积缓冲区保存尚未消费的数据。
流程可以抽象为:
第一次 channelRead
↓
收到 Header + 半个 Body
↓
decode() 发现帧不完整
↓
readerIndex 不动
↓
剩余数据保留在累积缓冲区
第二次 channelRead
↓
新数据进入
↓
与上次剩余数据累积
↓
decode() 再次判断
↓
完整帧被解析
因此,不需要自己在 Handler 成员变量中维护:
private ByteBuf remainingBuffer;
也不需要自己手动把两次读取的字节数组拼起来。
八、decode 为什么会循环?
假设一次网络读取包含三条完整消息:
Message A
+
Message B
+
Message C
自定义 decode() 只解析一条:
output.add(messageA);
返回后,ByteToMessageDecoder 会发现输入缓冲区仍然有可读数据,于是继续调用 decode()。
其内部 callDecode() 会在仍应继续解码时反复调用 decode();默认并不是一次 channelRead() 只解出一条消息。
过程如下:
channelRead()
↓
累积 ByteBuf
↓
decode() → Message A
↓
decode() → Message B
↓
decode() → Message C
↓
没有完整消息
↓
结束
所以通常不需要在 decode() 内再写一个复杂的外层死循环。
下面两种形式都能工作。
1. 每次解一条
@Override
protected void decode(
ChannelHandlerContext context,
ByteBuf input,
List<Object> output
) {
RpcFrame frame =
tryDecodeOne(input);
if (frame != null) {
output.add(frame);
}
}
让父类负责再次调用。
2. 一次解多条
@Override
protected void decode(
ChannelHandlerContext context,
ByteBuf input,
List<Object> output
) {
while (true) {
RpcFrame frame =
tryDecodeOne(input);
if (frame == null) {
return;
}
output.add(frame);
}
}
后一种方式需要更加谨慎,避免在:
没有消费字节
也没有产生消息
的情况下无限循环。
九、out 做了什么?
解码器完成一条消息后调用:
output.add(rpcFrame);
这里的 output 不是最终业务结果集合。
它表示:
当前解码器产生的下游消息
Pipeline:
原始 ByteBuf
↓
RpcFrameDecoder
↓
RpcMessageDecoder
↓
RpcRequestHandler
执行过程:
ByteBuf
↓
RpcFrameDecoder 输出 RpcFrame
↓
RpcMessageDecoder 输出 RpcRequest
↓
RpcRequestHandler 执行业务
如果一次读取解出三条消息:
out.add(requestA)
out.add(requestB)
out.add(requestC)
后面的 Handler 会分别收到它们,而不是再拿到一团无法区分的字节。
十、别急着自己拆
如果协议中已经存在明确的长度字段,通常可以直接使用:
LengthFieldBasedFrameDecoder
它专门根据消息中的长度字段动态拆分 ByteBuf,适合 Header 中携带 Body 长度或整帧长度的二进制协议。
我们的协议中:
Header 长度:20
Body Length 偏移:16
Body Length 宽度:4
Body Length 表示:仅 Body 长度
配置如下:
pipeline.addLast(
new LengthFieldBasedFrameDecoder(
8 * 1024 * 1024,
16,
4,
0,
0
)
);
参数依次表示:
maxFrameLength
最大帧长度
lengthFieldOffset
长度字段起始偏移
lengthFieldLength
长度字段占用字节数
lengthAdjustment
长度修正值
initialBytesToStrip
输出前跳过多少字节
对于当前协议:
完整帧长度
=
长度字段末尾位置
+
bodyLength
=
20 + bodyLength
因此 lengthAdjustment 为 0。
initialBytesToStrip 为 0,表示把完整 Header 和 Body 都交给下一个解码器。
十一、拆帧与解码分开
推荐把两件事拆开。
1. 帧解码器
只负责回答:
一条完整消息从哪里开始,
到哪里结束?
new LengthFieldBasedFrameDecoder(
MAX_FRAME_LENGTH,
BODY_LENGTH_OFFSET,
Integer.BYTES,
0,
0
)
2. 消息解码器
负责回答:
这条完整消息代表什么对象?
public final class RpcMessageDecoder
extends MessageToMessageDecoder<ByteBuf> {
private static final int MAGIC =
0x52504331;
private final SerializerRegistry serializers;
public RpcMessageDecoder(
SerializerRegistry serializers
) {
this.serializers = serializers;
}
@Override
protected void decode(
ChannelHandlerContext context,
ByteBuf frame,
List<Object> output
) throws Exception {
int magic = frame.readInt();
if (magic != MAGIC) {
throw new CorruptedFrameException(
"Invalid RPC magic"
);
}
byte version =
frame.readByte();
byte flags =
frame.readByte();
byte serializerId =
frame.readByte();
frame.skipBytes(1);
long requestId =
frame.readLong();
int bodyLength =
frame.readInt();
if (bodyLength != frame.readableBytes()) {
throw new CorruptedFrameException(
"Body length mismatch"
);
}
byte[] body =
new byte[bodyLength];
frame.readBytes(body);
RpcSerializer serializer =
serializers.get(serializerId);
Object message =
serializer.deserialize(
flags,
body
);
output.add(
new RpcMessage(
version,
flags,
requestId,
message
)
);
}
}
Pipeline:
pipeline.addLast(
new LengthFieldBasedFrameDecoder(
MAX_FRAME_LENGTH,
BODY_LENGTH_OFFSET,
Integer.BYTES,
0,
0
),
new RpcMessageDecoder(serializers),
new RpcRequestHandler(serviceInvoker)
);
这样每层职责更加明确:
LengthFieldBasedFrameDecoder
处理 TCP 消息边界
RpcMessageDecoder
处理协议字段和反序列化
RpcRequestHandler
处理 RPC 业务
十二、编码也要成帧
服务端能够正确拆帧的前提,是客户端按照同一协议编码。
public final class RpcMessageEncoder
extends MessageToByteEncoder<RpcMessage> {
private static final int MAGIC =
0x52504331;
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(MAGIC);
output.writeByte(message.version());
output.writeByte(message.flags());
output.writeByte(serializer.id());
output.writeByte(0);
output.writeLong(message.requestId());
output.writeInt(body.length);
output.writeBytes(body);
}
}
Header 与 Body 最好在同一个编码器中构造成一个逻辑帧:
一个 RpcMessage
↓
一次编码
↓
Header + Body
不要在业务代码中分散执行:
channel.write(header);
channel.write(bodyPart1);
channel.writeAndFlush(bodyPart2);
否则协议完整性会依赖更多外部状态,也更容易在异常、并发和对象生命周期处理中出错。
十三、两端都要解码
RPC 是双向通信。
客户端发送:
RpcRequest
服务端返回:
RpcResponse
两个方向都运行在 TCP 字节流之上,因此:
服务端需要拆请求帧
客户端需要拆响应帧
不能只在服务端配置:
LengthFieldBasedFrameDecoder
而让客户端的响应 Handler 直接对原始 ByteBuf 反序列化。
客户端 Pipeline:
pipeline.addLast(
new LengthFieldBasedFrameDecoder(
MAX_FRAME_LENGTH,
BODY_LENGTH_OFFSET,
Integer.BYTES,
0,
0
),
new RpcMessageDecoder(serializers),
new RpcMessageEncoder(serializer),
new RpcResponseHandler(pendingRequests)
);
服务端 Pipeline:
pipeline.addLast(
new LengthFieldBasedFrameDecoder(
MAX_FRAME_LENGTH,
BODY_LENGTH_OFFSET,
Integer.BYTES,
0,
0
),
new RpcMessageDecoder(serializers),
new RpcMessageEncoder(serializer),
new RpcRequestHandler(serviceInvoker)
);
两端可以复用同一种通信帧格式:
Magic
Version
Flags
Serializer
Request ID
Body Length
Body
再通过 flags 区分:
请求
响应
心跳
异常
十四、ID 不负责拆包
requestId 很重要,但它解决的不是消息边界。
它负责:
请求和响应的关联
例如:
请求 101 → Future A
请求 102 → Future B
请求 103 → Future C
响应可能乱序到达:
响应 103
响应 101
响应 102
客户端根据 requestId 找到对应 Future。
而长度字段负责:
从 TCP 字节流中找到完整帧
两者分工如下:
Body Length
解决消息从哪里结束
Request ID
解决响应属于哪个请求
如果没有长度字段:
连一条完整响应在哪里结束都不知道
此时有 requestId 也无法直接解析。
如果没有 requestId:
能够拆出一条条响应
但不知道应该完成哪个 Future
所以 RPC 长连接多路复用必须同时具备:
消息边界
+
请求关联
十五、一次读取不是一次事件
还有一个常见误解:
客户端发送一条消息
↓
服务端触发一次可读事件
↓
服务端读取一条消息
真实过程可能是:
消息 A 到达
↓
Socket 变为可读
↓
EventLoop 尚未立即处理
消息 B 到达
消息 C 到达
↓
数据继续积累在接收缓冲区
EventLoop 开始读取
↓
一次取得 A + B + C 的部分或全部字节
就绪事件只表示:
当前可以读取数据
不表示:
当前恰好存在一条完整业务消息
Selector、epoll 和 Netty EventLoop 关注的是 I/O 就绪状态;业务消息边界仍然必须由应用协议和解码器处理。
十六、别用另一个 I/O 线程跑业务
资料后半部分讨论了把业务任务提交给:
当前 EventLoop
另一个 EventLoop
独立线程池
这里需要区分。
一个 Channel 注册后,其 I/O 操作由对应 EventLoop 处理;一个 EventLoop 通常管理多条 Channel。
因此,解码器应该保持轻量:
读取协议字段
验证长度
反序列化
构造消息对象
不要在解码器中执行:
慢数据库查询
远程 HTTP 调用
大文件读写
长时间计算
Thread.sleep()
也不建议为了处理阻塞业务,随意从同一个 NioEventLoopGroup 中挑另一个 I/O EventLoop。
这样只是把阻塞从:
I/O 线程 A
搬到了:
I/O 线程 B
线程 B 管理的其他 Channel 仍然会受到影响。
更清晰的做法是准备专门的业务执行器:
EventExecutorGroup businessGroup =
new DefaultEventExecutorGroup(16);
然后只让业务 Handler 使用它:
pipeline.addLast(
new LengthFieldBasedFrameDecoder(
MAX_FRAME_LENGTH,
BODY_LENGTH_OFFSET,
Integer.BYTES,
0,
0
)
);
pipeline.addLast(
new RpcMessageDecoder(serializers)
);
pipeline.addLast(
businessGroup,
new RpcRequestHandler(serviceInvoker)
);
ChannelPipeline 支持在添加 Handler 时指定 EventExecutorGroup,该执行器组将负责调用这个 Handler 的方法;DefaultEventExecutorGroup 则使用一组普通 EventExecutor 处理任务。
线程链路变成:
NioEventLoop
↓
读取字节
↓
拆帧
↓
协议解码
↓
业务 EventExecutor
↓
反射调用服务
↓
writeAndFlush 响应
I/O 线程和业务线程的完整选择,将作为下一篇文章的主线。
十七、解码器的限制
1. 必须限制最大帧
不能完全相信网络中的:
bodyLength
攻击者可以构造:
bodyLength = 2 147 483 647
如果程序直接分配数组:
new byte[bodyLength];
可能造成严重内存压力。
因此需要:
最大帧长度
最大 Body 长度
负数检查
字段一致性检查
LengthFieldBasedFrameDecoder 的 maxFrameLength 就是第一道边界。
2. Header 应该定长
不要把一个 Java Header 对象通过 ObjectOutputStream 序列化后,再假设:
序列化结果永远是 110 字节
Java 对象序列化结果可能包含类描述、流信息和对象引用状态,不能作为稳定的固定长度协议头。
协议 Header 应使用明确的整数和字节字段逐项编码。
3. 不要从零读取
错误:
int length = input.getInt(16);
正确:
int start = input.readerIndex();
int length =
input.getInt(start + 16);
累积缓冲区经过多轮消费后,当前帧起点不一定为 0。Netty 官方也明确要求绝对读取时以当前 readerIndex 为基准。
4. 不要提前消费
错误:
RpcHeader header = readHeader(input);
if (!bodyComplete(input, header)) {
return;
}
正确思路:
先查看
后判断
完整后再消费
5. 不要共享实例
ByteToMessageDecoder 的子类不能标注为 @Sharable。它维护与当前连接相关的累积状态,因此通常应为每条 Channel 创建独立实例。
正确:
@Override
protected void initChannel(
SocketChannel channel
) {
channel.pipeline().addLast(
new RpcFrameDecoder()
);
}
不要创建一个全局解码器实例,然后放入所有 Channel。
6. 注意 ByteBuf 生命周期
如果调用:
ByteBuf frame =
input.readBytes(length);
会创建新的 ByteBuf。
若它既没有加入 output,也没有被正确释放,就可能泄漏。Netty 官方文档特别提醒,readBytes(int) 返回的新 Buffer 必须进入输出列表或被释放;也可以根据生命周期选择切片方法。
简单场景中,读取到普通 byte[] 可以降低引用计数管理复杂度:
byte[] body =
new byte[bodyLength];
input.readBytes(body);
十八、常见误区
1. 一次 write 对应一次 read
错误。
TCP 传输的是有序字节流,不保留应用写入边界。
2. TCP 把完整消息弄坏了
错误。
TCP 保证字节可靠、有序到达;消息边界本来就属于应用协议。
3. 半包表示数据丢失
错误。
半包通常只是当前读取时,完整消息尚未全部进入应用缓冲区。
4. 多写一个 while 就能解决
不完整。
循环只能消费同一次读取中的多条完整消息,还必须保存跨读取的剩余数据。
5. Body 不够时直接 return 即可
只有在没有错误移动 readerIndex 时才成立。
如果 Header 已经被消费,直接返回会破坏下一次解码起点。
6. readInt 和 getInt 一样
错误。
readInt() 推进读指针,getInt(index) 不推进。
7. getInt(0) 永远正确
错误。
当前帧可能并不从底层缓冲区的索引 0 开始。
8. ByteToMessageDecoder 只调用 decode 一次
错误。
只要仍然能够继续解码,它可能在一次读取处理中反复调用 decode()。
9. out 是普通临时集合
不完整。
加入 out 的对象会作为解码结果继续向 Pipeline 下游传播。
10. requestId 可以解决粘包
错误。
requestId 解决请求响应关联,长度字段或其他帧协议解决消息边界。
11. 只在服务端解码即可
错误。
客户端接收响应时,同样面对 TCP 字节流,也需要拆帧和协议解码。
12. 一条连接不能并发请求
错误。
具备消息帧和 requestId 后,一条连接可以存在多个未完成请求。
13. 客户端线程顺序决定响应顺序
错误。
请求可以按任意线程调度顺序发出,服务端也可能乱序完成;RPC 应通过 requestId 关联,而不是依赖顺序。
14. Header 序列化长度固定
错误。
协议头应该采用确定的字段布局,而不是把某次 Java 序列化后的长度写死。
15. 解码器可以执行业务
不推荐。
解码器运行在 I/O 链路上,耗时操作会阻塞同一 EventLoop 管理的其他 Channel。
16. 把业务交给另一个 NioEventLoop 就行
不推荐。
其他 NioEventLoop 仍是 I/O 线程。阻塞业务更适合独立的业务执行器。
17. ByteToMessageDecoder 可以共享
错误。
其子类不能标注 @Sharable。
18. LengthFieldBasedFrameDecoder 会完成反序列化
错误。
它只根据长度字段切出完整帧,不理解 RPC 服务、参数或业务对象。
总结
并发 RPC 请求复用一条连接后,客户端可能连续发送:
Request A
Request B
Request C
TCP 将它们视为:
一段连续字节流
服务端一次读取可能得到:
一条完整消息
多条完整消息
半条消息
半条消息加几条完整消息
因此:
一次 channelRead
≠
一条 RPC 消息
错误解码器通常这样工作:
先读取 Header
↓
readerIndex 向后移动
↓
发现 Body 不完整
↓
直接 return
↓
下一次从 Body 中间解析 Header
正确解码器必须:
Header 不够
↓
不移动指针,等待
Header 完整
↓
使用 get 查看 Body 长度
完整帧不够
↓
不移动指针,等待
完整帧到达
↓
读取 Header 和 Body
输出 RpcMessage
ByteToMessageDecoder 负责:
累积多次读取的数据
保存未消费字节
重复调用 decode
把完整消息传给下游
LengthFieldBasedFrameDecoder 进一步负责:
根据长度字段
从 TCP 字节流中切出完整帧
推荐 Pipeline:
ByteBuf
↓
LengthFieldBasedFrameDecoder
↓
完整 RPC 帧
↓
RpcMessageDecoder
↓
RpcRequest / RpcResponse
↓
业务 Handler
同时要明确:
Body Length
负责找到消息边界
requestId
负责匹配请求和响应
二者缺一不可。
解码器只应该处理:
帧边界
协议校验
对象转换
耗时业务应交给专门的业务执行器,而不是堵住 NioEventLoop。
最后,可以用一句话概括整篇文章:
TCP 只负责可靠地传输字节,Netty 解码器负责把跨越多次读取的字节重新累积,并按照应用协议切回一条条完整消息。
评论区