Skip to content

Selector 多路复用

Selector 是 Java NIO 网络编程的核心。它解决的问题是:

不想一个连接占一个线程,而是让一个线程监听很多连接,哪个连接有事件就处理哪个。

这就是多路复用。

学习目标

学完本页,你应该能说明:

  1. BIO 为什么连接多了会撑爆线程。
  2. Selector 为什么能让一个线程管理多个 Channel。
  3. OP_ACCEPTOP_READOP_WRITEOP_CONNECT 分别是什么。
  4. NIO 为什么是同步非阻塞。
  5. 为什么实际项目一般用 Netty,而不是手写复杂 NIO。
  6. OP_WRITE、空轮询、半包粘包这些坑怎么理解。

BIO 的连接模型

传统 Socket BIO 通常是一个连接分配一个线程。

mermaid
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() 等待事件,等到事件后遍历处理。

mermaid
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"]

核心思想:

  1. Channel 必须配置非阻塞:configureBlocking(false)
  2. Channel 注册到 Selector。
  3. Selector 监听事件。
  4. 有事件时才处理,没有事件时线程可以阻塞在 select(),而不是每个连接一个线程阻塞。

事件类型

事件含义常见 Channel
OP_ACCEPT服务端监听到新连接ServerSocketChannel
OP_CONNECT客户端连接建立完成SocketChannel
OP_READ有数据可读SocketChannel
OP_WRITE可以写数据SocketChannel

注意:OP_WRITE 不是“有数据要写”,而是“底层缓冲区现在大概率可写”。很多初学者一直注册 OP_WRITE,会导致 Selector 几乎一直被唤醒,CPU 飙高。

NIO 服务端流程

mermaid
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 --> F

Demo:最小 NIO Echo Server

这个 Demo 可以帮助理解 Selector,但真实项目不要直接复制它上生产。真实网络框架还要处理协议、半包、粘包、连接管理、异常、内存池、背压和线程模型。

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

测试:

bash
telnet 127.0.0.1 8080

输入文本后服务端会返回 echo: xxx

为什么这个 Demo 还不适合生产

问题说明生产怎么做
没处理半包一次 read 不一定读到完整业务消息自定义协议或长度字段
没处理粘包多条消息可能一起读到解码器按协议拆分
写操作简单粗暴非阻塞写可能只写一部分维护待写队列,必要时注册 OP_WRITE
单线程处理业务慢业务会卡住所有连接IO 线程和业务线程分离
Buffer 固定 1024消息超长会处理异常动态累积或协议限制
没有心跳死连接不能及时发现心跳检测和空闲关闭

这也是为什么商业项目通常用 Netty。Netty 把 Selector、事件循环、ChannelPipeline、ByteBuf、编解码、线程模型封装好了。

半包和粘包为什么一定会遇到

TCP 是字节流协议,不是消息协议。它只保证字节按顺序到达,不保证你一次 write 对应对方一次 read

业务上你以为发送了两条消息:

text
msg1 = "ORDER_CREATED"
msg2 = "ORDER_PAID"

网络和操作系统实际可能让接收端这样读到:

情况接收端一次 read 结果说明
正常看起来像一条ORDER_CREATED只是碰巧
半包ORDER_CREA,下一次 TED一条消息被拆开
粘包ORDER_CREATEDORDER_PAID多条消息粘在一起
半包 + 粘包ORDER_CREATEDORDER_,下一次 PAID最常见的复杂情况
mermaid
flowchart TD
    A["业务 write 消息 A"] --> C["TCP 字节流"]
    B["业务 write 消息 B"] --> C
    C --> D["接收端 read 一段字节"]
    D --> E{"这一段是否刚好一条业务消息"}
    E -- "不一定" --> F["需要协议解码"]

根因是:

  1. TCP 面向字节流,没有消息边界。
  2. 操作系统发送缓冲区可能合并多次写。
  3. 网络分片可能把一次写拆成多段。
  4. 接收端 Buffer 大小和读取时机不固定。
  5. 非阻塞 read 可能只读到当前已经到达的一部分字节。

所以 NIO 服务不能把“一次 read 返回的数据”当成“一条完整业务消息”。这就是为什么 Netty 里有各种 Decoder。

