Java NIO - Netty 的一些简单案例及半包问题的演示

EchoServer 案例

服务端的实现

NettyEchoServer:功能极其简单,服务端读取客户端输入的数据,然后将数据直接回显到控制台。

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
import io.netty.bootstrap.ServerBootstrap;
import io.netty.buffer.ByteBuf;
import io.netty.channel.*;
import io.netty.channel.nio.NioIoHandler;
import io.netty.channel.socket.SocketChannel;
import io.netty.channel.socket.nio.NioServerSocketChannel;
import java.net.InetSocketAddress;
import java.nio.charset.StandardCharsets;

/**
* @author KJ
* @description NettyEchoServer
*/
public class NettyEchoServer {
private final int port;
ServerBootstrap b = new ServerBootstrap();

public NettyEchoServer(int port) {
this.port = port;
}

public void runServer() {
// 1. 创建反应器轮询组:bossGroup 负责接收连接,workerGroup 负责处理具体的 I/O 读写
EventLoopGroup bossLoopGroup = new MultiThreadIoEventLoopGroup(1, NioIoHandler.newFactory());
EventLoopGroup workerLoopGroup = new MultiThreadIoEventLoopGroup(NioIoHandler.newFactory());
try {
// 2. 设置反应器轮询组
b.group(bossLoopGroup, workerLoopGroup);

// 3. 设置通道类型:使用 NIO 的服务端 TCP 通道
b.channel(NioServerSocketChannel.class);

// 4. 设置监听端口
b.localAddress(new InetSocketAddress(port));

// 5. 设置通道选项 (例如:允许端口复用,设置 TCP 积压队列大小)
b.option(ChannelOption.SO_BACKLOG, 128)
.childOption(ChannelOption.SO_KEEPALIVE, true)
.childOption(ChannelOption.TCP_NODELAY, true);

// 6. 装配子通道流水线 (Pipeline)
b.childHandler(new ChannelInitializer<SocketChannel>() {
@Override
protected void initChannel(SocketChannel ch) {
// 向子通道流水线添加共享的 Handler 实例
ch.pipeline().addLast(NettyEchoServerHandler.INSTANCE);
}
});

// 7. 绑定服务器,并异步地阻断直到绑定完成
ChannelFuture f = b.bind().sync();
System.out.println("Echo 服务器启动成功,监听端口:" + port);

// 8. 等待通道关闭的异步任务结束
f.channel().closeFuture().sync();
} catch (InterruptedException e) {
e.printStackTrace();
} finally {
// 9. 优雅关闭反应器轮询组,释放所有资源
bossLoopGroup.shutdownGracefully();
workerLoopGroup.shutdownGracefully();
}
}


/**
* @ChannelHandler.Sharable 用来标注一个Handler实例可以被多个通道安全地共享。
* 多个通道的流水线可以加入同一个Handler实例。这种共享操作,Netty 默认是不允许的。
* 很多应用场景都需要Handler实例能共享。例如,一个服务器处理十万以上的通道,如果
* 每一个通道都新建很多重复的Handler实例,就会浪费很多宝贵的空间,降低了服务器的性能。
* 所以,如果在Handler实例中没有与特定通道强相关的数据或者状态,建议设计成共享模式。
*
* ChannelHandlerAdapter.isSharable() 可以判断一个Handler是否为可共享。
* 如果其对应的实现加上了@Sharable注解,那么这个方法将返回 true。
*
* NettyEchoServerHandler没有保存与任何通道连接相关的数据,也没有内部的其他数据需要保存。
* 所以,该处理器不仅仅可以用来共享,而且不需要做任何同步控制。这里为它加上了@Sharable注解,
* 表示可以共享。更进一步,这里还设计为了INSTANCE静态实例,所有的通道直接使用这个实例即可。
*
*/
@ChannelHandler.Sharable // 未加该注解,试图将同一个Handler实例添加到多个ChannelPipeline,则会抛异常。
static class NettyEchoServerHandler extends ChannelInboundHandlerAdapter {
public static final NettyEchoServerHandler INSTANCE = new NettyEchoServerHandler();

/**
* 回显服务器处理器的逻辑分为两步:
*
* 第一步,读取从对端输入的数据。channelRead() 方法的 msg 参数的形参类型不是ByteBuf,而是Object,
* 这是由流水线的上一站决定的。一般而言,入站处理的流程是:Netty读取底层的二进制数据,填充到msg时,
* msg是ByteBuf类型,然后经过流水线,传入第一个入站处理器;每一个节点处理完后,将自己的处理结果作
* 为msg参数不断向后传递。因此,msg参数的形参类型只能是Object类型。第一个入站处理器的channelRead
* 方法的msg类型绝对是ByteBuf类型,因为它是Netty读取到的ByteBuf数据包。另外,从Netty 4.1开始,
* ByteBuf 的默认类型是 Direct ByteBuf。注意,Java不能直接访问Direct ByteBuf内部的数据,必须
* 通过调用 getBytes()、readBytes() 等方法将数据读入Java数组中才能继续进行处理。
*
* 第二步,将数据写回客户端。这一步很简单,直接复用前面的msg实例即可。不过要注意,如果上一步调用的是
* readBytes() 方法,那么这一步就不能直接将msg写回了,因为数据已经被readBytes()方法读完了。幸好,
* 上一步调用的读数据方法是 getBytes(),它不影响 ByteBuf 的数据指针,因此可以继续使用。这里除了调用
* ctx.writeAndFlush()方法把msg数据写回客户端之外,也可调用通道的ctx.channel().writeAndFlush()
* 方法发送数据。这两种方法在这里的效果是一样的,因为这个流水线上没有任何出站处理器。
*
* 注:假设你的 Pipeline 顺序是:
* Head ⇄ Encoder_A ⇄ Encoder_B ⇄ Your_Handler ⇄ Tail
* 调用 ctx.writeAndFlush():数据流向是 Your_Handler ➔ Encoder_A ➔ Head,它跳过了 Encoder_B。
* 调用 ctx.channel().writeAndFlush():数据流向是 Tail ➔ Encoder_B ➔ Encoder_A ➔ Head,它经过所有处理器。
*/
@Override
public void channelRead(ChannelHandlerContext ctx, Object msg) {
ByteBuf in = (ByteBuf) msg;

// 根据 ByteBuf 的存储属性判断是堆内存还是直接内存
System.out.println("msg type: " + (in.hasArray() ? "堆内存" : "直接内存"));

int len = in.readableBytes();
byte[] arr = new byte[len];
in.getBytes(0, arr);

System.out.println("server received: " + new String(arr, StandardCharsets.UTF_8));
System.out.println("写回前,msg.refCnt:" + in.refCnt());

// 零拷贝机制:直接写回原始数据包,Netty 会在发送完成后自动 release
ChannelFuture f = ctx.writeAndFlush(msg); // 从当前位置开始向 Header(头部)方向传播;而 ctx.channel().writeAndFlush() 从尾部开始向头部传播

f.addListener((ChannelFuture future) -> {
// 回调中再次查看引用计数(通常此时由于发送完毕已减 1)
System.out.println("写回任务状态:" + (future.isSuccess() ? "成功\n" : "失败\n"));
// 注意:在异步完成后访问计数仅用于观察,不应再进行读写操作
});
}

@Override
public void exceptionCaught(ChannelHandlerContext ctx, Throwable cause) {
// 发生异常时关闭连接
cause.printStackTrace();
ctx.close();
}
}

public static void main(String[] args) {
new NettyEchoServer(9000).runServer();
}
}


