Netty 快速入门
Netty 快速入门
Netty 简介
Netty 是一款基于 NIO(Nonblocking I/O,非阻塞 IO)开发的网络通信框架。
Netty 的特性
- 高并发:Netty 是一款基于 NIO(Nonblocking IO,非阻塞 IO)开发的网络通信框架,对比于 BIO(Blocking I/O,阻塞 IO),他的并发性能得到了很大提高。
- 传输快:Netty 的传输依赖于内存零拷贝特性,尽量减少不必要的内存拷贝,实现了更高效率的传输。
- 封装好:Netty 封装了 NIO 操作的很多细节,提供了易于使用调用接口。
核心组件
Channel:Netty 网络操作抽象类,它除了包括基本的 I/O 操作,如 bind、connect、read、write 等。EventLoop:主要是配合 Channel 处理 I/O 操作,用来处理连接的生命周期中所发生的事情。ChannelFuture:Netty 框架中所有的 I/O 操作都为异步的,因此我们需要 ChannelFuture 的 addListener() 注册一个 ChannelFutureListener 监听事件,当操作执行成功或者失败时,监听就会自动触发返回结果。ChannelHandler:充当了所有处理入站和出站数据的逻辑容器。ChannelHandler 主要用来处理各种事件,这里的事件很广泛,比如可以是连接、数据接收、异常、数据转换等。ChannelPipeline:为 ChannelHandler 链提供了容器,当 channel 创建时,就会被自动分配到它专属的 ChannelPipeline,这个关联是永久性的。
Netty 有两种发送消息的方式:
- 直接写入 Channel 中,消息从 ChannelPipeline 当中尾部开始移动;
- 写入和 ChannelHandler 绑定的 ChannelHandlerContext 中,消息从 ChannelPipeline 中的下一个 ChannelHandler 中移动。
高性能
Netty 高性能表现在哪些方面:
- NIO 线程模型:同步非阻塞,用最少的资源做更多的事。
- 内存零拷贝:尽量减少不必要的内存拷贝,实现了更高效率的传输。
- 内存池设计:申请的内存可以重用,主要指直接内存。内部实现是用一颗二叉查找树管理内存分配情况。
- 串形化处理读写:避免使用锁带来的性能开销。
- 高性能序列化协议:支持 protobuf 等高性能序列化协议。
零拷贝
传统意义的拷贝
是在发送数据的时候,传统的实现方式是:
File.read(bytes)
Socket.send(bytes)
这种方式需要四次数据拷贝和四次上下文切换:
数据从磁盘读取到内核的 read buffer
数据从内核缓冲区拷贝到用户缓冲区
数据从用户缓冲区拷贝到内核的 socket buffer
数据从内核的 socket buffer 拷贝到网卡接口(硬件)的缓冲区
零拷贝的概念
明显上面的第二步和第三步是非必要的,通过 java 的 FileChannel.transferTo 方法,可以避免上面两次多余的拷贝(当然这需要底层操作系统支持)
- 调用 transferTo,数据从文件由 DMA 引擎拷贝到内核 read buffer
- 接着 DMA 从内核 read buffer 将数据拷贝到网卡接口 buffer
上面的两次操作都不需要 CPU 参与,所以就达到了零拷贝。
Netty 中的零拷贝
主要体现在三个方面:
bytebuffer
Netty 发送和接收消息主要使用 bytebuffer,bytebuffer 使用对外内存(DirectMemory)直接进行 Socket 读写。
原因:如果使用传统的堆内存进行 Socket 读写,JVM 会将堆内存 buffer 拷贝一份到直接内存中然后再写入 socket,多了一次缓冲区的内存拷贝。DirectMemory 中可以直接通过 DMA 发送到网卡接口
Composite Buffers
传统的 ByteBuffer,如果需要将两个 ByteBuffer 中的数据组合到一起,我们需要首先创建一个 size=size1+size2 大小的新的数组,然后将两个数组中的数据拷贝到新的数组中。但是使用 Netty 提供的组合 ByteBuf,就可以避免这样的操作,因为 CompositeByteBuf 并没有真正将多个 Buffer 组合起来,而是保存了它们的引用,从而避免了数据的拷贝,实现了零拷贝。
对于 FileChannel.transferTo 的使用
Netty 中使用了 FileChannel 的 transferTo 方法,该方法依赖于操作系统实现零拷贝。
Netty 流程
graph TB
A[客户端发起连接] --> B[BossGroup 接收连接]
B --> C[NioEventLoop 处理 Accept]
C --> D[创建 NioSocketChannel]
D --> E[注册到 WorkerGroup 的 EventLoop]
E --> F[ChannelPipeline 初始化]
F --> G[ChannelHandler 链处理 I/O 事件]
G --> H{入站/出站}
H -->|入站| I[ChannelInboundHandler 处理读取]
H -->|出站| J[ChannelOutboundHandler 处理写入]
I --> K[业务逻辑处理]
J --> L[数据写回客户端]典型应用场景
- RPC 框架:Dubbo、gRPC 等 RPC 框架底层使用 Netty 实现高性能的网络通信,利用其 NIO 模型支撑高并发远程调用。
- 消息中间件:RocketMQ、Kafka 等消息中间件使用 Netty 作为网络传输层,实现低延迟、高吞吐的消息传递。
- 实时通信服务器:WebSocket 长连接服务器、即时通讯(IM)系统、在线游戏服务器等需要维护大量长连接的场景。
- HTTP/HTTPS 服务器:构建轻量级 HTTP 服务器或 API 网关,如 Spring WebFlux 底层即基于 Netty。
- 分布式协调服务:Elasticsearch、ZooKeeper 等分布式系统节点间通信使用 Netty 实现高效的数据同步。
最佳实践
- 合理配置 EventLoopGroup 线程数:BossGroup 通常配置 1 个线程处理连接接受,WorkerGroup 根据 CPU 核数配置(一般为 CPU 核数 × 2)。
- 使用池化的 Direct Buffer:开启
PooledByteBufAllocator减少内存分配开销,避免频繁 GC。 - 合理设置 Pipeline 中 Handler 顺序:编解码器在前、业务处理器在后,避免不必要的解码开销。
- 注意资源释放:确保
ByteBuf在使用后正确释放(调用ReferenceCountUtil.release()),防止内存泄漏。 - 使用心跳机制维持长连接:通过
IdleStateHandler检测空闲连接,配合心跳包保持连接活性。 - 优雅关闭:调用
shutdownGracefully()而非shutdown(),确保在关闭前处理完队列中的任务。
常见问题
Netty 的线程模型与 Tomcat 等传统 Web 容器有何区别?
传统 Tomcat 采用 BIO 模型,每个连接一个线程,线程数多、上下文切换开销大;Netty 基于 NIO Reactor 模型,用少量线程处理大量连接,通过事件驱动避免阻塞,在高并发场景下性能优势显著。
什么时候应该选择 Netty 而非 Spring MVC?
当应用需要支撑大量长连接(如 IM、IoT)、需要自定义协议、或对延迟和吞吐量有极致要求时,应选择 Netty。普通 Web 应用使用 Spring MVC / Spring Boot 更为便捷。
如何排查 Netty 内存泄漏?
启动时添加 JVM 参数 -Dio.netty.leakDetectionLevel=PARANOID 开启严格泄漏检测,结合日志中的 LEAK 关键字定位未释放的 ByteBuf。生产环境建议使用 ADVANCED 级别平衡性能和检测能力。
应用
Netty 是一个广泛使用的 Java 网络编程框架。很多著名软件都使用了它,如:Dubbo、Cassandra、Elasticsearch、Vert.x 等。
有了 Netty,你可以实现自己的 HTTP 服务器,FTP 服务器,UDP 服务器,RPC 服务器,WebSocket 服务器,Redis 的 Proxy 服务器,MySQL 的 Proxy 服务器等等。
public class NettyOioServer {
public void server(int port) throws Exception {
final ByteBuf buf = Unpooled.unreleasableBuffer(
Unpooled.copiedBuffer("Hi!\r\n", Charset.forName("UTF-8")));
EventLoopGroup group = new OioEventLoopGroup();
try {
ServerBootstrap b = new ServerBootstrap(); //1
b.group(group) //2
.channel(OioServerSocketChannel.class)
.localAddress(new InetSocketAddress(port))
.childHandler(new ChannelInitializer<SocketChannel>() {//3
@Override
public void initChannel(SocketChannel ch)
throws Exception {
ch.pipeline().addLast(new ChannelInboundHandlerAdapter() { //4
@Override
public void channelActive(ChannelHandlerContext ctx) throws Exception {
ctx.writeAndFlush(buf.duplicate()).addListener(ChannelFutureListener.CLOSE);//5
}
});
}
});
ChannelFuture f = b.bind().sync(); //6
f.channel().closeFuture().sync();
} finally {
group.shutdownGracefully().sync(); //7
}
}
}