Selector 多路复用
Selector 是 Java NIO 网络编程的核心。它解决的问题是:
不想一个连接占一个线程,而是让一个线程监听很多连接,哪个连接有事件就处理哪个。
这就是多路复用。
学习目标
学完本页,你应该能说明:
- BIO 为什么连接多了会撑爆线程。
- Selector 为什么能让一个线程管理多个 Channel。
OP_ACCEPT、OP_READ、OP_WRITE、OP_CONNECT分别是什么。- NIO 为什么是同步非阻塞。
- 为什么实际项目一般用 Netty,而不是手写复杂 NIO。
- OP_WRITE、空轮询、半包粘包这些坑怎么理解。
BIO 的连接模型
传统 Socket BIO 通常是一个连接分配一个线程。
flowchart TD
A["客户端 1"] --> B["服务端线程 1"]
C["客户端 2"] --> D["服务端线程 2"]
E["客户端 3"] --> F["服务端线程 3"]
G["客户端 N"] --> H["服务端线程 N"]连接少时很好理解;连接多时问题明显:
| 问题 | 后果 |
|---|---|
| 线程数量膨胀 | 内存占用高,线程上下文切换多 |
| 慢客户端占住线程 | 工作线程无法处理其他连接 |
| 阻塞读写不可控 | 超时、断连、半包处理复杂 |
| 高并发长连接成本高 | 聊天、网关、采集长连接场景很吃亏 |
这不是说 BIO 不能用。文件读写、内部工具、小并发 Socket,用 BIO 很正常。问题是高并发连接场景下,BIO 的线程模型成本太高。
Selector 模型
Selector 让 Channel 注册自己关心的事件。一个线程调用 select() 等待事件,等到事件后遍历处理。
flowchart TD
A["ServerSocketChannel"] --> D["Selector"]
B["SocketChannel 1"] --> D
C["SocketChannel 2"] --> D
E["SocketChannel N"] --> D
D --> F["select 等待事件"]
F --> G["selectedKeys"]
G --> H["处理 accept / read / write"]核心思想:
- Channel 必须配置非阻塞:
configureBlocking(false)。 - Channel 注册到 Selector。
- Selector 监听事件。
- 有事件时才处理,没有事件时线程可以阻塞在
select(),而不是每个连接一个线程阻塞。
事件类型
| 事件 | 含义 | 常见 Channel |
|---|---|---|
OP_ACCEPT | 服务端监听到新连接 | ServerSocketChannel |
OP_CONNECT | 客户端连接建立完成 | SocketChannel |
OP_READ | 有数据可读 | SocketChannel |
OP_WRITE | 可以写数据 | SocketChannel |
注意:OP_WRITE 不是“有数据要写”,而是“底层缓冲区现在大概率可写”。很多初学者一直注册 OP_WRITE,会导致 Selector 几乎一直被唤醒,CPU 飙高。
NIO 服务端流程
flowchart TD
A["创建 Selector"] --> B["打开 ServerSocketChannel"]
B --> C["配置非阻塞"]
C --> D["绑定端口"]
D --> E["注册 OP_ACCEPT"]
E --> F["select 等待事件"]
F --> G{"有事件吗"}
G -- "否" --> F
G -- "是" --> H["遍历 selectedKeys"]
H --> I{"事件类型"}
I -- "ACCEPT" --> J["接收连接并注册 OP_READ"]
I -- "READ" --> K["读取数据到 ByteBuffer"]
I -- "WRITE" --> L["写出待发送数据"]
J --> F
K --> F
L --> FDemo:最小 NIO Echo Server
这个 Demo 可以帮助理解 Selector,但真实项目不要直接复制它上生产。真实网络框架还要处理协议、半包、粘包、连接管理、异常、内存池、背压和线程模型。
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.ServerSocketChannel;
import java.nio.channels.SocketChannel;
import java.nio.charset.StandardCharsets;
import java.util.Iterator;
public class NioEchoServer {
public static void main(String[] args) throws IOException {
Selector selector = Selector.open();
ServerSocketChannel server = ServerSocketChannel.open();
server.configureBlocking(false);
server.bind(new InetSocketAddress(8080));
server.register(selector, SelectionKey.OP_ACCEPT);
System.out.println("NIO echo server started at 8080");
while (true) {
selector.select();
Iterator<SelectionKey> iterator = selector.selectedKeys().iterator();
while (iterator.hasNext()) {
SelectionKey key = iterator.next();
iterator.remove();
if (!key.isValid()) {
continue;
}
if (key.isAcceptable()) {
ServerSocketChannel serverChannel = (ServerSocketChannel) key.channel();
SocketChannel client = serverChannel.accept();
client.configureBlocking(false);
client.register(selector, SelectionKey.OP_READ, ByteBuffer.allocate(1024));
System.out.println("client connected: " + client.getRemoteAddress());
} else if (key.isReadable()) {
SocketChannel client = (SocketChannel) key.channel();
ByteBuffer buffer = (ByteBuffer) key.attachment();
int len;
try {
len = client.read(buffer);
} catch (IOException ex) {
key.cancel();
client.close();
continue;
}
if (len == -1) {
key.cancel();
client.close();
continue;
}
buffer.flip();
String message = StandardCharsets.UTF_8.decode(buffer).toString();
System.out.println("receive: " + message.trim());
ByteBuffer response = StandardCharsets.UTF_8.encode("echo: " + message);
while (response.hasRemaining()) {
client.write(response);
}
buffer.clear();
}
}
}
}
}测试:
telnet 127.0.0.1 8080输入文本后服务端会返回 echo: xxx。
为什么这个 Demo 还不适合生产
| 问题 | 说明 | 生产怎么做 |
|---|---|---|
| 没处理半包 | 一次 read 不一定读到完整业务消息 | 自定义协议或长度字段 |
| 没处理粘包 | 多条消息可能一起读到 | 解码器按协议拆分 |
| 写操作简单粗暴 | 非阻塞写可能只写一部分 | 维护待写队列,必要时注册 OP_WRITE |
| 单线程处理业务 | 慢业务会卡住所有连接 | IO 线程和业务线程分离 |
| Buffer 固定 1024 | 消息超长会处理异常 | 动态累积或协议限制 |
| 没有心跳 | 死连接不能及时发现 | 心跳检测和空闲关闭 |
这也是为什么商业项目通常用 Netty。Netty 把 Selector、事件循环、ChannelPipeline、ByteBuf、编解码、线程模型封装好了。
半包和粘包为什么一定会遇到
TCP 是字节流协议,不是消息协议。它只保证字节按顺序到达,不保证你一次 write 对应对方一次 read。
业务上你以为发送了两条消息:
msg1 = "ORDER_CREATED"
msg2 = "ORDER_PAID"网络和操作系统实际可能让接收端这样读到:
| 情况 | 接收端一次 read 结果 | 说明 |
|---|---|---|
| 正常看起来像一条 | ORDER_CREATED | 只是碰巧 |
| 半包 | ORDER_CREA,下一次 TED | 一条消息被拆开 |
| 粘包 | ORDER_CREATEDORDER_PAID | 多条消息粘在一起 |
| 半包 + 粘包 | ORDER_CREATEDORDER_,下一次 PAID | 最常见的复杂情况 |
flowchart TD
A["业务 write 消息 A"] --> C["TCP 字节流"]
B["业务 write 消息 B"] --> C
C --> D["接收端 read 一段字节"]
D --> E{"这一段是否刚好一条业务消息"}
E -- "不一定" --> F["需要协议解码"]根因是:
- TCP 面向字节流,没有消息边界。
- 操作系统发送缓冲区可能合并多次写。
- 网络分片可能把一次写拆成多段。
- 接收端 Buffer 大小和读取时机不固定。
- 非阻塞
read可能只读到当前已经到达的一部分字节。
所以 NIO 服务不能把“一次 read 返回的数据”当成“一条完整业务消息”。这就是为什么 Netty 里有各种 Decoder。
常见协议拆包方式
解决半包粘包,本质是给字节流加“消息边界”。
| 方式 | 思路 | 适合场景 | 风险 |
|---|---|---|---|
| 固定长度 | 每条消息固定 N 字节 | 定长报文、设备协议 | 浪费空间,扩展差 |
| 分隔符 | 用 \n、特殊字符表示结尾 | 文本协议、简单命令 | 内容里不能随便出现分隔符 |
| 长度字段 | 先写消息长度,再写消息体 | 二进制协议、RPC、网关 | 必须校验长度,防止恶意大包 |
| JSON 行 | 每行一个 JSON | 日志、调试协议 | 文本解析成本较高 |
商业系统最常见的是长度字段协议:
+------------+------------------+
| length=12 | body bytes |
+------------+------------------+读取时先看 Buffer 里够不够 4 字节长度字段;够了再看消息体是否完整;不完整就保留半包,等下一次 read。
flowchart TD
A["SocketChannel read 到 Buffer"] --> B{"够不够读取 length"}
B -- "不够" --> C["compact 保留半包"]
B -- "够" --> D["读取 length"]
D --> E{"body 是否完整"}
E -- "不完整" --> C
E -- "完整" --> F["取出一条完整消息"]
F --> G{"Buffer 中还有数据吗"}
G -- "有" --> B
G -- "无" --> H["clear 或 compact"]Demo:按长度字段拆包
下面这个 Demo 演示“从 ByteBuffer 中尽可能解析完整消息,半包保留到下次”。为了让原理清楚,代码只处理核心流程。
import java.nio.ByteBuffer;
import java.nio.charset.StandardCharsets;
public class LengthFieldFrameDemo {
public static void main(String[] args) {
ByteBuffer buffer = ByteBuffer.allocate(1024);
writeFrame(buffer, "ORDER_CREATED");
writeFrame(buffer, "ORDER_PAID");
buffer.flip();
readFrames(buffer);
}
private static void writeFrame(ByteBuffer buffer, String text) {
byte[] body = text.getBytes(StandardCharsets.UTF_8);
buffer.putInt(body.length);
buffer.put(body);
}
private static void readFrames(ByteBuffer buffer) {
while (true) {
if (buffer.remaining() < 4) {
buffer.compact();
return;
}
buffer.mark();
int length = buffer.getInt();
if (length < 0 || length > 1024 * 1024) {
throw new IllegalArgumentException("非法消息长度: " + length);
}
if (buffer.remaining() < length) {
buffer.reset();
buffer.compact();
return;
}
byte[] body = new byte[length];
buffer.get(body);
System.out.println(new String(body, StandardCharsets.UTF_8));
}
}
}这里有几个关键点:
| 代码 | 原理 |
|---|---|
remaining() < 4 | 连长度字段都不完整,不能解析 |
mark() | 先标记读取位置,后面发现 body 不完整要回退 |
| 校验 length | 防止异常数据申请超大内存 |
reset() | 半包时回到 length 之前 |
compact() | 保留未读字节,把 Buffer 切回写模式 |
这也是为什么前面说 compact() 不是可有可无。只要处理网络协议,就很容易遇到“读到一半”的情况。
IO 线程和业务线程为什么要分离
Selector 所在线程负责处理 IO 事件。如果你在这个线程里查数据库、调用 HTTP、写大文件,其他连接的读写事件都要等。
错误模型:
flowchart TD
A["Selector 线程"] --> B["读取消息"]
B --> C["直接查数据库"]
C --> D["数据库慢,线程阻塞"]
D --> E["其他连接事件无法处理"]推荐模型:
flowchart TD
A["Selector / EventLoop 线程"] --> B["读取和解码"]
B --> C["投递到业务线程池"]
C --> D["业务处理 DB / HTTP"]
D --> E["生成响应"]
E --> F["回到 IO 线程写出"]但这又引出背压问题:业务线程池必须有界。如果业务处理不过来,IO 线程不能无限把消息丢进队列,否则只是把连接上的压力搬到 JVM 内存里。
生产框架 Netty 会把这些问题系统化:EventLoop 管 IO,Pipeline 管处理链,Decoder 负责拆包,业务线程池处理慢业务,写队列和水位线负责背压。
NIO 为什么叫同步非阻塞
以常见 Java NIO 网络用法为例:
| 维度 | 说明 |
|---|---|
| 非阻塞 | SocketChannel.read 没数据时可以立即返回,不会一直挂住线程 |
| 同步 | 数据拷贝仍然由发起调用的线程完成,调用方要主动处理读写 |
flowchart TD
A["Selector 通知可读"] --> B["线程调用 channel.read"]
B --> C["内核把数据复制到 ByteBuffer"]
C --> D["read 返回读取字节数"]
D --> E["业务线程解析数据"]真正异步 IO 是调用后立即返回,操作系统完成后通过回调或 Future 通知。
OP_WRITE 为什么容易导致 CPU 高
很多人以为“我要写数据,所以注册 OP_WRITE”。这是错的。
OP_WRITE 表示 Socket 发送缓冲区可写。大多数时候它都是可写的,所以 Selector 会不断返回写事件。
错误思路:
channel.register(selector, SelectionKey.OP_READ | SelectionKey.OP_WRITE);正确思路:
- 平时只关心
OP_READ。 - 当确实有数据没写完,才临时加上
OP_WRITE。 - 待写数据写完后,移除
OP_WRITE。
flowchart TD
A["业务产生响应"] --> B["尝试直接 write"]
B --> C{"全部写完"}
C -- "是" --> D["继续只监听 OP_READ"]
C -- "否" --> E["保存剩余数据"]
E --> F["注册 OP_WRITE"]
F --> G["下次可写时继续写"]
G --> H{"写完了吗"}
H -- "是" --> I["移除 OP_WRITE"]
H -- "否" --> GSelector 空轮询
历史上 JDK 某些版本在特定系统条件下可能出现 Selector 空轮询问题:select() 没有事件也不断返回,导致 CPU 飙高。
生产排查思路:
top看 Java 进程 CPU。jstack看线程栈是否长期停在 selector 循环。- 查看 NIO/Netty 版本和 JDK 版本。
- 如果是 Netty,优先升级成熟版本,Netty 对这类问题有重建 Selector 等规避策略。
不要自己在业务里随便 while(true) 空转轮询。
商业场景:长连接采集网关
医疗设备、物联网设备、日志 Agent、IM 客户端,都可能和服务端保持长连接。
flowchart TD
A["设备连接"] --> B["网关 NIO / Netty"]
B --> C["心跳检测"]
B --> D["协议解码"]
D --> E["业务线程池处理"]
E --> F["写入 MQ"]
F --> G["后端服务异步消费"]为什么适合 NIO/Netty?
- 连接多,但每个连接不是一直有数据。
- 用一个连接一个线程会浪费资源。
- 需要心跳、拆包、限流、背压、连接状态管理。
- Netty 已经封装了成熟的事件循环和编解码模型。
常见坑
| 坑 | 后果 | 正确做法 |
|---|---|---|
| Channel 忘记非阻塞 | 注册 Selector 报错或行为异常 | configureBlocking(false) |
| selectedKey 不 remove | 重复处理同一个事件 | 遍历时 iterator.remove() |
| 一直监听 OP_WRITE | CPU 飙高 | 只有有剩余待写数据时才注册 |
| 在 IO 线程做慢业务 | 所有连接延迟变高 | 丢给业务线程池处理 |
| 不处理 read 返回 -1 | 连接关闭后资源泄漏 | cancel key 并 close channel |
| 不处理半包粘包 | 消息错乱 | 明确协议和解码器 |
面试标准回答
Selector 是 Java NIO 的多路复用器。传统 BIO 通常一个连接一个线程,连接很多时线程成本高;NIO 把 Channel 配置成非阻塞后注册到 Selector,由一个线程监听多个 Channel 的 accept、read、write 等事件,有事件时再处理。NIO 常见网络模型属于同步非阻塞:读写调用不会一直阻塞,但数据拷贝仍由调用线程完成。生产中一般不用手写复杂 NIO,而是使用 Netty,因为还要处理粘包拆包、内存池、线程模型、心跳、背压和异常关闭。