客户端的实现

NettyEchoClient

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
import io.netty.bootstrap.Bootstrap;
import io.netty.buffer.ByteBuf;
import io.netty.buffer.PooledByteBufAllocator;
import io.netty.channel.*;
import io.netty.channel.nio.NioIoHandler;
import io.netty.channel.socket.SocketChannel;
import io.netty.channel.socket.nio.NioSocketChannel;
import java.nio.charset.StandardCharsets;
import java.util.Scanner;

/**
* @author KJ
* @description NettyEchoClient
*/
public class NettyEchoClient {
private final int serverPort;
private final String serverIp;
Bootstrap b = new Bootstrap();

public NettyEchoClient(String ip, int port) {
this.serverPort = port;
this.serverIp = ip;
}

/**
* 客户端在成功连接到服务端后不断循环获取控制台的输入,通过与服务端之间的连接通道发送到服务器。
*/
public void runClient() {
// 创建反应器轮询组
EventLoopGroup workerLoopGroup = new MultiThreadIoEventLoopGroup(NioIoHandler.newFactory());
try {
// 1.设置反应器轮询组
b.group(workerLoopGroup);
// 2.设置nio类型的通道
b.channel(NioSocketChannel.class);
// 3.设置监听端口
b.remoteAddress(serverIp, serverPort);
// 4.设置通道的参数
b.option(ChannelOption.ALLOCATOR, PooledByteBufAllocator.DEFAULT);
// 5.装配子通道流水线
b.handler(new ChannelInitializer<SocketChannel>() {
// 有连接到达时会创建一个通道
protected void initChannel(SocketChannel ch) {
// 管理子通道中的Handler
// 向子通道流水线添加一个Handler
ch.pipeline().addLast(NettyEchoClientHandler.INSTANCE);
}
});

ChannelFuture future = b.connect();
future.addListener((ChannelFuture futureListener) -> {
if (futureListener.isSuccess()) {
System.out.println("EchoClient客户端连接成功!");
} else {
System.out.println("EchoClient客户端连接失败!");
}
});
future.sync(); // 阻塞,直到连接成功

Channel channel = future.channel();
Scanner scanner = new Scanner(System.in);
System.out.println("请输入发送内容:");
while (scanner.hasNext()) {
// 获取输入的内容
String next = scanner.next();
byte[] bytes = next.getBytes(StandardCharsets.UTF_8);
// 发送ByteBuf
ByteBuf buffer = channel.alloc().buffer();
buffer.writeBytes(bytes);
channel.writeAndFlush(buffer);
System.out.println("请输入发送内容:");
}
} catch (Exception e) {
e.printStackTrace();
} finally {
// 优雅关闭EventLoopGroup,
// 释放掉所有资源,包括创建的线程
workerLoopGroup.shutdownGracefully();
}
}

@ChannelHandler.Sharable
static class NettyEchoClientHandler extends ChannelInboundHandlerAdapter {
public static final NettyEchoClientHandler INSTANCE = new NettyEchoClientHandler();

@Override // 入站处理方法
public void channelRead(ChannelHandlerContext ctx, Object msg) {
ByteBuf byteBuf = (ByteBuf) msg;
int len = byteBuf.readableBytes();
byte[] arr = new byte[len];
byteBuf.getBytes(0, arr);
System.out.println("client received: " + new String(arr, StandardCharsets.UTF_8));
// 释放ByteBuf的两种方法
// 方法一:手动释放ByteBuf
byteBuf.release();
// 方法二:调用父类的入站方法,将msg向后传递
// super.channelRead(ctx,msg);
}
}

public static void main(String[] args) {
NettyEchoClient nettyEchoClient = new NettyEchoClient("127.0.0.1", 9000);
nettyEchoClient.runClient();
}
}


