Netty 是一个异步事件驱动的网络应用框架,用于快速开发可维护的高性能协议服务器和客户端。它是 Java 网络编程领域的事实标准,支撑了 Dubbo、RocketMQ、gRPC 等大量知名项目。
一、Netty 核心架构
1.1 Reactor 模型演进
单线程 Reactor(所有 I/O 在一个线程):
┌──────────────────────────────┐
│ Reactor Thread │
│ ┌─────────┐ ┌───────────┐ │
│ │ Selector │──→│ Dispatch │ │
│ └────┬────┘ └─────┬─────┘ │
│ │ │ │
│ ▼ ▼ │
│ Accept Handler Read/Write│
│ + Execute Handler │
└──────────────────────────────┘
多线程 Reactor(Boss 接收、Worker 处理):
┌────────────┐ ┌────────────────────────────┐
│ Boss Group │ │ Worker Group │
│ (1 thread) │ │ (N threads, default 2*CPU)│
│ │ │ │
│ Selector │─────→│ Thread 1: Selector + I/O │
│ Accept │ │ Thread 2: Selector + I/O │
└────────────┘ │ Thread N: Selector + I/O │
└────────────────────────────┘
主从 Reactor(Netty 默认):
┌────────────────────────────────────────────┐
│ Boss Group (监听端口) │
│ ├── Selector (accept 事件) │
│ └── 将 SocketChannel 注册到 Worker Group │
└────────────────────────────────────────────┘
│
▼
┌────────────────────────────────────────────┐
│ Worker Group (处理 I/O) │
│ ├── Thread 1: Selector (read/write) │
│ ├── Thread 2: Selector (read/write) │
│ └── Thread N: Selector (read/write) │
└────────────────────────────────────────────┘
1.2 核心组件关系
Channel ──→ 网络通道(Socket 抽象)
│
├── EventLoop ──→ 处理 I/O 事件的事件循环
│ │
│ ├── Selector ──→ 多路复用器
│ └── TaskQueue ──→ 用户任务队列
│
├── ChannelPipeline ──→ 处理器链
│ │
│ ├── ChannelInboundHandler ──→ 入站处理
│ └── ChannelOutboundHandler ──→ 出站处理
│
├── ChannelHandlerContext ──→ 处理器上下文
│
└── Unsafe ──→ 底层 I/O 操作
二、快速入门:Echo Server
2.1 服务端
public class NettyEchoServer {
public static void main(String[] args) throws Exception {
// Boss:处理 accept 事件
EventLoopGroup bossGroup = new NioEventLoopGroup(1);
// Worker:处理 read/write 事件
EventLoopGroup workerGroup = new NioEventLoopGroup();
try {
ServerBootstrap bootstrap = new ServerBootstrap();
bootstrap.group(bossGroup, workerGroup)
.channel(NioServerSocketChannel.class)
.option(ChannelOption.SO_BACKLOG, 1024)
.childOption(ChannelOption.TCP_NODELAY, true)
.childOption(ChannelOption.SO_KEEPALIVE, true)
.childHandler(new ChannelInitializer<SocketChannel>() {
@Override
protected void initChannel(SocketChannel ch) {
ch.pipeline()
.addLast(new LoggingHandler(LogLevel.INFO))
.addLast(new StringDecoder(CharsetUtil.UTF_8))
.addLast(new StringEncoder(CharsetUtil.UTF_8))
.addLast(new EchoServerHandler());
}
});
ChannelFuture future = bootstrap.bind(8080).sync();
System.out.println("Echo Server started on port 8080");
future.channel().closeFuture().sync();
} finally {
bossGroup.shutdownGracefully();
workerGroup.shutdownGracefully();
}
}
}
@ChannelHandler.Sharable
public class EchoServerHandler extends ChannelInboundHandlerAdapter {
@Override
public void channelRead(ChannelHandlerContext ctx, Object msg) {
System.out.println("Received: " + msg);
ctx.write(msg); // 写回客户端(出站事件)
}
@Override
public void channelReadComplete(ChannelHandlerContext ctx) {
ctx.flush(); // 将缓冲区的数据flush到Socket
}
@Override
public void exceptionCaught(ChannelHandlerContext ctx, Throwable cause) {
cause.printStackTrace();
ctx.close();
}
}
2.2 客户端
public class NettyEchoClient {
public static void main(String[] args) throws Exception {
EventLoopGroup group = new NioEventLoopGroup();
try {
Bootstrap bootstrap = new Bootstrap();
bootstrap.group(group)
.channel(NioSocketChannel.class)
.option(ChannelOption.TCP_NODELAY, true)
.handler(new ChannelInitializer<SocketChannel>() {
@Override
protected void initChannel(SocketChannel ch) {
ch.pipeline()
.addLast(new StringDecoder(CharsetUtil.UTF_8))
.addLast(new StringEncoder(CharsetUtil.UTF_8))
.addLast(new EchoClientHandler());
}
});
ChannelFuture future = bootstrap.connect("localhost", 8080).sync();
Channel channel = future.channel();
// 发送消息
for (int i = 0; i < 10; i++) {
channel.writeAndFlush("Hello Netty " + i + "\n");
}
channel.closeFuture().sync();
} finally {
group.shutdownGracefully();
}
}
}
public class EchoClientHandler extends ChannelInboundHandlerAdapter {
@Override
public void channelRead(ChannelHandlerContext ctx, Object msg) {
System.out.println("Server response: " + msg);
}
}
三、Channel Pipeline 详解
3.1 入站与出站事件流
ChannelPipeline 事件流向:
入站事件(Inbound):
Head ──→ Handler1 ──→ Handler2 ──→ Handler3 ──→ Tail
(read) (decode) (auth) (biz) (end)
↑
ctx.fireChannelRead(msg)
出站事件(Outbound):
Tail ←── Handler3 ←── Handler2 ←── Handler1 ←── Head
(end) (biz) (encode) (compress) (write)
↑
ctx.write(msg)
注意:
- Inbound 事件从前向后传递
- Outbound 事件从后向前传递
- 非常关键!Handler 的顺序决定处理流程
3.2 内置编解码器
// 处理 Redis 协议
ch.pipeline()
.addLast(new RedisDecoder())
.addLast(new RedisBulkStringAggregator())
.addLast(new RedisArrayAggregator())
.addLast(new RedisEncoder());
// 处理 HTTP
ch.pipeline()
.addLast(new HttpServerCodec()) // HTTP 编解码
.addLast(new HttpObjectAggregator(65536)) // 聚合请求
.addLast(new ChunkedWriteHandler()) // 支持 chunked
.addLast(new HttpServerHandler());
// 处理 TCP 粘包/拆包(LengthFieldBasedFrameDecoder)
ch.pipeline()
.addLast(new LengthFieldBasedFrameDecoder(
65536, // 最大帧长度
0, // lengthFieldOffset
4, // lengthFieldLength
0, // lengthAdjustment
4 // initialBytesToStrip
))
.addLast(new ProtobufDecoder(Message.getDefaultInstance()))
.addLast(new ProtobufEncoder())
.addLast(new BusinessHandler());
四、ByteBuf 内存管理
4.1 ByteBuf vs Java ByteBuffer
| 特性 | Java ByteBuffer | Netty ByteBuf |
|---|---|---|
| 读写模式 | 需 flip() 切换 | 读写索引分离 |
| 容量扩展 | 固定或重新分配 | 动态扩展 |
| 引用计数 | 无 | 有(支持池化) |
| 池化 | 无 | 有(PooledByteBufAllocator) |
| 复合缓冲区 | 无 | 有(CompositeByteBuf) |
4.2 核心 API
// 创建 ByteBuf
ByteBuf buf = Unpooled.buffer(1024); // 非池化堆内存
ByteBuf buf = Unpooled.directBuffer(1024); // 非池化直接内存
ByteBuf buf = PooledByteBufAllocator.DEFAULT.buffer(1024); // 池化
// 写入数据
buf.writeBytes("Hello".getBytes());
buf.writeInt(123);
buf.writeLong(System.currentTimeMillis());
// 读取数据
byte[] read = new byte[buf.readableBytes()];
buf.readBytes(read);
int value = buf.readInt();
// 零拷贝切片
ByteBuf sliced = buf.slice(0, 5); // 共享底层内存,无复制
ByteBuf copied = buf.copy(0, 5); // 深拷贝
// 组合多个 ByteBuf(零拷贝)
CompositeByteBuf composite = Unpooled.compositeBuffer();
composite.addComponents(true, headerBuf, bodyBuf, footerBuf);
4.3 引用计数(内存不泄露的关键)
public class SafeHandler extends ChannelInboundHandlerAdapter {
@Override
public void channelRead(ChannelHandlerContext ctx, Object msg) {
ByteBuf buf = (ByteBuf) msg;
try {
// 使用 buf
System.out.println(buf.toString(CharsetUtil.UTF_8));
// 传递给下一个 handler(引用计数 +1)
ctx.fireChannelRead(msg);
// 注意:如果 fireChannelRead,不要 release!
// 由 TailContext 或后续 handler 释放
} catch (Exception e) {
// 异常时需要手动释放
buf.release();
throw e;
}
}
}
// 更安全的做法:使用 ReferenceCountUtil
public class SaferHandler extends ChannelInboundHandlerAdapter {
@Override
public void channelRead(ChannelHandlerContext ctx, Object msg) {
try {
if (msg instanceof ByteBuf) {
ByteBuf buf = (ByteBuf) msg;
// 使用 buf
}
} finally {
// 自动处理引用计数
ReferenceCountUtil.release(msg);
}
}
}
五、零拷贝技术
5.1 传统数据传输
传统文件发送(4 次拷贝):
磁盘 ──→ Kernel Buffer ──→ User Buffer ──→ Socket Buffer ──→ 网卡
1 2 3 (应用中处理) 4 5
↑
CPU 参与拷贝
5.2 Netty 零拷贝
// 1. FileRegion(基于 sendfile)
public void sendFile(ChannelHandlerContext ctx, File file) throws IOException {
FileRegion region = new DefaultFileRegion(
new RandomAccessFile(file, "r").getChannel(),
0, // 偏移量
file.length() // 长度
);
ctx.writeAndFlush(region);
}
// 2. CompositeByteBuf(组合多个 ByteBuf,无需拷贝)
public ByteBuf createPacket(ByteBuf header, ByteBuf body, ByteBuf footer) {
CompositeByteBuf packet = ctx.alloc().compositeBuffer(3);
packet.addComponents(true, header, body, footer);
return packet; // 底层共享三个 ByteBuf 的内存
}
// 3. slice / duplicate(共享内存)
public void processHeader(ByteBuf fullPacket) {
ByteBuf header = fullPacket.slice(0, 16); // 只读视图
ByteBuf body = fullPacket.slice(16, fullPacket.readableBytes() - 16);
// header 和 body 共享 fullPacket 的内存
}
// 4. Unpooled.wrappedBuffer
byte[] array = new byte[1024];
ByteBuf buf = Unpooled.wrappedBuffer(array); // 包装现有数组,不复制
六、EventLoop 与线程模型
// EventLoop 执行逻辑
public void execute(Runnable task) {
if (inEventLoop()) {
// 当前线程就是 EventLoop 线程,直接执行
task.run();
} else {
// 将任务加入队列,等待 EventLoop 线程执行
taskQueue.offer(task);
}
}
// 关键原则:不要在 ChannelHandler 中做阻塞操作!
public class BadHandler extends ChannelInboundHandlerAdapter {
@Override
public void channelRead(ChannelHandlerContext ctx, Object msg) {
// ❌ 错误:阻塞 EventLoop 线程
Thread.sleep(1000);
// ❌ 错误:同步调用数据库
database.query("SELECT ...");
}
}
public class GoodHandler extends ChannelInboundHandlerAdapter {
private final EventExecutor executor; // 业务线程池
@Override
public void channelRead(ChannelHandlerContext ctx, Object msg) {
// ✅ 正确:异步提交到业务线程池
executor.execute(() -> {
// 阻塞操作在业务线程中执行
Result result = database.query("SELECT ...");
// 结果回写到 EventLoop
ctx.channel().eventLoop().execute(() -> {
ctx.writeAndFlush(result);
});
});
}
}
七、性能调优参数
| 参数 | 说明 | 推荐值 |
|---|---|---|
| SO_BACKLOG | 连接队列大小 | 1024-65535 |
| TCP_NODELAY | 禁用 Nagle | true(低延迟) |
| SO_KEEPALIVE | TCP 心跳 | true |
| SO_RCVBUF/SO_SNDBUF | 缓冲区大小 | 根据带宽时延积调整 |
| ALLOCATOR | ByteBuf 分配器 | PooledByteBufAllocator |
| RCVBUF_ALLOCATOR | 自适应接收缓冲区 | AdaptiveRecvByteBufAllocator |
| WRITE_BUFFER_WATER_MARK | 写缓冲区水位 | 32KB / 64KB |
| CHANNEL_OPTION.CONNECT_TIMEOUT_MILLIS | 连接超时 | 3000-10000 |
ServerBootstrap bootstrap = new ServerBootstrap();
bootstrap
.option(ChannelOption.SO_BACKLOG, 65535)
.option(ChannelOption.SO_REUSEADDR, true)
.childOption(ChannelOption.TCP_NODELAY, true)
.childOption(ChannelOption.SO_KEEPALIVE, true)
.childOption(ChannelOption.ALLOCATOR, PooledByteBufAllocator.DEFAULT)
.childOption(ChannelOption.RCVBUF_ALLOCATOR, new AdaptiveRecvByteBufAllocator())
.childOption(ChannelOption.WRITE_BUFFER_WATER_MARK, new WriteBufferWaterMark(32 * 1024, 64 * 1024));
八、总结
| 核心概念 | 说明 |
|---|---|
| EventLoop | 单线程执行所有 I/O 和任务,避免锁竞争 |
| ChannelPipeline | 责任链模式,灵活组装编解码与业务逻辑 |
| ByteBuf | 引用计数的缓冲区,支持池化和零拷贝 |
| 零拷贝 | FileRegion、CompositeByteBuf、slice |
| 线程安全 | ChannelHandler 可被多个 Channel 共享(加 @Sharable) |
Netty 的设计精髓在于将复杂的 NIO 操作封装为简洁的事件驱动模型,通过精心设计的内存管理和线程模型,实现了单机百万连接的高性能网络处理能力。
继续阅读
探索更多技术文章
浏览归档,发现更多关于系统设计、工具链和工程实践的内容。