Java网络编程[BIO][NIO][NIO主从Reactor模型]
● 网络编程相关
网络编程是大部分框架的基础tomcat、springboot、jdbc、消息队列rabbitMQ、redis、mysql等
网络编程核心组件:
- InetAddress:他的作用是代表一个地址,比如127.0.0.1就是代表本地地址,他就是这个127.0.0.1而
- TCP\UDP:代表传输数据信息的方式TCP:四挥三握、一对一安全可靠、UDP:广播速度快不可靠等
- Socket\ServerSocket:封装InetAddress与连接方式TCP\UDP,指定发送与接受的规则(TCP\UDP)、指定发送方(Socket)与接收方(ServerSocket)的地址
比如在springboot+vue环境:
- 页面发送一个js请求,浏览器底层会使用c++代码生成一个socket,这个socket封装了源ip与目标ip将数据信息(请求URL)通过IO发送给目标ip的指定端口(浏览器会将前端要发送请求携带的请求体的js对象转换为http协议文本application/json)然后碾碎成二进制字节码发送
- springboot初始化后会开启内部的tomcat服务器,tomcat会开一个ServerSocket(NIO模式)与springboot配置的端口号进行绑定并监听该端口(8080),前端发请求就会从这个端口接听到,之后调用servlet的DispatcherServlet来决定使用哪一个controller接受请求,如果遇到@RequestBody注解,就会自动调用Jackson(json解析)利用反射将数据实例化为user对象接口接受参数对象
- 等业务层处理完业务之后后端再将数据通过Socket(此时不是ServerSocket它只是一个服务器接待员,真正发送响应数据的是内部的Socket)把数据通过Jackson转化给json对象再封装成http协议文本通过二进制IO流传给浏览器的Socket接受
● InetAddresss
InetAddress相当与一个通讯录本,可以来记录IP地址
| 方法类型 | 方法签名 (Method Signature) | 功能描述 (Description) | 返回值类型 (Return Type) |
|---|---|---|---|
| 静态工厂 | static getByName(String host) |
在给定主机名的情况下确定主机的 IP 地址。 | InetAddress |
| 静态工厂 | static getAllByName(String host) |
在给定主机名的情况下,返回其绑定的所有 IP 地址。 | InetAddress[] |
| 静态工厂 | static getLocalHost() |
获取本地主机的地址对象。 | InetAddress |
| 实例方法 | getHostAddress() |
返回 IP 地址字符串(以文本表现形式)。 | String |
| 实例方法 | getHostName() |
获取此 IP 地址的主机名。 | String |
| 实例方法 | isReachable(int timeout) |
测试该地址在指定时间内是否可达。 | boolean |
package com.note;
import java.io.IOException;
import java.net.InetAddress;
import java.net.UnknownHostException;
public class SocketApplication {
public static void main(String[] args) throws IOException {
//知域名获取ip
System.out.println(InetAddress.getByName("www.baidu.com")); //www.baidu.com/39.156.70.239
//获取域名集群ip
InetAddress[] allByName = InetAddress.getAllByName("www.baidu.com");
for (InetAddress inetAddress : allByName) {
System.out.println(inetAddress);
//www.baidu.com/39.156.70.239
//www.baidu.com/39.156.70.46
//www.baidu.com/2409:8c00:6c21:118b:0:ff:b0e8:f003
//www.baidu.com/2409:8c00:6c21:11eb:0:ff:b0bf:59ca
}
System.out.println(InetAddress.getByName("www.baidu.com").getHostAddress()); //39.156.70.239
//获取本机ip
System.out.println(InetAddress.getLocalHost()); //username/192.168.1.47 可能是内网ip或虚拟网卡ip
System.out.println(InetAddress.getLoopbackAddress()); //获取回环地址
System.out.println(InetAddress.getByName("www.baidu.com").isReachable(3000)); //指定时间内是否可达 true
}
}
● Socket
○ TCP
socket相当于一个数据流管道的端点接口,当new出一个socket时就已经完成了三次握手的虚拟管道
作用就是将数据接受或是发出数据,他自己有一个socket.getInputStream或socket.getOutputStream的流,当一个发送数据时,监听端口会监听到流的声音然后执行发送或接受操作(他的流失字节流,建议试字符转换流)
- 输入流 (InputStream):从管道里往外拿数据(读)。
- 输出流 (OutputStream):往管道里塞数据(写)。
流程:
- 建立连接
- 获取流(输入输出)
- 接受\发送数据
- 关闭资源
Socket API
| 方法分类 | 方法签名 (Method Signature) | 功能描述 (Description) | 返回值 |
|---|---|---|---|
| 构造方法 | Socket() |
创建一个未连接的 Socket 对象。需配合 connect() 使用。 |
Socket |
| 构造方法 | Socket(String host, int port) |
创建一个流套接字并将其连接到指定主机上的指定端口号。 | Socket |
| 连接控制 | connect(SocketAddress endpoint) |
将此套接字连接到服务器。 | void |
| 连接控制 | connect(SocketAddress endpoint, int timeout) |
连接到服务器,并指定连接超时时间(毫秒)。 | void |
| I/O 获取 | getInputStream() |
返回此套接字的输入流(用于接收数据)。 | InputStream |
| I/O 获取 | getOutputStream() |
返回此套接字的输出流(用于发送数据)。 | OutputStream |
| 关闭操作 | close() |
关闭此套接字,释放所有相关资源。 | void |
| 关闭操作 | shutdownOutput() |
禁用此套接字的输出流(半关闭)。后续写入会抛出异常,通常用于告知对方数据发送完毕。 | void |
| 关闭操作 | shutdownInput() |
禁用此套接字的输入流。后续读取会返回 EOF(-1)。 | void |
| 状态查询 | isClosed() |
检查 Socket 是否已关闭。 | boolean |
| 状态查询 | isConnected() |
检查 Socket 是否曾经连接成功过。 | boolean |
| 设置参数 | setSoTimeout(int timeout) |
启用/禁用 SO_TIMEOUT(读取超时),以毫秒为单位。 |
void |
| 设置参数 | setTcpNoDelay(boolean on) |
启用/禁用 TCP_NODELAY(Nagle 算法)。true 表示禁用缓冲,立即发送。 |
void |
| 设置参数 | setKeepAlive(boolean on) |
启用/禁用 SO_KEEPALIVE。 |
void |
| 获取信息 | getInetAddress() |
返回套接字连接的远程 IP 地址。 | InetAddress |
| 获取信息 | getPort() |
返回套接字连接的远程端口号。 | int |
| 获取信息 | getLocalAddress() |
获取套接字绑定的本地 IP 地址。 | InetAddress |
| 获取信息 | getLocalPort() |
获取套接字绑定的本地端口号。 | int |
ServerSocket API
| 方法分类 | 方法签名 (Method Signature) | 功能描述 (Description) | 返回值 |
|---|---|---|---|
| 构造方法 | ServerSocket(int port) |
创建绑定到特定端口的服务器套接字。队列长度默认 50。 | ServerSocket |
| 构造方法 | ServerSocket(int port, int backlog) |
创建绑定到特定端口的服务器套接字,并指定连接队列长度(backlog)。 | ServerSocket |
| 构造方法 | ServerSocket() |
创建未绑定的服务器套接字。需配合 bind() 使用。 |
ServerSocket |
| 核心功能 | accept() |
监听并接受连接。此方法是阻塞的,直到建立连接。一旦连接建立,返回一个新的 Socket 对象。 |
Socket |
| 绑定操作 | bind(SocketAddress endpoint) |
将 ServerSocket 绑定到特定的地址(IP 和端口)。 | void |
| 绑定操作 | bind(SocketAddress endpoint, int backlog) |
绑定地址并指定请求队列的最大长度。 | void |
| 关闭操作 | close() |
关闭此套接字。停止监听,不再接受新的连接。 | void |
| 状态查询 | isClosed() |
检查 ServerSocket 是否已关闭。 | boolean |
| 状态查询 | isBound() |
检查 ServerSocket 是否已绑定到特定端口。 | boolean |
| 获取信息 | getInetAddress() |
返回此服务器套接字的本地地址。 | InetAddress |
| 获取信息 | getLocalPort() |
返回此服务器套接字侦听的端口。 | int |
| 配置参数 | setSoTimeout(int timeout) |
设置 accept() 的超时时间(毫秒)。若超时未连接,抛出 SocketTimeoutException。0 表示无限等待。 |
void |
| 配置参数 | setReuseAddress(boolean on) |
启用/禁用 SO_REUSEADDR。允许在关闭连接后立即重用该端口(防止重启服务时报错“端口被占用”)。 |
void |
| 配置参数 | setReceiveBufferSize(int size) |
设置接受缓冲区的大小(SO_RCVBUF)。 | void |
△ 模拟发送方
package com.note;
import java.io.IOException;
import java.io.OutputStreamWriter;
import java.io.PrintWriter;
import java.net.Socket;
import java.util.Scanner;
public class SocketApplication {
public static void main(String[] args) throws IOException {
System.out.println("【发送方】准备连接...");
// 1. 建立连接:指定目标的 IP 和 端口
// "127.0.0.1" 代表本机,如果对方在别的电脑,就填那台电脑的 IP
try (Socket socket = new Socket("127.0.0.1", 8888)) {
System.out.println("【发送方】连接成功!请在控制台输入消息(输入 bye 退出):");
// 2. 获取输出流(OutputStream):准备发送数据
// PrintWriter 是个好东西,第二个参数 true 表示自动刷新缓冲区 (Auto Flush)
// 如果不开启自动刷新,你的消息可能会卡在内存里发不出去!
PrintWriter writer = new PrintWriter(new OutputStreamWriter(socket.getOutputStream()), true);
// 3. 从键盘读取输入(模拟用户打字)
Scanner scanner = new Scanner(System.in);
while (true) {
String line = scanner.nextLine(); // 阻塞等待键盘输入
// 4. 发送数据
writer.println(line); // println 自带换行符,这点很重要!因为接收方用的是 readLine()
if ("bye".equals(line)) {
break;
}
}
} catch (IOException e) {
System.err.println("【发送方】连接失败,请检查接收方是否已启动!");
}
}
}
△ 模拟接收方
package com.note;
import java.io.BufferedReader;
import java.io.IOException;
import java.io.InputStreamReader;
import java.net.ServerSocket;
import java.net.Socket;
public class SocketApplication {
public static void main(String[] args) {
System.out.println("接收方启动");
try (ServerSocket serverSocket = new ServerSocket( 8888)) {
System.out.println("【接收方】监听端口 8888,等待连接...");
// 2. 阻塞等待(accept):程序卡在这里,直到有客户端连上来
Socket socket = serverSocket.accept();
System.out.println("【接收方】客户端连接成功!IP: " + socket.getInetAddress());
// 3. 获取输入流(InputStream):准备接收数据
// 包装成 BufferedReader 是为了可以使用 readLine() 一行一行地读,比较方便
BufferedReader reader = new BufferedReader(new InputStreamReader(socket.getInputStream()));
String msg;
// 4. 循环读取:只要对方不挂断,我就一直读
// readLine() 也是阻塞的,对方不发换行符,我就一直等
while ((msg = reader.readLine()) != null) {
System.out.println("【接收方】收到消息: " + msg);
if ("bye".equals(msg)) {
System.out.println("【接收方】对方下线了,我也撤了。");
break;
}
}
} catch (IOException e) {
System.out.println(e.getMessage());
}
}
}
△ UDP
UDP是广播发送,不需要建立连接管道,发送方只负责把数据发出去不需要确定数据是否到达
UDP协议有两个核心类:
DatagramSocket:代表接与发的地址可以有IP端口等
DatagramPacket:代表数据包,UDP不是字节发送,而是打成一个包群发
DatagramSocket API
| 方法分类 | 方法签名 (Method Signature) | 功能描述 (Description) | 返回值 |
|---|---|---|---|
| 构造方法 | DatagramSocket() |
创建一个数据报套接字,绑定到本地任意可用端口。通常用于客户端。 | DatagramSocket |
| 构造方法 | DatagramSocket(int port) |
创建一个数据报套接字,绑定到本地指定端口。通常用于服务端。 | DatagramSocket |
| 构造方法 | DatagramSocket(int port, InetAddress laddr) |
创建一个数据报套接字,绑定到指定的本地地址和端口。 | DatagramSocket |
| 核心功能 | send(DatagramPacket p) |
发送数据报。将打包好的数据包发送出去。 | void |
| 核心功能 | receive(DatagramPacket p) |
接收数据报。阻塞等待,直到接收到一个数据包并填充到参数 p 中。 |
void |
| 连接操作 | connect(InetAddress address, int port) |
连接到远程地址。注意:UDP 是无连接的,这里的 connect 只是为了限制后续 send/receive 的目标,提高安全性或性能,并不会像 TCP 那样握手。 | void |
| 连接操作 | disconnect() |
断开连接。 | void |
| 关闭操作 | close() |
关闭此数据报套接字。 | void |
| 设置参数 | setSoTimeout(int timeout) |
设置 receive() 的超时时间(毫秒)。超时抛出 SocketTimeoutException。 |
void |
| 设置参数 | setSendBufferSize(int size) |
设置发送缓冲区大小(SO_SNDBUF)。 | void |
| 设置参数 | setReceiveBufferSize(int size) |
设置接收缓冲区大小(SO_RCVBUF)。 | void |
DatagramPackat API
| 方法分类 | 方法签名 (Method Signature) | 功能描述 (Description) | 返回值 |
|---|---|---|---|
| 构造方法 (接收用) | DatagramPacket(byte[] buf, int length) |
创建一个用于接收数据的包。buf 是存放数据的缓冲区,length 是接收的最大长度。 |
DatagramPacket |
| 构造方法 (发送用) | DatagramPacket(byte[] buf, int length, InetAddress address, int port) |
创建一个用于发送数据的包。指定了目标地址和端口。 | DatagramPacket |
| 构造方法 (发送用) | DatagramPacket(byte[] buf, int offset, int length, SocketAddress address) |
创建一个用于发送数据的包,支持指定数据的偏移量(offset)。 | DatagramPacket |
| 获取数据 | getData() |
获取数据缓冲区(即构造时的 byte[])。 |
byte[] |
| 获取数据 | getLength() |
获取实际收到的数据长度(接收后)或将要发送的数据长度(发送前)。 | int |
| 获取数据 | getOffset() |
获取数据的偏移量。 | int |
| 获取地址 | getAddress() |
获取发送端或接收端的 IP 地址。 | InetAddress |
| 获取地址 | getPort() |
获取发送端或接收端的端口号。 | int |
| 获取地址 | getSocketAddress() |
获取完整的 Socket 地址(IP + 端口)。 | SocketAddress |
| 设置数据 | setData(byte[] buf) |
设置数据缓冲区。 | void |
| 设置数据 | setLength(int length) |
设置数据包的长度。 | void |
| 设置地址 | setAddress(InetAddress iaddr) |
设置目标 IP 地址。 | void |
| 设置地址 | setPort(int iport) |
设置目标端口号。 | void |
△ 模拟发送方
public static void main(String[] args) throws IOException {
System.out.println("UDP 发送端启动...");
// 1. 创建发送端 Socket (不需要指定端口,系统随机分配)
DatagramSocket socket = new DatagramSocket();
// 2. 准备数据
String msg = "Hello UDP!";
byte[] data = msg.getBytes();
// 3. 打包 (数据 + 长度 + 目标IP + 目标端口)
DatagramPacket packet = new DatagramPacket(
data,
data.length,
InetAddress.getByName("127.0.0.1"),
9999
);
// 4. 发射!(Fire and Forget)
socket.send(packet);
System.out.println("消息已发送");
socket.close();
}
△ 模拟接收方
public static void main(String[] args) throws IOException {
System.out.println("UDP 接收端启动,监听 9999...");
// 1. 创建接收端 Socket,占用端口 9999
DatagramSocket socket = new DatagramSocket(9999);
// 2. 准备一个“空袋子” (数据包容器)
// UDP 一个包最大 64KB (64 * 1024)
byte[] buffer = new byte[1024];
DatagramPacket packet = new DatagramPacket(buffer, buffer.length);
// 3. 阻塞接收 (类似于 accept,但这里是接包裹)
socket.receive(packet); // 此时数据已经被填入 buffer 里了
// 4. 拆包
// packet.getLength() 是实际收到的数据长度
String msg = new String(packet.getData(), 0, packet.getLength());
System.out.println("收到来自 " + packet.getAddress() + " 的消息: " + msg);
socket.close();
}
● BIO
BIO(Blocking IO)阻塞式IO
在一对多环境下,服务器可以接收很多浏览器发来的请求,每当一个请求进来时,就会给这个请求分配一个线程,这个线程将始终陪同这个请求,直到读写操作完毕后,这个线程才会释放,如果同时进来1w个请求就会同时创建1w个线程,这样的开销无比巨大,虽然可以通过伪异步操作(线程池)来缓解这个问题,但是治标不指本
○ BIO的C\S架构
服务器是一对多模式,也就是一个服务器可以同时接收多个请求也可以返回数据,客户端可以发请求也可以收到响应
C\S架构模式的服务器客户端收发
开始监听端口后,当一个请求进来之后,就会把这个请求扔给线程池来处理,主线程继续监听,子线程开始处理业务,处理完业务后返回数据给客户端,因为传进来的是Socket,这个Socket自带IP和端口所以就能找到要返回的那个客户端地址与端口
服务器:
package com.note;
import java.io.*;
import java.net.*;
import java.util.concurrent.ArrayBlockingQueue;
import java.util.concurrent.ThreadPoolExecutor;
import java.util.concurrent.TimeUnit;
@SuppressWarnings("all")
public class SocketApplication {
public static void main(String[] args) throws IOException {
//创建线程池 开始执行BIO模式的C/S模式
ThreadPoolExecutor threadPoolExecutor = new ThreadPoolExecutor(
5,
10,
20,
TimeUnit.SECONDS,
new ArrayBlockingQueue<>(10)
);
System.out.println("TCP 接收端启动,监听 9999...");
try (ServerSocket serverSocket = new ServerSocket(9999)) {
while (true) {
//每次接收到数据就扔给一个线程来处理
Socket socket = serverSocket.accept();
//扔给线程池阻塞队列来处理
threadPoolExecutor.execute(() -> {
try {
controller(socket);
} catch (IOException e) {
throw new RuntimeException(e);
}
});
}
} catch (IOException e) {
e.printStackTrace();
}
}
public static void controller(Socket socket) throws IOException {
System.out.println(Thread.currentThread().getName() + "已启动");
try (
BufferedReader br = new BufferedReader(new InputStreamReader(socket.getInputStream())); //负责处理请求
PrintWriter writer = new PrintWriter(new OutputStreamWriter(socket.getOutputStream()), true) //负责返回响应
) {
//读取
String msg = null;
while ((msg = br.readLine()) != null) {
System.out.println(msg);
//返回
writer.println(msg);
}
} catch (IOException e) {
e.printStackTrace();
}
System.out.println(Thread.currentThread().getName() + "已关闭");
}
}
客户端:
package com.note;
import java.io.*;
import java.net.InetAddress;
import java.net.Socket;
import java.util.Scanner;
public class SocketApplication {
public static void main(String[] args) throws IOException {
try (Socket socket = new Socket(InetAddress.getByName("127.0.0.1"), 9999)) {
// socket.setSoTimeout(5000);
//创建一个线程 专门用来接受响应数据
new Thread(() -> {
try {
listener(socket);
} catch (IOException e) {
throw new RuntimeException(e);
}
}).start();
PrintWriter printWriter = new PrintWriter(new OutputStreamWriter(socket.getOutputStream()), true);
Scanner scanner = new Scanner(System.in);
while (true) {
String line = scanner.nextLine();
printWriter.println(line);
}
} catch (IOException e) {
e.printStackTrace();
}
}
public static void listener(Socket socket) throws IOException {
try (BufferedReader br = new BufferedReader(new InputStreamReader(socket.getInputStream()))) {
String msg = null;
while ((msg = br.readLine()) != null) {
System.out.println(msg);
}
} catch (IOException e) {
e.printStackTrace();
}
}
}
○ BIO的B\S架构
B\S就是浏览器与服务器的互通
浏览器发送请求都是http协议文本接受也是所以数据格式要有规范
服务器会接收到浏览器发来的http协议文本如同:
POST /login HTTP/1.1
Content-Type: application/json
User-Agent: PostmanRuntime/7.51.1
Accept: */*
Postman-Token: 7609b09b-a21f-4183-8e14-4912b53ed268
Host: 127.0.0.1:9999
Accept-Encoding: gzip, deflate, br
Connection: keep-alive
Content-Length: 55
{
"username": "ry",
"password": "admin123"
给浏览器返回数据:
package com.note;
import java.io.*;
import java.net.*;
import java.nio.charset.StandardCharsets;
import java.util.concurrent.ArrayBlockingQueue;
import java.util.concurrent.ThreadPoolExecutor;
import java.util.concurrent.TimeUnit;
@SuppressWarnings("all")
public class SocketApplication {
public static void main(String[] args) throws IOException {
//创建线程池 开始执行BIO模式的B/S模式
ThreadPoolExecutor threadPoolExecutor = new ThreadPoolExecutor(
5,
10,
20,
TimeUnit.SECONDS,
new ArrayBlockingQueue<>(10)
);
System.out.println("TCP 接收端启动,监听 9999...");
try (ServerSocket serverSocket = new ServerSocket(9999)) {
while (true) {
//每次接收到数据就扔给一个线程来处理
Socket socket = serverSocket.accept();
//扔给线程池阻塞队列来处理
threadPoolExecutor.execute(() -> {
try {
controller(socket);
} catch (IOException e) {
throw new RuntimeException(e);
}
});
}
} catch (IOException e) {
e.printStackTrace();
}
}
public static void controller(Socket socket) throws IOException {
System.out.println(Thread.currentThread().getName() + "已启动");
try (
BufferedReader br = new BufferedReader(new InputStreamReader(socket.getInputStream())); //负责处理请求
PrintWriter writer = new PrintWriter(new OutputStreamWriter(socket.getOutputStream()), true) //负责返回响应
) {
//读取
String msg = null;
while ((msg = br.readLine()) != null && msg.length() > 0) {
System.out.println(msg);
}
String response = "<p>你好,世界</p>";
//返回 浏览器数据 封装http协议文本
writer.println("HTTP/1.1 200 OK"); // 状态行
writer.println("Content-Type: text/html; charset=utf-8"); // 响应头 设置文本格式与编码
writer.println("Content-Length: " + response.getBytes("UTF-8").length);
writer.println(); // 空行 表示响应体前的空行
writer.println(response); //响应体数据 html代码
} catch (IOException e) {
System.out.println(Thread.currentThread().getName() + "连接丢失或关闭");
}
System.out.println(Thread.currentThread().getName() + "连接正常关闭");
}
}
● NIO
NIO(Non-Blocking IO)非阻塞IO
相比于BIO的线程陪同(建立连接、)模式,NIO非常灵活它分为三个核心组件
- channel:相当于通讯管道,他不同于BIO这个管道既可以读,又可以写,他是浏览器前端与后端连接的桥梁,浏览器前端将http协议文本数据碾碎成二进制通过网络传入这个管道,到后端之后再将这个数据存到缓冲区(buffer),后端需要发数据时会将这个缓冲区的数据碾碎成二进制通过channel管道响应给浏览器前端
- buffer:缓冲区,它相当于装数据的集装箱,方便调度时,当做读写线程容易找到的一个区块,讲这个区块进行读写操作,如果是流的话会非常分散,读写线程不能灵活读取写入数据,BIO 的「流」= 水管直排:水(数据)只能单向流(读 / 写分开),要么从水管往外放(读),要么往水管里灌(写),而且必须等水流完才能停(阻塞),没法中途暂停;NIO 的「Buffer」= 蓄水池:先把水(数据)存到蓄水池里,想读就从池里取,想写就往池里灌,可暂停、可回退、可批量处理,还能切换读写方向,适配非阻塞的 “按需处理” 逻辑。
- selector:他有一个selector调度器也是监听器可以对建立连接、读写进行灵活调度,比如一个请求进来,selector(一个独立的线程或多个线程Reactor线程)就会被提醒有新的请求(监听到ACCEPT事件)调用accept创建连接,这时候,ACCEPT会注册一个READ事件让selector监听,当浏览器发送一个请求或数据时,这个READ事件就会监听到,然后selector的线程就会从channel读到数据给buffer,这时候将数据发给controller(后端线程池)进行处理,当处理完成后会注册一个WRITER事件监听会不会有数据返回,当返回数据时,selector线程就会将数据放进buffer通过channel返回浏览器前端
通俗来说ACCEPT、READ、WRITER都是传话的,真正干活的是selector线程(reactor)他一个人来处理建立连接,读写操作,他只负责快速的读写操作,真正复杂的业务逻辑就会丢给线程池(controller)来处理业务
假设有1万个并发请求同时连接后(ACCEPT忙完),用户待在页面内,这时候大部分用户都在浏览而不是发起操作请求,发请求的只有少部分或者是千分之一,所以reactor线程效率极高,因为大部分连接channel是空闲的,他就可以去处理那些真正有事的请求,真当reactor处理不过来时,就会实施多线程(reactor)分工来处理(tomcat NIO模型),这个多线程分工的线程池和业务线程池(controller)是完全解耦合的
与BIO的比较:
环境:成千上万的用户登录进页面后,执行操作时(这时已经accept建立了连接)
-
BIO阻塞式:监听线程在监听过程中会阻塞线程(accept),当他接收到一个请求后这个线程会一直陪同这个请求的读写操作,当他执行读写操作时,假如后面有成千上万的请求一起进,他就会无暇顾及,就是陪同一个请求走到死,无暇顾及其他请求
-
如果给BIO配个线程池也是治标不治本,BIO只能把读写操作交给线程池,交给线程池后,如果是短连接,请求完就断开了,那么用户在页面再次请求就会重新建立连接\断开连接,这样性能开销极大,如果是长连接,那么就会出现BIO经典阻塞问题,一个线程陪一个请求到死
-
长连接:每一个线程就会分配1MB那么1000个约等于1G,线程个开销极大
-
而NIO:每当一个连接进入时注册的ACCEPT通知selector,selector把连接注册到操作系统内,这是连接由操作系统来维护不会因为长连接占用线程而苦恼

○ NIO的C\S使用
△ API
Buffer API
| 类别 | 方法签名 | 说明 | 注意事项 |
|---|---|---|---|
| 静态工厂 | static ByteBuffer allocate(int capacity) |
在JVM堆内存分配指定容量的缓冲区 | 受GC管理,分配快,但数据在传输到I/O设备时可能需要复制到直接内存 |
static ByteBuffer allocateDirect(int capacity) |
分配直接内存缓冲区(操作系统内存) | 分配成本高,但能减少JVM与操作系统间的数据拷贝,适合长期存在或大文件操作 | |
static ByteBuffer wrap(byte[] array) |
将现有byte数组包装成缓冲区 | 缓冲区数据与数组共享同一内存,修改相互影响 | |
| 状态属性 | int capacity() |
返回缓冲区容量(不可变) | 创建后固定不变 |
int position() |
返回当前位置 | 下一个要读或写的位置索引 | |
Buffer position(int newPosition) |
设置新的位置 | 必须满足:0 ≤ newPosition ≤ limit | |
int limit() |
返回界限位置 | 第一个不能读或写的位置索引 | |
Buffer limit(int newLimit) |
设置新的界限 | 必须满足:0 ≤ newLimit ≤ capacity | |
| 核心操作 | Buffer flip() |
切换为读模式:limit = position; position = 0; |
最常用! 写入数据后必须调用才能读取 |
Buffer clear() |
"清空"缓冲区:position = 0; limit = capacity; |
不删除数据,只是重置指针,准备重新写入 | |
Buffer rewind() |
重绕缓冲区:position = 0; |
重新读取已写入的数据 | |
Buffer compact() |
压缩缓冲区:将未读数据复制到头部,position = remaining(); limit = capacity; |
处理粘包/拆包的关键方法,保留未处理数据 | |
Buffer mark() |
在当前position设置标记 |
可与reset()配合使用 |
|
Buffer reset() |
将position重置到上次标记的位置 |
如果未标记,抛出InvalidMarkException |
|
| 读写操作 | byte get() |
从当前位置读取一个字节,position加1 |
如果position ≥ limit,抛出BufferUnderflowException |
ByteBuffer get(byte[] dst) |
读取多个字节到目标数组 | 自动增加position,需要确保数组不会过大 |
|
byte get(int index) |
读取指定索引处的字节(绝对位置) | 不改变position,适合随机访问 |
|
ByteBuffer put(byte b) |
写入一个字节到当前位置,position加1 |
如果position ≥ limit,抛出BufferOverflowException |
|
ByteBuffer put(byte[] src) |
写入整个字节数组 | 自动增加position |
|
ByteBuffer put(int index, byte b) |
在指定索引处写入字节(绝对位置) | 不改变position |
|
| 辅助方法 | int remaining() |
返回剩余元素数:limit - position |
判断是否还有数据可读/可写 |
boolean hasRemaining() |
判断是否还有剩余元素 | position < limit |
|
boolean isDirect() |
判断是否为直接内存缓冲区 | 直接缓冲区可能带来性能优势 | |
boolean isReadOnly() |
判断是否为只读缓冲区 | 只读缓冲区不能修改内容 |
Channel API
| 类别 | 方法签名 | 说明 | 所属类/接口 |
|---|---|---|---|
| Channel基础 | boolean isOpen() |
判断通道是否打开 | Channel |
void close() |
关闭通道 | Channel |
|
| SelectableChannel | SelectableChannel configureBlocking(boolean block) |
设置阻塞模式 | SelectableChannel |
boolean isBlocking() |
判断是否为阻塞模式 | SelectableChannel |
|
SelectionKey register(Selector sel, int ops) |
向选择器注册,指定感兴趣的事件 | SelectableChannel |
|
SelectionKey register(Selector sel, int ops, Object att) |
注册时可附加一个附件对象 | SelectableChannel |
|
SelectionKey keyFor(Selector sel) |
返回此通道在指定选择器上的注册键 | SelectableChannel |
|
| ServerSocketChannel | static ServerSocketChannel open() |
打开服务器套接字通道 | ServerSocketChannel |
ServerSocketChannel bind(SocketAddress local) |
绑定到本地地址 | ServerSocketChannel |
|
SocketChannel accept() |
接受客户端连接,返回SocketChannel | ServerSocketChannel |
|
ServerSocket socket() |
获取对等的ServerSocket对象 | ServerSocketChannel |
|
| SocketChannel | static SocketChannel open() |
打开套接字通道 | SocketChannel |
boolean connect(SocketAddress remote) |
连接到远程地址 | SocketChannel |
|
boolean finishConnect() |
完成非阻塞连接过程 | SocketChannel |
|
boolean isConnected() |
判断是否已连接 | SocketChannel |
|
boolean isConnectionPending() |
判断连接是否正在进行中 | SocketChannel |
|
int read(ByteBuffer dst) |
从通道读取数据到缓冲区 | SocketChannel |
|
long read(ByteBuffer[] dsts) |
分散读取到多个缓冲区 | SocketChannel |
|
int write(ByteBuffer src) |
将缓冲区数据写入通道 | SocketChannel |
|
long write(ByteBuffer[] srcs) |
从多个缓冲区聚集写入 | SocketChannel |
|
| FileChannel | int read(ByteBuffer dst, long position) |
从文件指定位置读取 | FileChannel |
int write(ByteBuffer src, long position) |
向文件指定位置写入 | FileChannel |
|
long transferTo(long pos, long count, WritableByteChannel target) |
将数据从文件通道传输到其他通道 | FileChannel |
|
MappedByteBuffer map(MapMode mode, long position, long size) |
将文件区域映射到内存 | FileChannel |
Selector API
| 类别 | 方法签名 | 说明 | 注意事项 |
|---|---|---|---|
| 创建与关闭 | static Selector open() |
打开一个选择器 | 底层使用系统调用(如epoll) |
boolean isOpen() |
判断选择器是否打开 | ||
void close() |
关闭选择器并释放资源 | 同时取消所有注册的键 | |
| 核心选择操作 | int select() |
阻塞等待,直到至少一个通道就绪 | 返回就绪通道的数量 |
int select(long timeout) |
带超时的阻塞等待 | 超时时间以毫秒为单位 | |
int selectNow() |
非阻塞检查,立即返回 | 不阻塞,无论是否有通道就绪都立即返回 | |
Selector wakeup() |
使尚未返回的select()方法立即返回 |
另一个线程调用此方法可唤醒阻塞在select()的线程 |
|
| 键集合管理 | Set<SelectionKey> keys() |
返回所有已注册的键(不可修改) | 此集合为只读,直接修改会抛出异常 |
Set<SelectionKey> selectedKeys() |
返回已就绪的键集合 | 需要手动移除已处理的键,否则下次select()还会返回 |
|
| 其他 | SelectorProvider provider() |
返回创建此选择器的提供者 | 通常不需要直接使用 |
SelectionKey API
| 类别 | 方法签名 | 说明 | 注意事项 |
|---|---|---|---|
| 事件类型常量 | static int OP_ACCEPT |
连接接受事件(值:16) | ServerSocketChannel专用 |
static int OP_CONNECT |
连接完成事件(值:8) | SocketChannel连接时 |
|
static int OP_READ |
读就绪事件(值:1) | 最常见的事件类型 | |
static int OP_WRITE |
写就绪事件(值:4) | 通常只在发送缓冲区满时才关注 | |
| 核心属性 | SelectableChannel channel() |
返回关联的通道 | |
Selector selector() |
返回关联的选择器 | ||
| 事件操作 | int interestOps() |
返回当前感兴趣的事件集合 | 位掩码形式,如OP_READ | OP_WRITE |
SelectionKey interestOps(int ops) |
设置新的感兴趣事件 | 可以动态修改监听的事件类型 | |
int readyOps() |
返回已就绪的事件集合 | 由Selector.select()设置 |
|
| 状态检查 | boolean isAcceptable() |
是否连接可接受 | 等价于(readyOps() & OP_ACCEPT) != 0 |
boolean isConnectable() |
是否连接已完成 | ||
boolean isReadable() |
是否可读 | 注意:非阻塞模式下,即使可读,read()也可能返回0 |
|
boolean isWritable() |
是否可写 | 几乎总是返回true,除非发送缓冲区已满 |
|
| 附件与有效性 | Object attach(Object ob) |
附加一个对象到此键 | 常用于存储会话状态或缓冲区 |
Object attachment() |
获取附加的对象 | ||
boolean isValid() |
判断此键是否有效 | 通道关闭或取消注册后变为无效 | |
void cancel() |
取消此键的注册 | 在下一次select()时从选择器中移除 |
API作用:
- ServerSocketChannel:它相当于一个客服中心的转接台,负责监听有没有新连接,如果有新的连接进来之后,就把这个连接转接给专业的客服进行对接,他就相当于监听的作用,并没有处理事情的能力
- SocketChannel:他就相当于一个客户端与服务器的专用双向通道(全双工),这个通道只代表这个客户端与服务器的数据通道端口,相当于一个客户用一个聊天窗口和客户进行互动,这个窗口就是这个SocketChannel,但是这个客服可以处理多个窗口,这就是NIO典型的非阻塞模式,这个客服(reactor线程)不会死在一个通道上
- Buffer:缓冲区,因为客户端与服务器的网速不同,如果客户端慢,服务器快,客户端这边才传了一般数据服务器就开始处理了,那必然出错,所以先把数据放到缓冲区内,攒够了在处理,发送数据也是同理
- Selector:总调度中心,转接台(ServerSocketChannel)监听到事件了,就会通知selector调度器去分派客服处理这个事件,他负责将事件对应的事情分给某从reactor处理

NIO使用流程:
主线程管理组:
- 获取管道管理器(Selector)
- 创建管道监听器(ServerSocketChannel)
- 绑定端口
- 开启非阻塞模式
- 注册ACCEPT事件
- 循环监听干活(Selector.select())
- 如果请求到来判断监听类型执行相应处理
这个SelectionKey就是一个记录,内部:
- channel:记录了哪个客户端的连接
- selector:记录了主线程reactor
- InterestOps:事件类型
△ 服务器
reactor单线程使用
package com.note;
import com.note.properties.ServerProperties;
import java.io.IOException;
import java.net.InetSocketAddress;
import java.nio.channels.SelectionKey;
import java.nio.channels.Selector;
import java.nio.channels.ServerSocketChannel;
import java.util.Iterator;
import java.util.Set;
@SuppressWarnings("all")
public class ServerApplication {
public static void main(String[] args) throws IOException {
//创建selector多路复用:管理所有channel的
Selector selector = Selector.open();
//创建管道监听器
ServerSocketChannel serverSocketChannel = ServerSocketChannel.open();
//绑定端口
serverSocketChannel.bind(new InetSocketAddress(ServerProperties.SERVER_PORT));
//开启非阻塞模式
serverSocketChannel.configureBlocking(false);
//注册ACCEPT事件让selector可以监听连接请求
serverSocketChannel.register(selector, SelectionKey.OP_ACCEPT);
System.out.println("Server started on port " + ServerProperties.SERVER_PORT);
//reactor主线程循环监听事件
while (true) {
//让监听阻塞等待100ms 如果100ms没有东西进来 可以循环干点别的事 然后进入下次循环
int readyChannel = selector.select(100);
// ... 这中间可以干点别的事
if (readyChannel == 0) { //没有事件就进入下次循环
continue;
}
//如果有事件 这个事件就会进入selectedKeys内
Set<SelectionKey> selectionKeys = selector.selectedKeys();
//使用过iterator循环遍历
Iterator<SelectionKey> iterator = selectionKeys.iterator();
while (iterator.hasNext()) {
SelectionKey key = iterator.next();
//这个事件将要被处理了 那么处理完就要把这个事件删掉 不然会造成循环业务处理
//比如一个read事件 循环不断的通知selector去读数据处理数据,数据都为空了还处理就会出错
iterator.remove();
//开始处理事件
try {
if (key.isAcceptable()) { //ACCEPT事件
} else if (key.isReadable()) { //READ
} else if (key.isWritable()) { //WRITER
} else if (key.isConnectable()) { //CONNECT
}
} catch (Exception e) {
// 客户端断开连接时,释放资源 当客户端异常关闭时(断电什么的)
// 这时候再执行READ什么的 就会抛异常
System.err.println("客户端连接异常:" + e.getMessage());
key.cancel();
if (key.channel() != null) {
key.channel().close();
}
}
}
}
}
}
△ ACCEPT事件
处理ACCEPT事件:
这时selector监听到一个新的accept请求,开始建立新的客户端连接
流程:
- 从key获取本通道channel
- 非阻塞获取连接accept(默认长连接)
- 设置客户端通道为非阻塞模式(因为读写操作都是由reactor完成的,如果不设置成非阻塞模式当这个连接占着茅坑不拉屎的时候selector就没法工作了)
[非阻塞原理]:会进行内核检查,如果accept()发现目前没有连接,他不会傻傻的等待(占用线程),而是会直接返回null - 为客户端通道分配缓冲区大小
- 注册READ事件到selector
更新主线程代码
//开始处理事件
try {
if (key.isAcceptable()) { //ACCEPT事件
/**
* 这里key已经包含了selector为什么还要传一个selector
* 因为在NIO多线程模型下 多个reactor线程一起处理
* selector主线程只负责接待 selectorB负责建立新连接
* 如果都让主线程干活那么就变成单线程了,所以这里这样做是考虑了多线程环境
*/
handlerAccept(key, selector);
} else if (key.isReadable()) { //READ
} else if (key.isWritable()) { //WRITER
这里先不写WRITER事件 到后面线程池多模型Reactor再写
} else if (key.isConnectable()) { //CONNECT
}
} catch (Exception e) {
// 客户端断开连接时,释放资源 当客户端异常关闭时(断电什么的)
// 这时候再执行READ什么的 就会抛异常
System.err.println("客户端连接异常:" + e.getMessage());
key.cancel();
if (key.channel() != null) {
key.channel().close();
}
}
ACCEPT事件处理
/**
* 处理ACCEPT时间
*
* @param key 事件记录
* @param selector 要注册ACCEEPT的reactor线程
*/
private static void handlerAccept(SelectionKey key, Selector selector) throws IOException {
//获取selector监听的这一个channel连接通道
ServerSocketChannel serverSocketChannel = (ServerSocketChannel) key.channel();
//非阻塞获取客户端连接
SocketChannel clientChannel = serverSocketChannel.accept();
/**
* 这里判断clientChannel 为null是因为在reactor多模型下,会有很多线程(selector.select())来抢这个accept
* 如果发生线程安全问题,这里做一个保障
*/
if (clientChannel == null) {
return;
}
//为客户端设置非阻塞
clientChannel.configureBlocking(false);
//设置缓冲区大小
ByteBuffer buffer = ByteBuffer.allocate(1024);
//注册READ事件顺便设置该channel的缓冲区
clientChannel.register(selector, SelectionKey.OP_READ, buffer);
String clientAddress = clientChannel.getRemoteAddress().toString();
System.out.println("新客户端连接:" + clientAddress);
}
△ READ事件
READ事件处理
在主线程的try catch内添加handlerRead(key)
流程:
- 获取这个客户端的连接通道channel和缓冲区buffer
- 非阻塞计算数据指针 读到多少算多少
- 真正的读取数据flip()
- 根据请求来执行真正的业务丢给线程池(controller)
- 压缩数据 完成READ最后一步 compact()
读完数据正常完成的情况:
比如读数据hello:初始化指针[0, 1024] 网卡数据进入 [5, 1024]
flip()读数据 [0, 5] 慢慢读到[5 , 5]
compact()重置指针 [0, 1024] 准备下一次读数据
没读完数据情况:
假设hello这个数据非常大 一次读不完 或者客户端传的慢 现在读到了[5, 5]发现后面两个数据不完整,就把这个没读完的数据放到头部 将前面的数据读出去[ h e l l o ] -> [ l o l l o ] 坐标变为[ 2, 1024 ]这样下次读数据就可以凑一个完整包了 - 清空缓冲区clear() 【和compact二选一 大部分compact 这个不能存数据继续读】
这个清空缓冲区并不是真的清空而是重置缓冲区指针pos
比如上一次读数据进来一个"hello" 那么初始pos的初始坐标与最大坐标为[0, 1024] 网卡数据进来之后是[5, 1024]
在进行读数据flip(),这个坐标就会变成[0, 5]开始读 慢慢读到[5, 5]
如果不进行clear,那么数据再进来就变成[6, 5]显然超标了
重置之后变为[0, 1024],网卡开始进数据"world"变为[5, 1024] 这时候原本的hello就会被world覆盖掉
在进行读flip()就可以读到想要的数据了
其实还可以再优化一下reactor不让他处理截取数据发送请求操作,而是扔给线程池去做
/**
* 处理READ事件
*
* @param key 事件对象
*/
private static void handlerRead(SelectionKey key) throws IOException {
//获取该连接的channel对象
SocketChannel socketChannel = (SocketChannel) key.channel();
//获取缓冲区
ByteBuffer buffer = (ByteBuffer) key.attachment();
//非阻塞获取缓冲区指针
int read = socketChannel.read(buffer);
if (read == 0) { //没数据就return
return;
}
if (read < 0) { //小于0代表客户端断开连接
System.out.println(socketChannel.getRemoteAddress() + "客户端断开连接");
key.cancel(); //消除channel的事件key
socketChannel.close(); //关闭这个客户端连接通道
return;
}
//开始读数据
buffer.flip();
String clientMsg = StandardCharsets.UTF_8.decode(buffer).toString().trim();
String clientAddress = socketChannel.getRemoteAddress().toString();
System.out.println("收到客户端[" + clientAddress + "]消息:" + clientMsg);
//获取请求
String[] requestBody = clientMsg.split(" ");
if (requestBody[0].equals("POST") && requestBody[1].equals("/login")) {
//登录请求 把请求体扔进去
String[] split = clientMsg.split("\r\n\r\n"); //截取请求体 找\r\n\r\n空白行一下就是
String usernameAndPassword = split[1];
String[] usernameOrPassword = usernameAndPassword.split("\""); //截出用户名 密码
ServerConfig.THREAD_POOL_EXECUTOR.execute(new TaskController(usernameOrPassword[3], usernameOrPassword[7], key));
} else {
socketChannel.write(ByteBuffer.wrap(clientMsg.getBytes()));
}
//完成数据读取 压缩数据 等待下一次数据
buffer.compact();
}
简单的Task
package com.note.controller;
import com.fasterxml.jackson.core.JsonProcessingException;
import com.fasterxml.jackson.databind.ObjectMapper;
import com.note.properties.ServerProperties;
import java.io.IOException;
import java.nio.ByteBuffer;
import java.nio.channels.SelectionKey;
import java.nio.channels.SocketChannel;
import java.util.HashMap;
@SuppressWarnings("all")
public class TaskController implements Runnable {
private final String username;
private final String password;
private final SelectionKey key;
public TaskController(String username, String password, SelectionKey key) {
this.username = username;
this.password = password;
this.key = key;
}
/**
* 执行controller分配
*/
@Override
public void run() {
//准备写入返回数据
String body = "";
HashMap<String, Object> resultMap = new HashMap<>();
//登录
int login = new LoginController().login(username, password);
if (login == 0) {
resultMap.put("code", "401");
resultMap.put("msg", "登录失败");
} else {
resultMap.put("code", "200");
resultMap.put("msg", "登录成功");
}
//把返回数据封装为json格式
ObjectMapper mapper = new ObjectMapper();
try {
body = mapper.writeValueAsString(resultMap);
} catch (JsonProcessingException e) {
throw new RuntimeException(e);
}
String response = "HTTP/1.1 200 OK\r\n" + // 状态行
"Content-Type: application/json; charset=utf-8\r\n" + // 头信息
"Content-Length: " + body.getBytes().length + "\r\n" +
"\r\n" + // 空行(必须有!)
body; // 真正的身体
try {
//获取管道与缓冲区
SocketChannel socketChannel = (SocketChannel) key.channel();
ByteBuffer buffer = (ByteBuffer) key.attachment();
//写回数据
socketChannel.write(ByteBuffer.wrap(response.getBytes()));
//判断是否写完
if (buffer.hasRemaining()) {
//走到这里表示缓冲区满了 写不下了 这时候就需要注册WRITER事件了
System.out.println("注册WRITER事件");
socketChannel.register(key.selector(), SelectionKey.OP_WRITE);
key.attach(ByteBuffer.wrap(body.getBytes()));
}
} catch (IOException e) {
throw new RuntimeException(e);
}
}
}
简单的Controller
package com.note.controller;
public class LoginController {
/**
* 简易的登录逻辑
*
* @param username 用户名
* @param password 密码
* @return 结果 1成功 0失败
*/
public int login(String username, String password) {
System.out.println(username + " " + password);
if (username == null && password == null) {
return 0;
}
return 1;
}
}
△ WRITER事件
在数据写回时极大部分情况用的是直接channel.writer()直接写回,而极少情况用WRITER,如果一上来就注册WRITER他就会一直通知selector可以写,无限循环速度极快,CPU撑不住,除非是客户端网速极慢,比如一个200M文件需要发给客户端,客户端是2G网,他接受数据很慢,导致channel缓冲区满了,这时候就需要用到WRITER了,但是注册完后,数据处理完之后要立马取消注册WRITER不然就会出现CPU负载过高的情况,在这不过多赘述,下面主从reactor再写
△ 客户端
客户端的流程:
- 获取通道管理器
- 设置管道非阻塞
- 打开selector
- 非阻塞连接服务器(如果连上了就注册READ事件,如果没有连上就注册CONNECT事件等到连上了再通知selector完成连接)
- 实时监听READ事件或channel.writer()写满缓冲区注册WRITER事件
- 大体流程与服务器相同
package com.note;
import java.io.IOException;
import java.net.InetSocketAddress;
import java.nio.ByteBuffer;
import java.nio.channels.SelectionKey;
import java.nio.channels.Selector;
import java.nio.channels.SocketChannel;
import java.util.Iterator;
import java.util.Scanner;
public class NioClient {
// 选择器,类似“事件监听器”
private Selector selector;
public void start(String ip, int port) {
try {
// 1. 打开一个通道(Channel)
SocketChannel socketChannel = SocketChannel.open();
// > 架构师划重点:必须设置为非阻塞模式,否则就退化成普通 BIO 了
socketChannel.configureBlocking(false);
// 2. 打开选择器(Selector)
this.selector = Selector.open();
// 3. 连接服务器
// 注意:非阻塞模式下,connect 方法可能在连接未建立时就返回 false
if (socketChannel.connect(new InetSocketAddress(ip, port))) {
// 如果运气好直接连上了,注册“读”事件
socketChannel.register(selector, SelectionKey.OP_READ);
} else {
// 如果还没连上,注册“连接”事件,等连上了通知我
socketChannel.register(selector, SelectionKey.OP_CONNECT);
}
System.out.println("客户端启动,准备连接服务器...");
// 4. 开启事件循环(类似前端的 Event Loop)
// 在实际项目中,这通常会放在一个独立的线程中运行
new Thread(() -> {
while (true) {
try {
// 阻塞等待,直到至少有一个事件发生(比如收到数据了)
selector.select();
// 获取所有发生的事件
Iterator<SelectionKey> iterator = selector.selectedKeys().iterator();
while (iterator.hasNext()) {
SelectionKey key = iterator.next();
// > 架构师划重点:处理完必须移除,否则下次循环还会拿到这个已处理的事件
iterator.remove();
handleKey(key);
}
} catch (IOException e) {
e.printStackTrace();
}
}
}).start();
// 模拟发送数据
Scanner scanner = new Scanner(System.in);
while (scanner.hasNextLine()) {
String line = scanner.nextLine();
// 发送数据到服务器
socketChannel.write(ByteBuffer.wrap(line.getBytes()));
}
} catch (IOException e) {
e.printStackTrace();
}
}
private void handleKey(SelectionKey key) throws IOException {
if (key.isConnectable()) {
// 处理连接建立事件
SocketChannel channel = (SocketChannel) key.channel();
// 完成连接
if (channel.isConnectionPending()) {
channel.finishConnect();
}
// 连接成功后,改为关注“读”事件,准备接收服务器消息
channel.configureBlocking(false);
channel.register(selector, SelectionKey.OP_READ);
System.out.println("成功连接到服务器!");
} else if (key.isReadable()) {
// 处理读取数据事件
SocketChannel channel = (SocketChannel) key.channel();
// 分配缓冲区,假设一次读取 1024 字节
ByteBuffer buffer = ByteBuffer.allocate(1024);
// 将通道数据读入缓冲区
int len = channel.read(buffer);
if (len > 0) {
// > 架构师划重点:buffer 刚写完,现在要读它,必须 flip() 切换模式
// 就像把杯子倒过来才能把水倒出来
buffer.flip();
String content = new String(buffer.array(), 0, len);
System.out.println("收到服务器消息: " + content);
// 这里不需要 buffer.clear(),因为是局部变量,下次循环会 new 新的
} else if (len == -1) {
// 服务器关闭了连接
channel.close();
System.out.println("服务器关闭了连接");
}
}
}
public static void main(String[] args) {
// 确保你有一个服务器跑在 8080 端口,否则这里会一直尝试连接
new NioClient().start("127.0.0.1", 8080);
}
}
● NIO 主从模型Reactor
单reactor模型是一个reactor线程来处理(ACCEPT、READ、WRITER)事件,在高并发场景下难以支撑(适合rps<1000)
而多reactor模型是一个主线程MainReactor负责监听ACCEPT,从reactor(SubReactor)负责处理事件(rps>1000)
半包拆包&粘包问题:
- 半包拆包:当客户端发送给服务器数据时,通过TCP发送数据可能会将数据分成多段发送,这会导致服务器分批次收到数据,比如本来要发一个登录请求(/login用户名密码)数据达到服务端,可是被TCP拆包发送了,导致服务端不知道接受数据包的数据边界,比如第一次接收到用户名的包数据就开始处理了,密码数据包还没发过来,这时候处理就会出错。
半包拆包原因:MTU 限制:网络层有最大传输单元(比如 1500 字节),超过的数据包会被拆分 - 粘包问题:就是有多个请求发送时,服务端一次性收到了两个请求数据包,一个完整的包带着另一个完整的包,这时候处理就会出问题
粘包原因:Nagle 算法:TCP 会合并小数据包再发送,减少网络交互次数,或是服务端缓冲区数据没及时读取
轮询选择与分工:
-
轮询选择:为了保证从reactor线程池的负载均衡,使用轮询策略轮询分发任务,保证均衡处理量,这里使用AtomicInteger index = new AtomicInteger(0);初始化时值为0然后取模运算
比如
int currentIndex = index.getAndIncrement(); //先返回值再执行自增类似int++但这个getAndIncrement是线程安全的cpu内核级操作读+增一步完成
// ② 取模实现循环(0→1→2→0→1→2…)
int targetIndex = Math.abs(currentIndex) % workerList.size(); 这里用取绝对值是为了防止index增到超过int的最大类型开始溢出变成int最小值开始 -
分工:主reactor负责监听accpet事件,从reactor负责监听read、writer事件,监听到之后将任务分发给相应的selector处理