半包问题的复现

问题演示

改造一下前面的 NettyEchoClient 实例,通过循环的方式向 NettyEchoServer 回显服务器写入大量的 ByteBuf,然后看看实际的服务器响应结果。注意:服务器类不需要改造,直接使用之前的回显服务器即可。改造好的客户端类——叫 NettyDumpSendClient。在客户端建立连接成功之后,使用一个 for 循环不断通过通道向服务端发送ByteBuf, 一直写到1000次,这些ByteBuf的内容相同,都是相同的字符串内容。

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
public class NettyDumpSendClient {
private final String serverIp;
private final int serverPort;
Bootstrap b = new Bootstrap();

public NettyDumpSendClient(String ip, int port) {
this.serverPort = port;
this.serverIp = ip;
}

/**
* 客户端在成功连接到服务端后不断循环获取控制台的输入,通过与服务端之间的连接通道发送到服务器。
*/
public void runClient() {
// 创建反应器轮询组
EventLoopGroup workerLoopGroup = new MultiThreadIoEventLoopGroup(NioIoHandler.newFactory());
try {
// 1.设置反应器轮询组
b.group(workerLoopGroup);
// 2.设置nio类型的通道
b.channel(NioSocketChannel.class);
// 3.设置监听端口
b.remoteAddress(serverIp, serverPort);
// 4.设置通道的参数
b.option(ChannelOption.ALLOCATOR, PooledByteBufAllocator.DEFAULT);
// 5.组装处理流水线
b.handler(new ChannelInitializer<NioSocketChannel>() {
@Override
protected void initChannel(NioSocketChannel ch) {
// 哪怕现在逻辑为空,也必须初始化 Pipeline
// 后续如果你有入站处理逻辑,可以在这里 addLast
ch.pipeline().addLast(NettyEchoClientHandler.INSTANCE);
}
});

ChannelFuture future = b.connect();
future.addListener((ChannelFuture futureListener) -> {
if (futureListener.isSuccess()) {
System.out.println("NettyDumpSendClient客户端连接成功!");
} else {
System.out.println("NettyDumpSendClient客户端连接失败!");
}
});
future.sync(); // 阻塞,直到连接成功

Channel channel = future.channel();
//发送大量的文字
String content = "密涅瓦的猫头鹰在黄昏起飞。";
byte[] bytes = content.getBytes(StandardCharsets.UTF_8);
for (int i = 0; i < 1000; i++) {
// 发送ByteBuf
ByteBuf buffer = channel.alloc().buffer();
buffer.writeBytes(bytes);
channel.writeAndFlush(buffer);
}
} catch (Exception e) {
e.printStackTrace();
} finally {
// 优雅关闭EventLoopGroup,
// 释放掉所有资源,包括创建的线程
workerLoopGroup.shutdownGracefully();
}
}

@ChannelHandler.Sharable
static class NettyEchoClientHandler extends ChannelInboundHandlerAdapter {
public static final NettyEchoClient.NettyEchoClientHandler INSTANCE = new NettyEchoClient.NettyEchoClientHandler();

@Override // 入站处理方法
public void channelRead(ChannelHandlerContext ctx, Object msg) {
ByteBuf byteBuf = (ByteBuf) msg;
int len = byteBuf.readableBytes();
byte[] arr = new byte[len];
byteBuf.getBytes(0, arr);
System.out.println("client received: " + new String(arr, StandardCharsets.UTF_8));
byteBuf.release();
}
}

public static void main(String[] args) {
NettyDumpSendClient client = new NettyDumpSendClient("127.0.0.1", 9000);
client.runClient();
}
}

