Netty-202509120209
Netty
一、Netty概述
1. 什么是Netty
Netty是一款基于Java NIO的网络编程、高性能、异步事件驱动的网络应用框架。它的核心
目标是让开发者能够轻松、快速地构建高性能、高可靠、可扩展的网络服务器和客户端程
序。它通过精巧的抽象,极大地简化了TCP/UDP套接字服务器等网络编程的复杂性。
2. Netty的优势
高性能
异步非阻塞I/O:基于NIO,用少量线程即可处理海量连接,资源利用率极高。
零拷贝:支持在文件传输等场景中避免不必要的内存拷贝,显著提升性能。
内存池化:通过重用ByteBuf对象,减少GC压力,提升内存分配效率。
设计与拓展性
模块化与组件化:所有核心组件(如ChannelHandler)都是可插拔的,用户可以
像搭积木一样自定义协议栈。
丰富的协议支持:内置了HTTP、WebSocket、SSL/TLS等多种协议编解器,并易
于扩展自定义协议。
易用性
强大的API:提供了比JDK NIO更友好、更健壮的API,将复杂的Selector、Buffer
操作封装到底层。
详尽的文档和社区:拥有活跃的社区和丰富的案例,学习成本相对较低。
稳定性和可靠性
历经多年发展,被众多知名项目(如Dubbo、RocketMQ、Elasticsearch等)选作
底层通信框架,久经生产环境考验
跨平台支持
基于 Java:Netty 基于 Java 构建,具有良好的跨平台性,可以在任何支持 Java
的操作系统上运行。
3. 使用场景
高性能RPC框架:构建分布式服务的底层通信层,如Apache Dubbo、gRPC-Java。
Netty.pdf
实时通信系统:游戏服务器、弹幕系统、IM聊天应用。
代理与网关:API网关、反向代理、负载均衡器。
大数据处理:节点间高速数据传输。
二、IO模型演进与Netty线程模型
1. IO模型对比
1.1 BIO(阻塞IO)
同步阻塞IO。服务器模式为 一个连接对应一个线程。
当连接数暴涨时,线程数线性增长,巨大的线程上下文切换开销和内存消耗将成为系统瓶
颈,无法应对万级甚至十万级连接。
1.2 NIO(非阻塞IO)
同步非阻塞IO。服务器模式为 一个线程处理多个连接,其核心是IO多路复用技术。
线程通过一个Selector(选择器)轮询注册在其上的多个连接,只有当某个连接发生真正的
读写事件时,线程才进行处理。这使得单个线程可以高效管理成千上万个连接。
IO多路复用技术的演进:
select/poll: 早期实现,通过线性扫描所有文件描述符来检测就绪事件,时间复杂度
为O(n)。select还有最大描述符数量限制(1024)。
epoll (Linux): 现代高效实现,采用事件驱动方式,通过回调机制直接通知就绪事件,
时间复杂度为O(1),性能不受连接数增长影响,是Netty在Linux下高性能的基石。
1.3 AIO(异步IO)
异步非阻塞IO。应用程序发起IO操作后立即返回,由操作系统完成所有IO操作后再通知应
用程序处理。理论模型更先进,但Netty并未采用,原因如下:
特性
BIO(阻塞IO)
NIO (非阻塞IO)
Netty (基于NIO)
线程模型
1连接1线程
1线程处理多连接
主从Reactor多线程
阻塞方式
同步并阻塞
同步非阻塞
同步非阻塞(异步AP
I)
编程复杂度
简单
非常复杂
简单(高度封装)
核心机制
流(Stream)
选择器(Selector)
事件循环(EventLoop)
适用场景
低并发、短连接
高并发、长连接
超高并发网络通信
Linux实现不佳: 在Linux上,AIO底层实现仍使用epoll,并未实现真正的异步IO,性能
优势不明显。
架构不匹配: Netty的整体架构是Reactor模型(同步),而AIO是Proactor模型(异
步),混合使用会使架构变得复杂和混乱。
成熟度与可控性: NIO模型更成熟,且给予开发者更高的可控性,Netty在其上进行的
封装已经足够优秀。
2. Netty线程模型
主从Reactor多线程模型。Netty的线程模型是其架构的核心,它清晰地将事件处理职责分
离,以实现高并发和高性能。
核心角色与分工:
BossGroup (主Reactor):
职责:专责连接接入。内部包含一个或多个NioEventLoop,监听服务器端口的Ser
verSocketChannel,接受(accept)客户端的连接请求。
行为:成功建立连接后,将新创建的SocketChannel注册到WorkerGroup中的一个
NioEventLoop上。
WorkerGroup (从Reactor):
职责:专责I/O处理。内部包含多个NioEventLoop,处理由主Reactor分配下来的
所有连接的数据读写、编解码等网络I/O操作。
行为:每个NioEventLoop都独立运行,负责处理注册在其上的所有SocketChannel
的I/O事件。
NioEventLoop:事件循环引擎
本质:一个执行无限循环的线程,每个循环主要执行两大任务:
select:检查其负责的所有Channel是否有已就绪的I/O事件。
processSelectedKeys:处理所有就绪的I/O事件;并runAllTasks:执行队列中
的异步任务。
设计精髓:串行化处理与线程绑定
Netty通过线程局部性(Thread Local) 原则,保证一个Channel在其整个生命周
期内,其所有的I/O事件都由同一个NioEventLoop线程处理。这从根本上避免
了多线程并发操作同一个Channel所带来的锁竞争和上下文切换开销,实现了
无锁化设计,极大提升了性能。
三、Netty核心组件
1. Channel:网络连接的抽象
代表一个开放的网络连接(如TCP Socket),是执行所有网络I/O操作(读、写、连接、绑
定)的顶层接口。每个Channel都关联着一个唯一的ChannelPipeline和EventLoop。
2. EventLoop & EventLoopGroup:事件处理引擎
EventLoop:事件循环。负责处理在其上注册的所有Channel的I/O事件和执行异步任
务。其本质是一个线程执行一个循环。
EventLoopGroup:事件循环组。是多个EventLoop的集合,用于管理EventLoop的生命
周期并提供分配策略。
3. ChannelHandler 与 ChannelPipeline:处理责任链
ChannelPipeline:是一个由ChannelHandler实例组成的双向链表,是处理入站出站事
件的责任链。它为每个新创建的Channel分配一个独立的pipeline。
ChannelHandler:处理I/O事件或拦截I/O操作的单元。分为两类:
ChannelInboundHandler:处理入站事件(如数据到达、连接激活)。
ChannelOutboundHandler:处理出站事件(如数据写入、连接关闭)。
数据流规则:
入站数据沿Pipeline从头节点向尾节点流动,依次由InboundHandler处理。
出站数据沿Pipeline从尾节点向头节点流动,依次由OutboundHandler处理。
4. ByteBuf :字节数据容器
Netty提供的字节缓冲区,是对JDK ByteBuffer的全面增强。
读写索引分离:通过readerIndex和writerIndex分离读写位置,无需调用flip()方法切换
模式。
池化(Pooling):通过重用已分配的ByteBuf实例,减少JVM垃圾回收频率和压力,提升
性能。
零拷贝支持:提供CompositeByteBuf,允许将多个ByteBuf逻辑组合为一个视图,避免
内存拷贝;支持文件传输的零拷贝机制。
5. Bootstrap:应用启动引导
ServerBootstrap:用于服务端,需配置主(BossGroup)从(WorkerGroup)两个Eve
ntLoopGroup。
Bootstrap:用于客户端,只需配置一个EventLoopGroup。
6. ChannelFuture:异步操作结果
Netty中所有I/O操作都是异步的。调用如write(), connect()等方法会立即返回一个ChannelF
uture。开发者可以通过为其添加GenericFutureListener来在操作完成(成功或失败)时获
取通知并执行回调逻辑,这是Netty异步编程模型的基础。
工作流程
服务端启动:ServerBootstrap 绑定端口,BossGroup 接受连接,将通道分配给 Worke
rGroup,WorkerGroup 的每个 EventLoop 负责监听并处理 I/O。
事件处理:数据到达 Channel 后,由 EventLoop 驱动,经过 ChannelPipeline 中的 Ch
annelHandler 链处理。
资源释放:连接关闭时,自动释放 ByteBuf 内存(当 ByteBuf 的 release() 被调用并且
引用计数减为 0 时,会自动释放它所占的内存),调用 ChannelPipeline 中每个 Channe
lHandler 的 handleRemoved() 方法清理资源。
四、粘包与半包
粘包和半包通常出现在网络编程中,尤其是在使用 TCP 协议进行数据传输时。它们是由于
TCP 是面向字节流的协议,数据发送和接收的边界并不一定和应用层数据的边界一致,从
而导致的问题。
1. 定义
粘包:多个小的数据包被合并在一起,接收方无法区分每个数据包的边界。通常是因为接
收端的缓冲区足够大、Nagle 算法等原因。
例如,发送方依次发送两个数据包:第一个包:abc、第二个包:def,如果发生粘包,接
收方可能会一次性接收到合并后的数据 abcdef,无法知道原来是两个独立的数据包。
原因:
TCP 是面向字节流的:TCP 协议在发送数据时,并没有定义数据包的边界。多个小的
数据包可能会被合并为一个大的数据包一起发送给接收端,接收端就无法正确识别每
个数据包的开始和结束。
接收方的缓冲区足够大:如果接收端的缓冲区(如 ByteBuf)很大,可以一次性接收多
个数据包而没有及时处理,就会导致多个报文被“粘”在一起。
Nagle 算法:Nagle 算法将小的数据包合并,以减少网络中的小包数。如果发送的数据
量比较小(比如 1 字节),Nagle 算法可能会将多个小包合并为一个大包来发送,导致
接收方一次接收到多个数据包,产生粘包。
半包:一个数据包被拆分成多个 TCP 包,接收方一次接收到的数据不完整。通常是由于接
收缓冲区不足、MSS 限制、滑动窗口等原因。
例如,发送方发送的数据包 abcdef,由于大小问题,接收方可能只接收到 abc 或 def 的一
部分数据,剩下的部分需要等待后续的包才能拼接完整。
原因:
接收方缓冲区不够大:如果接收方的缓冲区(如 ByteBuf)小于发送的数据包大小,接
收方无法一次性接收完整的数据包,只能接收部分数据。此时,剩余的数据就会继续
等待。
滑动窗口机制:TCP 协议有一个滑动窗口机制控制流量。如果接收方的窗口(缓冲区)
空间不够,它可能会分批接收数据。例如,接收方的窗口剩余空间仅有 128 字节,但
发送方发送的数据包大小为 256 字节,接收方就只能先接收 128 字节,剩余部分在窗
口足够时继续接收。
MSS 限制:MSS(最大报文段大小)是 TCP 协议中规定的最大单个数据包大小。发送
的数据如果超过了 MSS,TCP 会将数据切割成多个小包进行传输。接收方可能需要多
次接收才能拼接出完整的数据包。
2. 解决粘包和半包问题
2.1 固定消息长度
FixedLengthFrameDecoder 用来指定服务端每次接受消息的长度len,如果接受到的消息小
于 len,那么它会等待下个消息,并把这两个消息合并成 len 长度,然后发送;如果发送的
消息长度超过 len,那么这个消息会被切割,先发送 len 长度的消息。
所以,FixedLengthFrameDecoder 适合定长消息的场景,对于定长消息的场景下可以解决
粘包和半包问题。
public class Service {
private static final Logger log =
LoggerFactory.getLogger(Main.class);
public static void main(String[] args) throws InterruptedException
{
ServerBootstrap serverBootstrap = new ServerBootstrap();
NioEventLoopGroup boss = new NioEventLoopGroup(1);
NioEventLoopGroup worker = new NioEventLoopGroup();
try {
ChannelFuture bind = serverBootstrap
.group(boss, worker)
.channel(NioServerSocketChannel.class)
.childHandler(new
ChannelInitializer
@Override
protected void initChannel(NioSocketChannel
ch) throws Exception {
ChannelPipeline pip = ch.pipeline();
pip.addLast(new
FixedLengthFrameDecoder(10));
// <----- 指定每次接受的长度是10
pip.addLast(new
LoggingHandler(LogLevel.DEBUG));
}
}).bind(8888).sync();
bind.channel().closeFuture().sync();
} finally {
2.2 分割符
LineBasedFrameDecoder :可以根据 \n 来作为消息的分隔符,只有遇到 \n 时才会发
送和接受消息,同时,它可以设置一个最大消息长度,当消息长度超过这个值时会抛
出异常。
boss.shutdownGracefully();
worker.shutdownGracefully();
}
}
}
public class Service {
private static final Logger log =
LoggerFactory.getLogger(Main.class);
public static void main(String[] args) throws InterruptedException
{
ServerBootstrap serverBootstrap = new ServerBootstrap();
NioEventLoopGroup boss = new NioEventLoopGroup(1);
NioEventLoopGroup worker = new NioEventLoopGroup();
try {
ChannelFuture bind = serverBootstrap
.group(boss, worker)
.channel(NioServerSocketChannel.class)
// .childOption(ChannelOption.RCVBUF_ALLOCATOR, new
AdaptiveRecvByteBufAllocator(16, 16, 16))
.childHandler(new
ChannelInitializer
@Override
protected void initChannel(NioSocketChannel
ch) throws Exception {
ChannelPipeline pip = ch.pipeline();
pip.addLast(new
LineBasedFrameDecoder(16)); // <----- 指定最大长度
pip.addLast(new
LoggingHandler(LogLevel.DEBUG));
}
}).bind(8888).sync();
bind.channel().closeFuture().sync();
} finally {
boss.shutdownGracefully();
worker.shutdownGracefully();
}
DelimiterBasedFrameDecoder :跟 LineBasedFrameDecoder类似,只不过它可以指
定分隔符是什么。
}
}
public class Service {
private static final Logger log =
LoggerFactory.getLogger(Main.class);
public static void main(String[] args) throws InterruptedException
{
ServerBootstrap serverBootstrap = new ServerBootstrap();
NioEventLoopGroup boss = new NioEventLoopGroup(1);
NioEventLoopGroup worker = new NioEventLoopGroup();
try {
ChannelFuture bind = serverBootstrap
.group(boss, worker)
.channel(NioServerSocketChannel.class)
// .childOption(ChannelOption.RCVBUF_ALLOCATOR, new
AdaptiveRecvByteBufAllocator(16, 16, 16))
.childHandler(new
ChannelInitializer
@Override
protected void initChannel(NioSocketChannel
ch) throws Exception {
ChannelPipeline pip = ch.pipeline();
pip.addLast(new
DelimiterBasedFrameDecoder(16,
Unpooled.wrappedBuffer("A".getBytes())));
// <---- 指定最大长度和分割符
pip.addLast(new
LoggingHandler(LogLevel.DEBUG));
}
}).bind(8888).sync();
bind.channel().closeFuture().sync();
} finally {
boss.shutdownGracefully();
worker.shutdownGracefully();
}
}
}
2.3 LengthFieldBasedFrameDecoder
LengthFieldBasedFrameDecoder 通过指定每个消息前面携带一些额外的信息来解决粘包和
半包问题,是比较常用的处理粘包和半包的处理器
LengthFieldBasedFrameDecoder 是要求每个单独的消息前面都加上一个len,来指定 cont
ent 有多长(content是我们真正需要的消息)。当content长度小于len时,消息不会继续向
下传递,而是在 pip.addLast(new LengthFieldBasedFrameDecoder(128, 0, 4, 0, 4)); 这里等
待,直到content的长度等于len时,才会继续向下传递,所以 LengthFieldBasedFrameDec
oder 可以解决粘包和半包问题。
LengthFieldBasedFrameDecoder 接受5个参数,第一个参数是消息的最大长度,超过最大
长度会抛出异常。其余四个参数含义如下
五、使用实例
1. 导包
2. 创建NettyServer,用于启动服务端
public class NettyServer {
public static void main(String[] args) throws Exception {
// 创建 boss 线程组,处理 连接请求,个数代表有 几主,每个 主 都需要配置一
个单独的监听端口
NioEventLoopGroup bossGroup = new NioEventLoopGroup(1);
// 创建 worker 线程组,处理 具体业务,个数代表有
NioEventLoopGroup workerGroup = new NioEventLoopGroup();
try {
ServerB ootstrap serverBootstrap = new
ServerBootstrap();
// 创建 Server 启动器,配置必要参数:bossGroup, workerGroup,
channel, handler
serverBootstrap.group(bossGroup, workerGroup)
.channel(NioServerSocketChannel.class)
.childHandler(new ChannelInitializer
() {
@Override
protected void initChannel(SocketChannel
socketChannel) throws Exception {
3. 创建NettyServerHandler,实现具体业务
// 对workerGroup的SocketChannel设置 个性化业务
Handler
socketChannel.pipeline().addLast(new
NettyServerHandler());
}
});
System.out.println("Netty Server start ...");
ChannelFuture channelFuture =
serverBootstrap.bind(9000).sync();
// 给 channelFuture 添加监听器,监听是否启动成功
/*channelFuture.addListener(new ChannelFutureListener() {
@Override
public void operationComplete(ChannelFuture
channelFuture) throws Exception {
if (channelFuture.isSuccess()) {
System.out.println("启动成功");
} else {
System.out.println("启动失败");
}
}
});*/
channelFuture.channel().closeFuture().sync();
} finally {
bossGroup.shutdownGracefully();
workerGroup.shutdownGracefully();
}
}
}
/**
并回复消息
*/
public class NettyServerHandler extends ChannelInboundHandlerAdapter {
// 当客户端连接服务器完成就会触发该方法
@Override
public void channelActive(ChannelHandlerContext ctx) throws
Exception {
4. 创建NettyClient,启动客户端
System.out.println("客户端建立连接成功");
}
@Override
public void channelInactive(ChannelHandlerContext ctx) throws
Exception {
System.out.println("客户端断开连接");
}
// 读取客户端发送的数据
@Override
public void channelRead(ChannelHandlerContext ctx, Object msg)
throws Exception {
ByteBuf byteBuf = (ByteBuf) msg;
System.out.println("收到客户端的消息是:" +
byteBuf.toString(CharsetUtil.UTF_8));
}
// 数据读取完毕时触发该方法
@Override
public void channelReadComplete(ChannelHandlerContext ctx) throws
Exception {
ByteBuf buf =
Unpooled.copiedBuffer("HelloClient".getBytes(CharsetUtil.UTF_8));
ctx.writeAndFlush(buf);
}
// 处理异常, 一般是需要关闭通道
@Override
public void exceptionCaught(ChannelHandlerContext ctx, Throwable
cause) throws Exception {
cause.printStackTrace();
ctx.close();
}
}
public class NettyClient {
public static void main(String[] args) throws InterruptedException
{
NioEventLoopGroup group = new NioEventLoopGroup();
try {
5. 创建 NettyClientHandler,实现客户端业务逻辑
Bootstrap bootstrap = new Bootstrap();
// 创建 客户端启动器,配置必要参数
bootstrap.group(group)
.channel(NioSocketChannel.class)
.handler(new ChannelInitializer
@Override
protected void initChannel(SocketChannel ch)
throws Exception {
ch.pipeline().addLast(new
NettyClientHandler());
}
});
System.out.println("Netty Client start ...");
ChannelFuture channelFuture =
bootstrap.connect("127.0.0.1", 9000).sync();
channelFuture.channel().closeFuture().sync();
} finally {
/**
*
*/
public class NettyClientHandler extends ChannelInboundHandlerAdapter {
// 当客户端连接服务器完成就会触发该方法
@Override
public void channelActive(ChannelHandlerContext ctx) {
System.out.println("连接建立成功");
ByteBuf buf =
Unpooled.copiedBuffer("HelloServer".getBytes(CharsetUtil.UTF_8));
ctx.writeAndFlush(buf);
}
//当通道有读取事件时会触发,即服务端发送数据给客户端
@Override
public void channelRead(ChannelHandlerContext ctx, Object msg) {
ByteBuf buf = (ByteBuf) msg;
运行
先执行 NettyServer 的 main 方法,启动服务端,查看日志
执行 NettyClient 的 main 方法,启动客户端,查看日志
关闭 NettyClient 服务,查看日志
重新执行 NettyClient 的 main 方法,启动客户端,查看日志
System.out.println("收到服务端的消息:" +
buf.toString(CharsetUtil.UTF_8));
}
@Override
public void exceptionCaught(ChannelHandlerContext ctx, Throwable
) {