常见协议拆包方式

解决半包粘包,本质是给字节流加“消息边界”。

方式思路适合场景风险
固定长度每条消息固定 N 字节定长报文、设备协议浪费空间,扩展差
分隔符\n、特殊字符表示结尾文本协议、简单命令内容里不能随便出现分隔符
长度字段先写消息长度,再写消息体二进制协议、RPC、网关必须校验长度,防止恶意大包
JSON 行每行一个 JSON日志、调试协议文本解析成本较高

商业系统最常见的是长度字段协议:

text
+------------+------------------+
| length=12  | body bytes       |
+------------+------------------+

读取时先看 Buffer 里够不够 4 字节长度字段;够了再看消息体是否完整;不完整就保留半包,等下一次 read。

mermaid
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 中尽可能解析完整消息,半包保留到下次”。为了让原理清楚,代码只处理核心流程。

java
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、写大文件,其他连接的读写事件都要等。

错误模型:

mermaid
flowchart TD
    A["Selector 线程"] --> B["读取消息"]
    B --> C["直接查数据库"]
    C --> D["数据库慢,线程阻塞"]
    D --> E["其他连接事件无法处理"]

推荐模型:

mermaid
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 没数据时可以立即返回,不会一直挂住线程
同步数据拷贝仍然由发起调用的线程完成,调用方要主动处理读写
mermaid
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 会不断返回写事件。

错误思路:

java
channel.register(selector, SelectionKey.OP_READ | SelectionKey.OP_WRITE);

正确思路:

  1. 平时只关心 OP_READ
  2. 当确实有数据没写完,才临时加上 OP_WRITE
  3. 待写数据写完后,移除 OP_WRITE
mermaid
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 -- "否" --> G

Selector 空轮询

历史上 JDK 某些版本在特定系统条件下可能出现 Selector 空轮询问题:select() 没有事件也不断返回,导致 CPU 飙高。

生产排查思路:

  1. top 看 Java 进程 CPU。
  2. jstack 看线程栈是否长期停在 selector 循环。
  3. 查看 NIO/Netty 版本和 JDK 版本。
  4. 如果是 Netty,优先升级成熟版本,Netty 对这类问题有重建 Selector 等规避策略。

不要自己在业务里随便 while(true) 空转轮询。

商业场景:长连接采集网关

医疗设备、物联网设备、日志 Agent、IM 客户端,都可能和服务端保持长连接。

mermaid
flowchart TD
    A["设备连接"] --> B["网关 NIO / Netty"]
    B --> C["心跳检测"]
    B --> D["协议解码"]
    D --> E["业务线程池处理"]
    E --> F["写入 MQ"]
    F --> G["后端服务异步消费"]

为什么适合 NIO/Netty?

  1. 连接多,但每个连接不是一直有数据。
  2. 用一个连接一个线程会浪费资源。
  3. 需要心跳、拆包、限流、背压、连接状态管理。
  4. Netty 已经封装了成熟的事件循环和编解码模型。

常见坑

后果正确做法
Channel 忘记非阻塞注册 Selector 报错或行为异常configureBlocking(false)
selectedKey 不 remove重复处理同一个事件遍历时 iterator.remove()
一直监听 OP_WRITECPU 飙高只有有剩余待写数据时才注册
在 IO 线程做慢业务所有连接延迟变高丢给业务线程池处理
不处理 read 返回 -1连接关闭后资源泄漏cancel key 并 close channel
不处理半包粘包消息错乱明确协议和解码器

面试标准回答

Selector 是 Java NIO 的多路复用器。传统 BIO 通常一个连接一个线程,连接很多时线程成本高;NIO 把 Channel 配置成非阻塞后注册到 Selector,由一个线程监听多个 Channel 的 accept、read、write 等事件,有事件时再处理。NIO 常见网络模型属于同步非阻塞:读写调用不会一直阻塞,但数据拷贝仍由调用线程完成。生产中一般不用手写复杂 NIO,而是使用 Netty,因为还要处理粘包拆包、内存池、线程模型、心跳、背压和异常关闭。

关联知识点