仔细观察服务端的控制台输出,可以看出存在三种类型的输出:

  • 读到一个完整的客户端输入ByteBuf。
  • 读到多个客户端的ByteBuf输入,但是“粘”在了一起。
  • 读到部分ByteBuf的内容,并且有乱码。

除了观察服务端的输出之外,再仔细观察客户端的输出,可以看 到客户端也存在以上三种类型的输出。对应于第1种情况接收到的完整的ByteBuf,这里称为“全包”。 对应于第2种情况,多个发送端的输入ByteBuf“粘”在了一起,这里 称为“粘包”。对应于第3种情况,一个发送过来的ByteBuf被“拆 开”接收,接收端读取到一个破碎的包,这里称为“半包”。为了简单起见,也可以将“粘包”的情况看成特殊的“半包”。 “粘包”和“半包”可以统称为传输的“半包问题”。


半包问题的本质

半包问题包含了 “粘包” 和 “半包” 两种情况:

  • 粘包:接收端(Receiver)收到一个ByteBuf,包含了发送端(Sender)的多个ByteBuf,发送端的多个ByteBuf在接收端“粘” 在了一起。
  • 半包:Receiver将Sender的一个ByteBuf“拆”开了收,收 到多个破碎的包。换句话说,Receiver收到了Sender的一个ByteBuf的 一小部分。

无论是粘包还是半包都不是一次正常的ByteBuf缓存区接收,具体如图所示:

粘包和半包的来源得从操作系统底层说起。我们知道,底层网络是以二进制字节报文的形式来传输数据的。读数据的过程大致为:当IO可读时,Netty 会从底层网络将二进制数据读到ByteBuf缓冲区中,再交给Netty程序转成Java POJO对象。写数据的过程大致为:编码器将一个Java类型的数据转换成底层能够传输的二进制ByteBuf缓冲数据。

在发送端 Netty 的应用层进程缓冲区中,程序以 ByteBuf 为单位来发送数据,但是到了底层操作系统内核缓冲区,底层会按照协议的规范对数据包进行二次封装,封装成传输层的协议报文,再进行发送。 在接收端收到传输层的二进制包后,首先复制到内核缓冲区,Netty读取ByteBuf时才复制到应用的用户缓冲区。在接收端,当Netty程序将数据从内核缓冲区复制到用户缓冲区的 ByteBuf时,问题来了:

  • 每次读取底层缓冲的数据容量是有限制的,当TCP内核缓冲区的数据包比较大时,可能会将一个底层包分成多次ByteBuf进行复制,进而造成用户缓冲区读到的是半包。
  • 当TCP内核缓冲区的数据包比较小时,一次复制的是不止一个内核缓冲区包,进而会造成用户缓冲区读到粘包。

