Administrator
发布于 2022-01-13 / 0 阅读
0
0

Netty-202509120209

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,用于启动服务端

io.netty

netty-all

4.1.92.Final

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();

}

}

}

/**

  • @Description: 自定义Handler需要继承netty规定好的某个HandlerAdapter (规范)
  • @BusinessDesc 业务逻辑:连接成功后,打印 连接成功; 收到客户端的消息时,打印,
  • 并回复消息

  • @Author: sxl
  • @Date:
  • */

    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 {

    /**

  • 自定义Handler需要继承netty规定好的某个HandlerAdapter (规范)
  • 业务逻辑:连接成功后,给服务端发送一条消息
  • *

    */

    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

    ) {


    评论