○ 实现
以下全部流程根据此笔记项目来实现:
[Java主从Reactor模型]
重要代码在server模块内
△ 主Reactor监听ACCEPT事件
- 先打开ServerSocketChannel负责监听客户端连接(channel事件),他只负责监听连接不干其他的
- 打开Selector调度器,然后ServerSocketChannel把accept事件注册到selector内:告诉selector去监听accpet事件,在执行register时会返回一个SelectionKey,这个SelectionKey同时也会被这个selector维护起来,这个key包含channel、关联的selector、兴趣事件key.interestOps()、key.readyOps()已触发事件等
- 这时候register事件会返回的accpetKey(SelectionKey)这个key会放到selector.keys()内,这个acceptkey绑定了ServerSocketChannel和这个selector,如果有新的连接进来这个acceptkey就会被放到selector.selectedKeys()内也就是下面阻塞select()那个itrator循环里。
Ps:只有第一次注册accept事件时生成一个accpet类型的key以后生成的要么是readKey或者是writerKey - 非阻塞式监听accept事件,如果有新连接从selector.selectedKeys()内拿到accpetKey(不管以后多少个新连接都用这一个accpetKey来拿),通过这个accpetKey拿到ServerSocketChannel,通过这个ServerSocketChannel.accept()接收新的连接SocketChannel,然后为这个新连接注册read事件生成read类型Key放到selector.keys()内,如果下次有read事件则会从keys放到selectedKeys内用iterator循环处理事件,处理完就remove

