Netty 高性能网络编程框架

掌握 Netty 的 Reactor 模型、Channel Pipeline、ByteBuf 内存管理与零拷贝技术,构建百万级并发连接的网络服务

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 ByteBufferNetty 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禁用 Nagletrue(低延迟)
SO_KEEPALIVETCP 心跳true
SO_RCVBUF/SO_SNDBUF缓冲区大小根据带宽时延积调整
ALLOCATORByteBuf 分配器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 操作封装为简洁的事件驱动模型,通过精心设计的内存管理和线程模型,实现了单机百万连接的高性能网络处理能力。

继续阅读

探索更多技术文章

浏览归档,发现更多关于系统设计、工具链和工程实践的内容。

全部文章 返回首页

「network」更多文章

  1. 网络安全:TLS/SSL、证书与加密通信
  2. 负载均衡算法与高可用架构
  3. DNS 系统与智能解析