如何解决呢?基本思路是,在接收端,Netty程序需要根据自定义协议将读取到的进程缓冲区ByteBuf在应用层进行二次组装,重新组装应用层的数据包。接收端的这个过程通常也称为分包或者拆包。在Netty中分包的方法主要有以下两种:

  • 可以自定义解码器分包器:基于 ByteToMessageDecoder 或者 ReplayingDecoder,定义自己的用户缓冲区分包器。
  • 使用 Netty 内置的解码器。例如,可以使用 Netty 内置的 LengthFieldBasedFrameDecoder 自定义长度数据包解码器对用户缓冲区 ByteBuf 进行正确的分包。


自定义协议解决粘包和拆包问题

在 Netty 中,虽然官方提供了像 LengthFieldBasedFrameDecoder(自定义长度解码器)这样的现成组件,但为了彻底理解粘包、拆包的底层原理,最好的办法是自己设计一个最简版的“定长”或“长度字段”自定义协议。下面我们设计一个业界最常用的 [Length] + [Body](长度 + 内容) 自定义协议:

  • 前 4 个字节(int):代表后面真实数据的字节长度。
  • 后面的字节:代表真实的文本内容。

当 Netty 读取到前 4 个字节时,就知道后面还要等多少字节才是一条完整的消息。这样无论底层网络怎么粘包(多条合在一起)或拆包(一条碎成多次),我们的协议都能精准将其还原。


核心协议编解码器

由于服务端和客户端都需要对该协议进行编码和解码,我们直接编写两个通用的处理器。

编码器:MyProtocolEncoder。将要发送的字符串转换成 [4字节长度] + [真实数据] 的二进制流。

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
import io.netty.buffer.ByteBuf;
import io.netty.channel.ChannelHandlerContext;
import io.netty.handler.codec.MessageToByteEncoder;
import java.nio.charset.StandardCharsets;

public class MyProtocolEncoder extends MessageToByteEncoder<String> {

@Override
protected void encode(ChannelHandlerContext channelHandlerContext, String msg, ByteBuf out) throws Exception {
if (msg == null) return;

// 1. 将字符串转为字节数组
byte[] bytes = msg.getBytes(StandardCharsets.UTF_8);

// 2. 写入 4 字节的长度标头
out.writeInt(bytes.length);

// 3. 写入真实的业务数据
out.writeBytes(bytes);
}
}

解码器:MyProtocolDecoder。负责在网络流中精准裁剪出一条条完整的消息。

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
import io.netty.buffer.ByteBuf;
import io.netty.channel.ChannelHandlerContext;
import io.netty.handler.codec.ByteToMessageDecoder;
import java.nio.charset.StandardCharsets;
import java.util.List;

public class MyProtocolDecoder extends ByteToMessageDecoder {

@Override
protected void decode(ChannelHandlerContext ctx, ByteBuf in, List<Object> out) throws Exception {
// 1. 如果可读字节少于 4 个,说明连“长度标头”都没收齐,继续等待
if (in.readableBytes() < 4) {
return;
}

// 2. 标记当前读取光标(用于数据不够时回滚)
in.markReaderIndex();

// 3. 读取 4 字节的长度
int length = in.readInt();

// 4. 判断剩余可读字节是否足够这条消息的长度
if (in.readableBytes() < length) {
// 数据还没收齐,回滚光标到刚才标记的位置,等待下次数据流入
in.resetReaderIndex();
return;
}

// 5. 数据够了,读取真实内容
byte[] bytes = new byte[length];
in.readBytes(bytes);

// 6. 将解析出的干净字符串传递给后续的业务 Handler
String content = new String(bytes, StandardCharsets.UTF_8);
out.add(content);
}
}


服务端主类

ProtocolServer

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
import io.netty.bootstrap.ServerBootstrap;
import io.netty.channel.*;
import io.netty.channel.nio.NioEventLoopGroup;
import io.netty.channel.socket.SocketChannel;
import io.netty.channel.socket.nio.NioServerSocketChannel;