/**
* 处理accept事件
*/
public void acceptHandler(SelectionKey key) throws IOException {
//获取新连接的通道
ServerSocketChannel serverSocketChannel = (ServerSocketChannel) key.channel();
//创建新连接
SocketChannel clientChannel = serverSocketChannel.accept();
if (clientChannel == null) {
return;
}
//设置非阻塞客户端连接
clientChannel.configureBlocking(false);
//轮询器迭代
int currentCount = index.getAndIncrement() % NIOModel.subReactors.length;
SubReactor targetReactor = NIOModel.subReactors[currentCount];
//设置上下文对象
ChannelContext channelContext = new ChannelContext();
//让从reactor注册read事件准备监听对应连接的数据
clientChannel.register(targetReactor.selector, SelectionKey.OP_READ, channelContext);
log.info("{} 客户端建立连接", clientChannel.getRemoteAddress());
}
△ 主Reactor分发任务
- 从reactor线程池初始化,初始化每一个subreactor的selector,这个selector每个从reactor独有并专属
- 把每一个从reactor通过execute提交给线程池,这时候run方法内开始执行非阻塞循环监听select(100),此时selector.keys()内无事件,selector.selectedKeys()内也没有东西,所以是空循环等待
- 直到主reactor接收到accept事件,通过accpetkey拿到ServerSocketChannel,通过这个ServerSocketChannel.accpet()拿到新连接的通道,通过这个新通道注册一个read事件给从reactor(轮询机制),这时候从reactor内部的selector.keys()才有一个事件叫read事件,如果以后这个channel有数据进入后,不会再麻烦主reactor因为他keys里面没有read,从reactor才有read,这时候从reactor把keys里的read事件放到selectedKeys()内,开始iterator循环处理read事件,处理完在删掉,这个read事件的key就代表着那个新连接的channel和这个从reactor的关联凭证

