侧边栏壁纸
  • 累计撰写 151 篇文章
  • 累计创建 21 个标签
  • 累计收到 3 条评论

目 录CONTENT

文章目录

TCP 为什么会拆开 RPC 消息?Netty 如何拼回一条完整消息?

YaFuX
2026-06-28 / 0 评论 / 0 点赞 / 2 阅读 / 0 字
温馨提示:
部分素材来自网络,若不小心影响到您的利益,请联系我们删除。

前言

上一篇文章中,我们已经写出了一个最小 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 长度
负数检查
字段一致性检查

LengthFieldBasedFrameDecodermaxFrameLength 就是第一道边界。

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 解码器负责把跨越多次读取的字节重新累积,并按照应用协议切回一条条完整消息。

0

评论区