public class ProtocolServer {
public static void main(String[] args) throws InterruptedException {
EventLoopGroup boss = new NioEventLoopGroup(1);
EventLoopGroup worker = new NioEventLoopGroup();
try {
ServerBootstrap b = new ServerBootstrap()
.group(boss, worker)
.channel(NioServerSocketChannel.class)
.childHandler(new ChannelInitializer<SocketChannel>() {
@Override
protected void initChannel(SocketChannel ch) {
ChannelPipeline pipeline = ch.pipeline();
// 挂载自定义的协议解密器和加密器
pipeline.addLast(new MyProtocolDecoder());
pipeline.addLast(new MyProtocolEncoder());
// 业务 Handler
pipeline.addLast(new SimpleChannelInboundHandler<String>() {
@Override
protected void channelRead0(ChannelHandlerContext ctx, String msg) {
System.out.println("【服务端收到真实消息】: " + msg);
ctx.writeAndFlush("服务器响应: " + msg);
}
});
}
});
ChannelFuture f = b.bind(9999).sync();
System.out.println("自定义协议服务端已启动,监听 9999...");
f.channel().closeFuture().sync();
} finally {
boss.shutdownGracefully();
worker.shutdownGracefully();
}
}
}


客户端主类

ProtocolClient:为了验证自定义协议真的防粘包,我们让客户端在 for 循环中不间断、无延迟地瞬间发送 100 条消息。如果不用自定义协议,这 100 条消息必然会在底层网络缓冲区中揉成一团(粘包)。

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
import io.netty.bootstrap.Bootstrap;
import io.netty.channel.*;
import io.netty.channel.nio.NioEventLoopGroup;
import io.netty.channel.socket.SocketChannel;
import io.netty.channel.socket.nio.NioSocketChannel;

public class ProtocolClient {
public static void main(String[] args) throws InterruptedException {
EventLoopGroup group = new NioEventLoopGroup();
try {
Bootstrap b = new Bootstrap()
.group(group)
.channel(NioSocketChannel.class)
.handler(new ChannelInitializer<SocketChannel>() {
@Override
protected void initChannel(SocketChannel ch) {
ChannelPipeline pipeline = ch.pipeline();
// 客户端也挂载相同的编解码器
pipeline.addLast(new MyProtocolDecoder());
pipeline.addLast(new MyProtocolEncoder());

pipeline.addLast(new SimpleChannelInboundHandler<String>() {
@Override
public void channelActive(ChannelHandlerContext ctx) {
System.out.println("【客户端】连接成功,开始极速无间隔群发 100 条消息模拟粘包环境...");
// 瞬间冲刷 100 条消息,极其容易引发物理网络粘包
for (int i = 1; i <= 100; i++) {
ctx.writeAndFlush("Hello Netty, 这是第 [" + i + "] 条消息!");
}
}

@Override
protected void channelRead0(ChannelHandlerContext ctx, String msg) {
System.out.println("【客户端收到回显】: " + msg);
}
});
}
});
ChannelFuture f = b.connect("127.0.0.1", 9999).sync();
f.channel().closeFuture().sync();
} finally {
group.shutdownGracefully();
}
}
}


测试验证

运行服务端,再运行客户端。观察服务端和客户端的控制台,输出结果一定是整整齐齐的 100 行,没有任何两条消息粘在一起,也没有任何一条消息中途断开:

1
2
3
4
【服务端收到真实消息】: Hello Netty, 这是第 [1] 条消息!
【服务端收到真实消息】: Hello Netty, 这是第 [2] 条消息!
...
【服务端收到真实消息】: Hello Netty, 这是第 [100] 条消息!
1
2
3
4
5
6
【客户端收到回显】: 服务器响应: Hello Netty, 这是第 [1] 条消息!
【客户端收到回显】: 服务器响应: Hello Netty, 这是第 [2] 条消息!
...
【客户端收到回显】: 服务器响应: Hello Netty, 这是第 [100] 条消息!

ReplayingDecoder 的核心思想是:它内部使用了一个特殊的 ReplayingDecoderBuffer。当你在里面读取数据(比如 in.readInt())而底层字节又不够时,它会抛出一个特殊的 Signal 错误,拦截并自动帮你把光标回滚到本次解码开始前的状态,静静等待下次数据流入。

虽然客户端的 for 循环发送得极快,操作系统底层的 TCP 缓冲区可能一次性把 5 条甚至 10 条消息拼在一起塞给了服务端的网卡。但因为我们的管道最前方有 MyProtocolDecoder:

  • 它每次都雷打不动地先抠出 4 个字节,得知当前这条消息只有 30 个字节长。
  • 它就只裁剪出后面的 30 个字节丢给业务 Handler。
  • 至于缓冲区里剩下的其余字节,它会留在 ByteBuf 里面,等下一次循环继续如法炮制。这就从底层用数学边界彻底根治了网络粘包与拆包。