△ 从Reactor解决粘包半包问题[读][与writer完全读写分离]
用标准的做法判断:判断http协议的请求头内存储的请求体长度也就是content-length
考虑如果第一次读没读到content-length的半包情况,就判断读请求头,判断请求数据是否包含\r\n\r\n这个就是空白行,请求头与请求体的分界线,如果读到这个\r\n\r\n就表示请求头完整,没有就存起来下次在拼起来
- 流程前置:创建channel上下文对象
/**
* 读写分离的写的队列 所有返回响应都从这个队列中拿buffer
* 没写完就把没写完buffer重新放到队列头
*/
private Deque<ByteBuffer> readyWriter = new LinkedList<>();
/**
* 读取数据时的缓冲区 所有读操作都从这个缓冲区读
*/
private ByteBuffer readyBuffer = ByteBuffer.allocate(SubReactorProperties.BUFFER_READ_SIZE);
/**
* 历史记录buffer用于存留粘包半包的残缺数据包
*/
private ByteBuffer historyBuffer = ByteBuffer.allocate(SubReactorProperties.BUFFER_READ_SIZE);
/**
* 请求体的长度
*/
private int contentLength = 0;
/**
* 记录已读的请求体数据指针
*/
private int bodyBytesRead = 0;
/**
* 状态枚举
*/
private StatusEnum sate = StatusEnum.HEADERS;
/**
* 轮询器 它的作用就是拼装请求头+请求体
*/
private AtomicInteger requestIndex = new AtomicInteger(0);
/**
* 存储 轮询器+ [请求头+请求体] 让外部可以拿到整体 他只存完整请求
*/
private HashMap<Integer, String> requestMap = new HashMap<>();
有限状态机FSM
有限状态机FSM(Finite State Machine)
就是一个对象他会按照一定规则在指定的几个状态下变来变去,不会出现额外的状态,而且在初始化时就会有一个状态
比如枚举类Enum
就是按照状态执行某方法来处理数据->切换状态接着换方法执行而已
/**
* 状态机状态枚举
*/
public enum StatusEnum {
HEADERS, // 正在解析请求头
BODY, // 正在解析请求体
COMPLETE, // 请求解析完成
ERROR // 解析出错
}