附1.使用 ReplayingDecoder 简化解码器

在刚才的 MyProtocolDecoder 中,我们为了防止拆包(数据没收齐),必须写大量的 in.readableBytes() < 4 和 in.resetReaderIndex() 这种防御性代码。对于开发者来说,你完全可以假定“网络数据已经百分之百收齐了”,闭着眼睛直接读。

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
import io.netty.buffer.ByteBuf;
import io.netty.channel.ChannelHandlerContext;
import io.netty.handler.codec.ReplayingDecoder;
import java.nio.charset.StandardCharsets;
import java.util.List;

// 继承 ReplayingDecoder,泛型 Void 代表不需要复杂的内部状态管理
public class MyProtocolDecoder extends ReplayingDecoder<Void> {

@Override
protected void decode(ChannelHandlerContext ctx, ByteBuf in, List<Object> out) throws Exception {
/*
* 不需要 if (in.readableBytes() < 4)
* 如果此时 4 个字节没收齐,Netty 底层会自动拦截、回滚光标,并退出本次循环
*/
int length = in.readInt();

/*
* 不需要 if (in.readableBytes() < length)
* 同样,如果后面的业务字节不够长,底层依然会自动安全回滚
*/
byte[] bytes = new byte[length];
in.readBytes(bytes);

// 轻松拿到干净的数据,抛给业务层
String content = new String(bytes, StandardCharsets.UTF_8);
out.add(content);
}
}

虽然 ReplayingDecoder 看起来像魔法一样好用,但天下没有免费的午餐。在生产环境的高并发压测下,它有两大致命缺陷:

  • 极端情况下的性能暴跌:如果一条消息有 1MB,由于网络原因,每次只挤过来 1KB 的碎包。ReplayingDecoder 在数据不够时,每次都会抛出异常回滚,并从头开始解码。这意味着为了这 1MB 的数据,它会把前面的数据重复解析上千次,导致 CPU 瞬间飙高。
  • 并不支持所有的 ByteBuf 操作:如果你在 decode 里调用了 in.indexOf()、in.getByte() 或者一些不支持的特定指针操作,它内部那层包装的 ReplayingDecoderBuffer 会直接抛出 UnsupportedOperationException 异常。

所以,最佳的实践建议是:

  • 小项目:用 ReplayingDecoder 快速实现,代码赏心悦目。
  • 高并发/高吞吐大型项目:老老实实像我们上一版本那样,继承 ByteToMessageDecoder 亲手用 markReaderIndex() 搞定边界;或者直接使用 Netty 官方推荐的工业级定海神针 —— LengthFieldBasedFrameDecoder(长度字段预检解码器)。


附2.万能解码器

如果要用 Netty 官方最强大的工具来彻底解决粘包和拆包,那必然是 LengthFieldBasedFrameDecoder(万能长度字段解码器)。它是工业界长连接协议的 “定海神针”。它功能极度强大,通过配置几个核心的参数偏移量,就能几乎兼容所有业界主流的自定义二进制协议(包括 Dubbo、RocketMQ 甚至 HTTP2 的帧解析)。下面我们来看如何用它来替换掉我们手写的 MyProtocolDecoder。


五个顶配参数

LengthFieldBasedFrameDecoder 的构造函数有 5 个最核心的参数,理解它们是掌握 Netty 协议开发的关键。我们以最常见的 [Length] + [Body] 协议为例:

1
2
3
4
5
6
7
public LengthFieldBasedFrameDecoder(
int maxFrameLength,
int lengthFieldOffset,
int lengthFieldLength,
int lengthAdjustment,
int initialBytesToStrip
)
  • maxFrameLength(最大帧长度):如果单条消息超过这个长度(例如 10MB),解码器会直接抛出异常并关闭连接,防止恶意客户端发送超大报文把服务器内存撑爆。
  • lengthFieldOffset(长度字段偏移量):长度字段从哪个字节开始。如果你的协议开头就是长度,那就是 0。如果开头有 2 字节的 Magic Number(魔数),那它就是 2。
  • lengthFieldLength(长度字段自身占用的字节数):你的长度是用什么类型存的?byte 是 1,short 是 2,int 是 4,long 是 8。
  • lengthAdjustment(长度补偿值):你的“长度”包含标头自身吗?如果长度值只代表 Body 的大小,这里填 0;如果长度值包含了【标头 + Body】的总大小,这里需要填负数进行补偿。
  • initialBytesToStrip(跳过/裁剪的字节数):解析完成后,传给下一个业务 Handler 的数据要不要剥掉前面的长度标头?如果填 4,业务层拿到的就直接是干净的 Body 字节,连手里的 in.readInt() 都可以省了!


重写上述案例

我们将刚才手写的 MyProtocolDecoder 废弃,直接在服务端和客户端的流水线(Pipeline)上挂载 Netty 官方的 LengthFieldBasedFrameDecoder。

升级服务端流水线:ProtocolServer

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
// ... ServerBootstrap 配置保持不变 ...
.childHandler(new ChannelInitializer<SocketChannel>() {
@Override
protected void initChannel(SocketChannel ch) {
ChannelPipeline pipeline = ch.pipeline();

/*
* 核心:配置官方长度解码器彻底解决粘包/拆包
* 参数配置含义:
* 1. maxFrameLength: 1024 * 1024 (最大支持 1MB 的消息)
* 2. lengthFieldOffset: 0 (长度字段就在最开头)
* 3. lengthFieldLength: 4 (用 int 存储长度,占 4 字节)
* 4. lengthAdjustment: 0 (长度值只代表 Body 大小,不需要补偿)
* 5. initialBytesToStrip: 4 (神技:自动剥离前 4 个字节的长度标头!)
*/
pipeline.addLast(new LengthFieldBasedFrameDecoder(
1024 * 1024, 0, 4, 0, 4
));

/*
* 注意:因为上面配置了 initialBytesToStrip = 4,
* 此时流到下一个 Handler 的 ByteBuf 里,就已经没有【长度标头】了,只剩下纯粹的 Body 二进制。
* 所以我们不需要自定义解码器了,直接挂一个 Netty 自带的 StringDecoder 就能转成字符串!
*/
pipeline.addLast(new StringDecoder(StandardCharsets.UTF_8));

// 编码器依然用我们之前写的(向外发送时,自动加上 4 字节的长度头)
pipeline.addLast(new MyProtocolEncoder());

// 业务 Handler
pipeline.addLast(new SimpleChannelInboundHandler<String>() {
@Override
protected void channelRead0(ChannelHandlerContext ctx, String msg) {
// 此时 msg 已经是被 StringDecoder 翻译好的干净字符串了
System.out.println("【LengthField 收到干净消息】: " + msg);
ctx.writeAndFlush("服务器响应: " + msg);
}
});
}
});

升级客户端流水线:ProtocolClient

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
// ... Bootstrap 配置保持不变 ...
.handler(new ChannelInitializer<SocketChannel>() {
@Override
protected void initChannel(SocketChannel ch) {
ChannelPipeline pipeline = ch.pipeline();

// 客户端入站也同样挂载官方解码器,剥离 4 字节标头
pipeline.addLast(new LengthFieldBasedFrameDecoder(
1024 * 1024, 0, 4, 0, 4
));
// 自动将解完边界的 Body 字节转为字符串
pipeline.addLast(new StringDecoder(StandardCharsets.UTF_8));

// 发送时用的编码器
pipeline.addLast(new MyProtocolEncoder());

pipeline.addLast(new SimpleChannelInboundHandler<String>() {
@Override
public void channelActive(ChannelHandlerContext ctx) {
System.out.println("【客户端】连接成功,开始极速群发 100 条消息...");
for (int i = 1; i <= 100; i++) {
ctx.writeAndFlush("Hello LengthField! 这是第 [" + i + "] 条消息");
}
}

@Override
protected void channelRead0(ChannelHandlerContext ctx, String msg) {
System.out.println("【客户端收到回显】: " + msg);
}
});
}
});

使用 LengthFieldBasedFrameDecoder 配合 StringDecoder,我们甚至连一行自定义解码代码都不用写,就完成了对自定义协议的高性能解析。相比于我们自己写 in.readableBytes() < 4 或使用 ReplayingDecoder,官方这个组件强在哪里?

  • 零复制与极致内存优化:它内部经过了极高并发的工业打磨,在裁剪字节(initialBytesToStrip)时,使用的是 ByteBuf.slice() 虚拟切片技术,不会在内存中发生真实的字节数组拷贝,性能拉满。
  • 完美抵御内存撑爆攻击:由于有 maxFrameLength 的保护,当黑客恶意向你的端口发送没有边界的垃圾流量时,它在读取到设定的阈值后就会立刻切断连接并报警,保护后面的业务代码不被拖垮。