Read事件整体流程
- HandlerRead()执行:new Handler().readHandler(key);
- 获取该对应的readKey(在分发任务提过的)的管道channel通过这个获取SocketChannel
- 获取key的附件ChannelContext上下文对象
- 获取上下文的读缓冲区
- 开始读数据channel.read()返回读取的字节数,判断客户端断开连接否
- flip()开始读数据这时候应该从[XXX, 8192]到[0, XXX]了
- buffer读取真正的字节数据传给有限状态机FSM开始状态解析循环
- 一开始是Headers状态解析请求头
从请求头中寻找空白行[\r\n\r\n]找到返回true找不到返回false并设置content-length将buffer的postion移到空白行后方便读取请求体
● 如果没找到返回0然后HandlerRead()事件不做任何处理 把读缓冲区compact之后等待下次read事件拼凑数据
● 如果找到了就把状态从Headers设置为Body开始解析请求体,然后把请求头存入requestmap集合 - 开始解析请求体
记录需要读取的请求体长度(或许上次已经读过了没读完没达到content-length需求、半包)
记录读取位置方便下次半包合整包也就是content-length缺少的长度读完数据后
将Body状态转换为Complete状态,并将请求体与request的请求头合并 - 开始执行Complete分支记录已经处理的完整请求,执行将本次的content-length清空,读请求体的也清空,然后判断还有没有数据(粘包问题)如果有就在设置为state为Headers状态继续循环读
- 整体完事后返回完整处理请求的个数,外部Read事件开始循环个数请求将requestmap里完整的请求数据信息丢给业务线程池去处理
- buffer.compact();

Read事件重要代码:
/**
* 处理read事件
*
* @param key read事件
*/
public void readHandler(SelectionKey key) throws Exception {
log.info("{}号 从reactor处理数据", index);
//获取连接对象
SocketChannel clientChannel = (SocketChannel) key.channel();
//获取上下文对象
ChannelContext channelContext = (ChannelContext) key.attachment();
//从上下文对象获取read事件的专属读缓冲区
ByteBuffer buffer = channelContext.getReadyBuffer();
//读取数据指针
int read = clientChannel.read(buffer); //返回已经读取了的字节个数
if (read == 0) return; //代表没数据可以读
//-1代表客户端断开连接
if (read < 0) {
log.info("{} 客户端断开连接", clientChannel.getRemoteAddress());
key.attach(null);
key.cancel();
clientChannel.close();
return;
}
//走到这里代表有数据 进行读数据操作
buffer.flip();
//开始读数据
byte[] data = new byte[read]; //设置目标字节数组 传参用的
buffer.get(data); //读取数据到data数组
int handlerRequest = channelContext.feedAndParse(data);//有限状态机FSM开始处理粘包 半包问题 并读取请求体 返回完整的请求处理结果
//处理完整请求 他有可能是粘包数据(里面可能有好几个完整请求) 所以要循环处理
for (int i = 0; i < handlerRequest; i++) {
//解析http协议包 + 丢给业务线程池处理业务
log.info("{} 收到客户端完整数据\n {}", clientChannel.getRemoteAddress(), null);
log.info("{}号 从reactor处理完毕", index);
}
buffer.compact(); //重置指针开始读下一次数据
}
/**
* 接受数据到累积缓冲区 然后状态机循环进行状态更迭
*
* @return 已处理的请求个数 如果有一个完整的请求头+请求体则会返回>0
*/
public int feedAndParse(byte[] data) {
//已解析的完整请求个数
int parseCount = 0;
// 将数据放入历史累加缓存区
historyBuffer.put(data);
//准备读数据 开始解析请求头或请求体
historyBuffer.flip();
//hasRemaining表示历史累积缓冲区是否有数据 第二个判断表述状态还没崩溃
while (historyBuffer.hasRemaining() && sate != StatusEnum.ERROR) {
switch (sate) {
case HEADERS: //【解析请求头】
if (parseHeader()) {
sate = (contentLength > 0) ? StatusEnum.BODY : StatusEnum.COMPLETE; //如果这个请求不存在请求体则直接complete
} else {
log.info("触发残包事件");
historyBuffer.compact(); //保留当前数据 继续读下次数据 这个之后才会将buffer改为写模式 不然都是flip的读模式 新数据进不来
return parseCount; //没找到就继续监听read事件 等数据拼凑
}
break;
case BODY: //【解析请求体】
if (parseBody()) {
sate = StatusEnum.COMPLETE; //完整读完请求体了
} else {
log.info("触发残包事件");
historyBuffer.compact(); //保留当前数据 继续读下次数据 这个之后才会将buffer改为写模式 不然都是flip的读模式 新数据进不来
return parseCount; //请求体数据不符合content-length长度继续监听read事件 等数据拼凑
}
break;
case COMPLETE: //【处理已完成状态】
parseCount++; //请求个数增加
contentLength = 0; //初始化请求头长度 因为已经读完了请求体 不需要他了 要是有粘包 从0接着来
historyBuffer.position(historyBuffer.position() + 1); //防止半包后读脏数据 因为之前减1了
if (historyBuffer.hasRemaining()) { //这里代表读完一个完整的请求后 请求体后面还有东西,这说明有粘包存在
sate = StatusEnum.HEADERS; //切换为请求头接着读 直到read事件读完这个请求
} else {
//走到这里说明全部完整请求皆以处理完成 清了历史记录返回请求个数即可
historyBuffer.clear();
return parseCount;
}
break;
}
}
//返回完整解析结果数 如果返回0则表示有数据没读完 这时候外部不做任何处理静等下次read事件
return parseCount;
}
解析请求头重要代码:
/**
* 解析请求头 找空白行
*
* @return true代表找到了 false代表没找到
*/
private boolean parseHeader() {
//先虚找 不动指针 用临时指针代替 找到了 在移动指针位置
int startPosition = historyBuffer.position();
for (int i = startPosition; i < historyBuffer.limit(); i++) {
//判断空白行
if (historyBuffer.get(i) == '\r' && historyBuffer.get(i + 1) == '\n'
&& historyBuffer.get(i + 2) == '\r' && historyBuffer.get(i + 3) == '\n') {
//找到了 去根据这个请求头找content-length长度
contentLength = findContentLength(startPosition, i); //找length 传i 也就是请求头的指针位置
//找到了 返回true 移动指针位置
historyBuffer.position(i + 4);
if (contentLength == 0) {
historyBuffer.position(historyBuffer.position() - 1);
}
return true;
}
}
//没找到 返回false
return false;
}
/**
* 根据请求头指针 找context-length
* 为什么要传这个startPosition 因为这次的包可能是粘包 前面已经读了一个完整的请求包了
*
* @param i 请求头指针
*/
private int findContentLength(int startPosition, int i) {
//开读
byte[] data = new byte[i];
historyBuffer.position(startPosition);
historyBuffer.get(data); //把数据读给data
String headers = new String(data).toLowerCase(); //将读到的请求头转为字符串
requestMap.put(requestIndex.get(), headers); //将请求头加到map集合 方便外部查找
//找到content-length对应的数字
String[] headerLine = headers.split("\n");
for (String line : headerLine) {
if (line.startsWith("content-length:")) {
return Integer.parseInt(line.split(":")[1].trim()); //返回请求体长度
}
}
//兜底
log.error("没有请求体或出错");
return 0;
}
△ WRITER事件触发机制[写][与read完全读写分离]
- read事件解析完数据后将请求体传给业务层线程池处理
- 业务线程池处理完后将响应数据写入附件包上下文的readyWriter队列(加锁队列)尾部(每次写入响应就是一个ByteBuffer包)
- 调用selector.wakeup唤醒调度器(如果调度器醒着不用管),这一步主要是减少发送延迟 写不写都行,唤醒后优先处理writer事件 下面有写流程
- 到这里read的活就干完了 他只负责读数据处理粘包半包的读 返回数据概不负责也不处理
- 然后在selector的iterator那里循环writer-read-channel.writer三个步骤顺序执行
- 之后将上下文的响应数据(readyWriter)在第三步channel.writer(ByteBuffer)写回,写回之后
- 判断key的附件里面的readyWriter队列还有没有数据
- 如果有数据,就说明这个byteBuffer还有东西 重新把他放入队列头部
- 代表数据没写完,这时候使用key.interestOps() & SelectionKey.OP_WRITE判断当前key有read事件的同时还有没有注册writer事件
- 如果没有则注册writer事件,注册时候用key.interestOps() | SelectionKey.OP_WRITE,意思是保留read事件并加一个writer事件
- 这时候继续iterator循环读到writer事件执行writer处理:
判断有无数据可写,没有的话就注销writer事件key.interestOps() & ~SelectionKey.OP_WRITE
如果有的话就用channel写入,注:这个channel并不会和read的channel写入冲突,虽然是同一个channel,但是iterator循环有先后顺序,这个从reactor处理只是单个线程不会有多个线程冲突的情况 - channel写入完后判断还有没有数据 如果这时候真没了,就清空key附件的readyWriter缓冲区,然后注销writer事件

大体来说:
read事件读数据(半包粘包处理)
丢给业务线程池响应数据写入readyWriter队列唤醒selector调度器
这时候会出现两种情况:
1主线程本来就醒着
2主线程阻塞中
不管是哪一种情况都会循环执行到writer事件或者channel.writer()写回操作
如果channel.writer()写回之后还有剩余数据则说明缓冲区满了就注册writer事件
如果是writer事件写回之后还剩数据就继续循环(这个writer如果不注销就一直叫唤)
直到数据清理完之后注销writer

更多推荐




所有评论(0)