前言
在前面的文章中,我们已经用 Java NIO 写出了一个单线程服务器:
while (true) {
selector.select();
Iterator<SelectionKey> iterator =
selector.selectedKeys().iterator();
while (iterator.hasNext()) {
SelectionKey key = iterator.next();
iterator.remove();
if (key.isAcceptable()) {
handleAccept(key);
}
if (key.isReadable()) {
handleRead(key);
}
}
}
一个线程配合一个 Selector,就可以管理大量连接:
一个线程
↓
一个 Selector
↓
多个 Channel
相比传统 BIO 的“一连接一线程”,线程数量已经大幅减少。
但新的问题随之出现:
一个线程既要接收连接,又要读取数据,还要执行编解码和业务逻辑,它忙得过来吗?
某个客户端的处理逻辑阻塞五秒,其他连接会不会一起等待?
能不能创建多个 Selector,让每个线程处理一部分连接?
新连接由监听线程接收后,应该交给哪个 Selector?
Worker 线程正在阻塞执行
select(),Boss 怎样把新 Channel 注册进去?
直接在 Boss 线程调用
channel.register(workerSelector)是否安全?
Selector.wakeup()到底唤醒了什么?
应该先注册、先入队,还是先调用
wakeup()?
为什么需要一条线程安全队列?
一个 Channel 可以在多个线程间来回迁移吗?
Netty 的
NioEventLoop、NioEventLoopGroup、Boss 和 Worker,分别对应什么?
这些问题的背后,是从单线程多路复用到多线程 Reactor 的关键一步:
连接由谁接收?
↓
连接由谁管理?
↓
连接怎样跨线程交接?
这一步看起来只是在单线程代码外面多加了几个线程,实际却引入了:
线程归属
任务调度
跨线程通信
Selector 唤醒
Channel 注册
负载分配
等一整套并发问题。
一、单线程的边界
1. 一个事件循环
单线程 Selector 的核心可以概括为:
select
↓
处理就绪事件
↓
执行其他任务
↓
再次 select
这就是最简单的 EventLoop:
while (!stopped) {
selector.select();
processSelectedKeys();
runTasks();
}
只要每次事件处理都足够快,一个线程确实可以管理很多连接。
2. 一个慢任务
假设这个 EventLoop 管理一万个客户端。
连接 A 发送了一条消息,处理函数开始查询数据库:
private void handleRead(SelectionKey key) {
byte[] request = readRequest(key);
// 阻塞 5 秒
Result result = database.query(request);
writeResponse(key, result);
}
由于整个 EventLoop 只有一个线程:
连接 A 查询数据库 5 秒
↓
连接 B 的读事件无法处理
连接 C 的写事件无法处理
连接 D 的新数据只能继续等待
Selector 可以发现很多 Channel 已经就绪,但真正处理这些事件的仍然只有一个线程。
因此:
多路复用解决的是“怎样等待大量连接”,并没有自动解决“怎样并行处理大量任务”。
3. 多线程方向
可以创建多个 SelectorThread:
SelectorThread 0
└── Selector 0
├── Channel A
├── Channel B
└── Channel C
SelectorThread 1
└── Selector 1
├── Channel D
├── Channel E
└── Channel F
SelectorThread 2
└── Selector 2
├── Channel G
├── Channel H
└── Channel I
不同线程可以同时运行在不同 CPU 核心上。
但这也意味着:
新连接建立后,必须选择其中一个 Selector,并把 Channel 安全地交给它。
二、三种 Reactor
为了理解多线程结构,可以把 Reactor 简化成三种形态。
这是一种便于学习的分类,并不是所有框架都必须使用完全相同的名称。
1. 单线程模型
一个 EventLoop
├── accept
├── read
├── write
└── task
图示:
客户端
↓
SelectorThread
├── ServerSocketChannel
├── SocketChannel A
├── SocketChannel B
└── SocketChannel C
优点:
结构简单
没有跨 EventLoop 交接
大部分状态不需要加锁
缺点:
只能使用一个线程
任何慢任务都会拖住所有连接
2. 混合模型
创建多个 EventLoop,每个 EventLoop 既可以监听连接,也可以处理客户端 I/O:
EventLoop 0
├── accept
├── read
└── write
EventLoop 1
├── accept
├── read
└── write
EventLoop 2
├── accept
├── read
└── write
这种模型没有严格区分 Boss 和 Worker。
它可以利用多个线程,但监听 Channel 和客户端 Channel 的职责会混在一起,管理和负载分配也更复杂。
3. 主从模型
将连接接收和数据读写分开:
Boss Group
↓
只负责 accept
↓
把 SocketChannel 交给 Worker
Worker Group
↓
负责 read / write
完整结构:
┌── Worker EventLoop 0
客户端 ─→ Boss ───┼── Worker EventLoop 1
├── Worker EventLoop 2
└── Worker EventLoop 3
原素材也把 Boss 定义为接收连接的一方,把已建立连接交给 Worker 的 Selector 处理读写。
Netty 的官方示例同样使用两个 NioEventLoopGroup:第一个通常被称为 Boss,负责接受连接;第二个通常被称为 Worker,在连接被接受并注册后负责其后续流量。
三、Boss 做什么?
Boss 的职责应该尽可能简单:
监听端口
↓
发现 OP_ACCEPT
↓
调用 accept()
↓
得到 SocketChannel
↓
选择一个 Worker
↓
提交注册任务
对应代码:
private void handleAccept(SelectionKey key) throws IOException {
ServerSocketChannel server =
(ServerSocketChannel) key.channel();
SocketChannel client;
while ((client = server.accept()) != null) {
client.configureBlocking(false);
workerGroup.register(client);
}
}
1. 只处理连接
Boss 不应该直接执行:
消息解码
数据库查询
文件读取
复杂计算
长时间业务逻辑
否则 Boss 无法及时再次调用 accept(),新连接只能继续堆积在内核连接队列中。
2. accept 后仍需设置非阻塞
即使 ServerSocketChannel 已经设置为非阻塞:
server.configureBlocking(false);
accept() 返回的新 SocketChannel 仍然默认处于阻塞模式。
Java 官方文档明确规定,ServerSocketChannel.accept() 返回的连接 Channel 始终为阻塞模式,与监听 Channel 的模式无关。
因此必须再次调用:
client.configureBlocking(false);
否则将它注册到 Selector 时会抛出:
IllegalBlockingModeException
可选择 Channel 必须处于非阻塞模式,才能注册到 Selector。
四、Worker 做什么?
Worker 负责管理已经建立的连接:
注册 OP_READ
↓
等待可读事件
↓
读取数据
↓
协议解析
↓
提交业务任务
↓
按需注册 OP_WRITE
一个 Worker EventLoop 通常管理多个 Channel:
Worker 0
├── Channel A
├── Channel B
├── Channel C
└── Channel D
Netty 的 EventLoop 文档明确说明:一个 Channel 注册后,其 I/O 操作由对应 EventLoop 处理;一个 EventLoop 通常会管理多个 Channel。
Netty 的 NioEventLoop 本身也是一个 SingleThreadEventLoop,内部使用 Selector 注册并多路复用多个 Channel。
1. Channel 归属
在本文模型中,一条连接被分配给某个 Worker 后:
Channel A
↓
固定归属 Worker 1
后续操作尽量都在 Worker 1 中完成:
注册事件
修改 interestOps
读取数据
写出数据
关闭连接
更新连接状态
这种线程归属能够减少:
锁竞争
状态并发修改
执行顺序混乱
内存可见性问题
需要注意,Java 的 SelectableChannel 本身允许被多个线程并发使用,也允许注册到不同 Selector;“一个 Channel 归属一个 EventLoop”是 Reactor 和 Netty 的架构约束,不是 JDK 强制规定。
五、注册难在哪?
现在出现了核心问题。
Boss 线程执行:
SocketChannel client = server.accept();
选择了 Worker 1。
但 Worker 1 此时可能正在执行:
workerSelector.select();
结构如下:
Boss 线程 Worker 线程
accept()
↓
得到 client selector.select()
↓ ↓
想把 client 注册进去 正在内核中阻塞等待
应该怎样交接?
1. 直接注册
一种直觉写法是:
client.register(
workerSelector,
SelectionKey.OP_READ
);
JDK 允许在任意时间调用 register()。
但如果 Selector 当前正在执行一次选择操作,新注册或修改的 interest set 不会影响这一次正在进行的选择,只会被后续选择操作看到。
这会带来两个问题:
当前 select 什么时候返回?
新 Channel 什么时候开始被监听?
如果没有其他网络事件,Worker 可能继续阻塞。
此外,直接从外部线程修改 Worker 的 Selector 状态,也会破坏“所有状态由所有者线程串行管理”的设计。
2. 使用 sleep
下面这种代码看似能运行:
selector.wakeup();
Thread.sleep(100);
client.register(
selector,
SelectionKey.OP_READ
);
但它并不能保证正确。
线程调度可能是:
Boss 调用 wakeup
↓
Worker 从 select 返回
↓
Worker 检查任务,发现为空
↓
Worker 再次进入 select
↓
Boss 才完成 register
下一次运行时序又可能不同。
sleep() 只能让某次实验“更容易成功”,不能建立线程之间的先后关系。
3. 正确方向
Boss 不应该直接操作 Worker 的内部状态。
它应该向 Worker 提交一个任务:
“请在你自己的线程中,
把这个 Channel 注册到你的 Selector。”
于是需要:
线程安全任务队列
+
Selector.wakeup()
六、wakeup 是什么?
Selector.wakeup() 的作用是:
让当前正在阻塞的选择操作立即返回。
如果另一个线程正在执行:
selector.select();
调用:
selector.wakeup();
会使这次选择操作尽快返回。
如果当前没有线程执行选择操作,那么唤醒效果会保留,使下一次 select() 立即返回;之后的选择操作仍会恢复正常阻塞。
可以把它理解为一个一次性的唤醒标记:
当前正在 select
↓
立即返回
当前没有 select
↓
下一次 select 立即返回
1. wakeup 不会关闭 Selector
调用:
selector.wakeup();
不会:
关闭 Selector
取消 SelectionKey
关闭 Channel
中断 EventLoop 线程
清空事件集合
它只是让一次选择操作结束等待。
2. 返回值可能是零
Selector 被 wakeup() 唤醒后:
int count = selector.select();
count 可能为零,也可能非零。
因此不能这样写:
if (selector.select() == 0) {
continue;
}
否则可能跳过本轮需要执行的任务队列。
更合理的 EventLoop 是:
selector.select();
runTasks();
processSelectedKeys();
无论返回多少,都检查跨线程任务。
七、先入队还是先唤醒?
正确的交接顺序通常是:
1. 把任务放入队列
2. 调用 selector.wakeup()
对应代码:
public void execute(Runnable task) {
taskQueue.offer(task);
selector.wakeup();
}
为什么不能反过来?
1. 先唤醒的风险
假设先执行:
selector.wakeup();
Worker 可能立刻返回,并检查任务队列:
任务队列为空
随后 Worker 再次进入 select()。
Boss 此时才执行:
taskQueue.offer(registerTask);
注册任务便可能在队列里等待,直到下一次网络事件或下一次唤醒。
2. 先入队更安全
先入队:
任务已经可见
↓
再执行 wakeup
Worker 无论是在当前 select() 中,还是即将进入下一次 select(),都会被唤醒,并能取到任务。
Java Selector 的 wakeup() 具有“当前没有选择操作时,让下一次选择立即返回”的语义,因此这里不会因为 Worker 尚未进入 select() 而丢失唤醒。
八、队列如何交接?
完整流程如下:
Boss 线程
↓
accept 得到 SocketChannel
↓
选择 Worker 2
↓
构造注册任务
↓
任务放入 Worker 2 的队列
↓
调用 Worker 2.selector.wakeup()
Worker 2
↓
select() 返回
↓
从任务队列取出注册任务
↓
在自己的线程中执行 register()
↓
下一轮开始监听 OP_READ
原笔记也采用了这种两步法:先把 Channel 放入目标线程的队列,再调用目标 Selector 的 wakeup(),最后由目标线程自己完成注册。
1. 任务不只是注册
队列中可以放入:
注册新 Channel
修改 interestOps
发送数据
关闭连接
定时任务
用户提交的 EventLoop 任务
统一抽象成:
Runnable
以后所有跨线程操作都可以变成:
eventLoop.execute(() -> {
// 由 EventLoop 所属线程执行
});
2. Netty 也有任务队列
Netty 的单线程执行器内部持有任务队列,并由 EventLoop 线程取出任务执行;NioEventLoop 同时负责 Selector I/O 和非 I/O 任务。
Netty 还提供 I/O 时间比例配置,用于平衡处理 I/O 事件和执行普通任务的时间;当前 NioEventLoop 默认比例为 50。
这说明 EventLoop 并不只是:
select + read
而是:
I/O 事件循环
+
任务执行器
+
定时任务调度器
九、Worker 怎么选?
最简单的是轮询:
private final AtomicInteger index =
new AtomicInteger();
private SelectorThread next() {
int position = Math.floorMod(
index.getAndIncrement(),
children.length
);
return children[position];
}
分配结果:
连接 1 → Worker 0
连接 2 → Worker 1
连接 3 → Worker 2
连接 4 → Worker 0
连接 5 → Worker 1
1. 为什么用 floorMod?
不要简单写:
Math.abs(index.getAndIncrement()) % length
因为:
Math.abs(Integer.MIN_VALUE)
结果仍然是负数。
Math.floorMod() 可以正确处理计数器溢出后的负数。
2. 轮询的局限
轮询只能大致均衡:
连接数量
它不知道每条连接的负载。
例如:
Worker 0:
1000 条空闲连接
Worker 1:
1000 条高频连接
两者连接数相同,但实际压力完全不同。
可以进一步考虑:
最少连接数
待执行任务数量
近期 I/O 活跃度
两次随机选择后取较轻者
不过,复杂调度本身也有成本。
原笔记中的轮询方案同样指出:连接数量均匀不代表真实负载均匀,可以根据各 Selector 的连接数量或其他指标改进。
十、事件循环顺序
一个简化的 SelectorThread 可以这样运行:
while (running) {
selector.select();
runTasks();
processSelectedKeys();
}
也可以:
while (running) {
selector.select();
processSelectedKeys();
runTasks();
}
两种顺序都可以工作,但需要考虑公平性。
1. 任务优先
先执行任务
再处理 I/O
优点:
新 Channel 能尽快完成注册
跨线程命令延迟较低
缺点:
任务队列持续堆积时
I/O 可能得不到及时处理
2. I/O 优先
先处理 I/O
再执行任务
优点:
网络事件响应及时
缺点:
持续发生 I/O 时
普通任务可能长期等待
成熟 EventLoop 通常会限制单轮处理量,或按时间比例分配 I/O 和任务执行预算。
3. 不要一次清空一切
假设任务队列中突然进入十万个任务:
一次循环全部执行
↓
Selector 长时间无法处理网络事件
更合理的策略可能是:
每轮最多执行 N 个任务
或
每轮最多执行 T 纳秒
同样,如果某个连接一直有大量数据,也不能永远只处理这一条连接。
这本质上是在解决:
公平性
饥饿
延迟
吞吐量
之间的平衡。
十一、代码骨架
下面实现一个精简版主从 Reactor。
它只用于展示:
Boss 接收连接
Worker 选择
任务队列
wakeup
线程内注册
没有实现完整协议、写缓冲、背压和异常恢复。
import java.io.IOException;
import java.net.InetSocketAddress;
import java.nio.ByteBuffer;
import java.nio.channels.Channel;
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;
import java.util.Objects;
import java.util.Queue;
import java.util.concurrent.ConcurrentLinkedQueue;
import java.util.concurrent.atomic.AtomicInteger;
public final class ReactorServer {
public static void main(String[] args) throws Exception {
SelectorThreadGroup bossGroup =
new SelectorThreadGroup(1, "boss");
SelectorThreadGroup workerGroup =
new SelectorThreadGroup(4, "worker");
bossGroup.setWorkerGroup(workerGroup);
bossGroup.bind(9090);
System.out.println("Server started on port 9090");
}
private static final class SelectorThreadGroup {
private final SelectorThread[] children;
private final AtomicInteger next =
new AtomicInteger();
private volatile SelectorThreadGroup workerGroup;
private SelectorThreadGroup(
int threadCount,
String namePrefix
) throws IOException {
if (threadCount <= 0) {
throw new IllegalArgumentException(
"threadCount must be greater than zero"
);
}
this.children =
new SelectorThread[threadCount];
this.workerGroup = this;
for (int i = 0; i < threadCount; i++) {
SelectorThread child =
new SelectorThread(
namePrefix + '-' + i
);
children[i] = child;
child.start();
}
}
private void setWorkerGroup(
SelectorThreadGroup workerGroup
) {
this.workerGroup =
Objects.requireNonNull(workerGroup);
}
private void bind(int port) throws IOException {
ServerSocketChannel server =
ServerSocketChannel.open();
boolean success = false;
try {
server.configureBlocking(false);
server.bind(
new InetSocketAddress(port)
);
next().registerServer(
server,
workerGroup
);
success = true;
} finally {
if (!success) {
closeQuietly(server);
}
}
}
private void register(SocketChannel client) {
next().registerClient(client);
}
private SelectorThread next() {
int position = Math.floorMod(
next.getAndIncrement(),
children.length
);
return children[position];
}
}
private static final class SelectorThread
implements Runnable {
private final Selector selector;
private final Queue<Runnable> taskQueue =
new ConcurrentLinkedQueue<>();
private final Thread thread;
private volatile boolean running = true;
private SelectorThread(String name)
throws IOException {
this.selector = Selector.open();
this.thread = new Thread(this, name);
}
private void start() {
thread.start();
}
private void execute(Runnable task) {
taskQueue.offer(
Objects.requireNonNull(task)
);
selector.wakeup();
}
private void registerServer(
ServerSocketChannel server,
SelectorThreadGroup workerGroup
) {
execute(() -> {
try {
server.register(
selector,
SelectionKey.OP_ACCEPT,
workerGroup
);
} catch (IOException exception) {
closeQuietly(server);
exception.printStackTrace();
}
});
}
private void registerClient(
SocketChannel client
) {
execute(() -> {
try {
ByteBuffer readBuffer =
ByteBuffer.allocate(8192);
client.register(
selector,
SelectionKey.OP_READ,
readBuffer
);
} catch (IOException exception) {
closeQuietly(client);
exception.printStackTrace();
}
});
}
@Override
public void run() {
while (running) {
try {
selector.select();
runTasks();
processSelectedKeys();
} catch (IOException exception) {
exception.printStackTrace();
}
}
closeQuietly(selector);
}
private void runTasks() {
Runnable task;
while ((task = taskQueue.poll()) != null) {
try {
task.run();
} catch (RuntimeException exception) {
exception.printStackTrace();
}
}
}
private void processSelectedKeys() {
Iterator<SelectionKey> iterator =
selector.selectedKeys().iterator();
while (iterator.hasNext()) {
SelectionKey key = iterator.next();
iterator.remove();
if (!key.isValid()) {
continue;
}
try {
if (key.isAcceptable()) {
handleAccept(key);
}
if (key.isReadable()) {
handleRead(key);
}
} catch (IOException exception) {
closeKey(key);
}
}
}
private void handleAccept(SelectionKey key)
throws IOException {
ServerSocketChannel server =
(ServerSocketChannel) key.channel();
SelectorThreadGroup workerGroup =
(SelectorThreadGroup) key.attachment();
SocketChannel client;
while ((client = server.accept()) != null) {
client.configureBlocking(false);
workerGroup.register(client);
}
}
private void handleRead(SelectionKey key)
throws IOException {
SocketChannel client =
(SocketChannel) key.channel();
ByteBuffer buffer =
(ByteBuffer) key.attachment();
while (true) {
int count = client.read(buffer);
if (count > 0) {
buffer.flip();
byte[] bytes =
new byte[buffer.remaining()];
buffer.get(bytes);
buffer.clear();
String message =
new String(
bytes,
StandardCharsets.UTF_8
);
System.out.printf(
"[%s] %s%n",
Thread.currentThread().getName(),
message
);
continue;
}
if (count == 0) {
return;
}
closeKey(key);
return;
}
}
private void closeKey(SelectionKey key) {
key.cancel();
closeQuietly(key.channel());
}
}
private static void closeQuietly(Channel channel) {
if (channel == null) {
return;
}
try {
channel.close();
} catch (IOException ignored) {
}
}
}
代码链路
服务启动:
创建 Boss Group
↓
创建 Worker Group
↓
Boss 持有 Worker 引用
↓
ServerSocketChannel 注册到 Boss
客户端连接:
Boss 收到 OP_ACCEPT
↓
accept 得到 SocketChannel
↓
选择一个 Worker
↓
注册任务进入 Worker 队列
↓
Worker Selector 被唤醒
↓
Worker 执行注册任务
↓
SocketChannel 关注 OP_READ
客户端发数据:
Worker Selector 返回
↓
发现 OP_READ
↓
Worker 读取数据
最关键的设计是:
Boss 不直接修改 Worker Selector
Boss 只提交任务
真正的注册操作
由 Worker 自己执行
十二、Netty 如何对应?
上面的教学代码可以映射到 Netty:
| 教学代码 | Netty 概念 |
|---|---|
SelectorThread | NioEventLoop |
SelectorThreadGroup | NioEventLoopGroup |
Selector | Java NIO Selector |
taskQueue | EventLoop 任务队列 |
next() | EventLoop 选择器 |
| Boss Group | Parent/Acceptor Group |
| Worker Group | Child Group |
registerClient() | Channel 注册到 EventLoop |
while + select | EventLoop 事件循环 |
Netty 的 NioEventLoopGroup 是面向 NIO Selector Channel 的多线程 EventLoopGroup,并允许配置线程数量和 SelectorProvider。
ServerBootstrap.group(parentGroup, childGroup) 分别设置负责父级 ServerChannel 和子级客户端 Channel 的 EventLoopGroup。
结构可以概括为:
ServerBootstrap
↓
Boss EventLoopGroup
↓
ServerChannel
↓ accept
Child Channel
↓
Worker EventLoopGroup
↓
某个固定 EventLoop
1. Channel 固定归属
Netty Channel 注册到某个 EventLoop 后,通常由这个 EventLoop 处理它后续的全部 I/O 事件。
这让大量 Handler 可以依赖线程内串行执行,而不必为每个连接状态都加锁。
2. Group 不等于线程池任务池
NioEventLoopGroup 虽然包含多个线程,但不能简单理解成普通业务线程池。
它管理的是:
多个 EventLoop
↓
每个 EventLoop 一个执行线程
↓
每个 EventLoop 管理一组 Channel
同一 Channel 的 I/O 不会在每次事件到来时随机选择一个 Worker 线程执行。
十三、ThreadLocal 要吗?
原素材最后提到了使用 ThreadLocal 给每个线程绑定一条队列。
在当前设计中,通常没有必要。
1. 实例字段已经隔离
每个 SelectorThread 都是独立对象:
private final Queue<Runnable> taskQueue =
new ConcurrentLinkedQueue<>();
结构为:
SelectorThread 0
└── taskQueue 0
SelectorThread 1
└── taskQueue 1
SelectorThread 2
└── taskQueue 2
队列已经通过对象实例完成隔离。
不需要再写:
ThreadLocal<Queue<Runnable>> queues;
2. ThreadLocal 的适用场景
ThreadLocal 更适合:
数据逻辑上属于当前线程
调用链又不方便层层传参
例如:
线程上下文
诊断信息
特定线程的临时编码器
但 EventLoop 的 Selector、任务队列和定时任务队列,本身就是 EventLoop 对象的重要状态。
明确写成实例字段通常更加直观。
3. ThreadLocal 的风险
使用 ThreadLocal 还要注意:
数据生命周期可能与线程一样长
线程池线程不会频繁销毁
忘记 remove 可能保留对象
跨线程后读取到的是另一份数据
因此:
线程独享对象不等于必须使用 ThreadLocal。
优先考虑清晰的对象归属关系。
十四、常见误区
1. 一个 Selector 必须对应一个线程
不是 JDK 的硬性规定。
这是 Reactor 和 Netty 常用的线程归属设计,用于减少并发修改和锁竞争。
2. 一个 Worker 只能管理一个连接
错误。
一个 EventLoop 通常会管理多个 Channel。
3. Boss 负责全部网络 I/O
错误。
主从模型中,Boss 主要负责接受连接,Worker 负责已连接 Channel 的后续流量。
4. accept 返回的 Channel 自动非阻塞
错误。
返回的新 SocketChannel 默认是阻塞模式,必须再次配置。
5. wakeup 会关闭 Selector
错误。
它只使一次阻塞选择操作立即返回。
6. wakeup 只能唤醒当前 select
不完整。
如果当前没有选择操作,它会让下一次选择立即返回。
7. 必须先 wakeup 再注册
错误。
跨线程交接更稳妥的方式是:
先把任务放入队列
再调用 wakeup
8. sleep 可以解决注册竞态
错误。
sleep 没有建立可靠的线程间先后关系。
9. Boss 可以直接操作 Worker 的内部集合
不推荐。
即使底层 API 允许并发调用,也会让状态归属和执行顺序变得难以控制。
10. 轮询一定负载均衡
错误。
轮询只能均衡连接数量,不能保证连接负载相同。
11. selectedKeys 会自动清空
错误。
经典迭代写法中,事件处理完后需要调用:
iterator.remove();
Selector 的 selected-key set 由应用读取和删除,但应用不能直接向其中添加 Key。
12. read 返回负数只是普通异常
不准确。
SocketChannel.read() 返回 -1 表示流已经结束,通常需要取消 Key 并关闭 Channel。
13. read 返回 0 表示连接断开
错误。
非阻塞模式下,0 通常表示当前没有可立即读取的数据。
14. ThreadLocal 是线程隔离的必选方案
错误。
对象实例字段本身就可以表达每个 EventLoop 独享的状态。
15. EventLoop 可以执行任意业务
错误。
一个阻塞任务会延迟该 EventLoop 管理的所有连接。
总结
单线程 Selector 的结构是:
一个线程
↓
一个 Selector
↓
同时负责 accept、read、write
它解决了一连接一平台线程的问题,但无法并行使用多个 CPU 核心。
多线程 Reactor 将连接分散到多个 EventLoop:
多个线程
↓
每个线程一个 Selector
↓
每个 Selector 管理一部分 Channel
主从模型进一步拆分职责:
Boss
负责 accept
Worker
负责客户端 read / write
真正困难的不是创建几个线程,而是:
Boss 接收的 SocketChannel
怎样安全地交给 Worker
核心方案是:
Boss 选择 Worker
↓
把注册任务放入 Worker 队列
↓
调用 Worker.selector.wakeup()
↓
Worker 从 select 返回
↓
Worker 在自己的线程中执行 register()
wakeup() 不是关闭 Selector,也不是处理网络事件。
它只是结束一次选择等待,让 EventLoop 有机会处理:
新注册任务
事件修改任务
写出任务
关闭任务
每个 Channel 一旦归属某个 EventLoop,后续 I/O 和状态修改尽量都由这个 EventLoop 串行执行。
这让系统从:
多个线程共同修改连接状态
变成:
跨线程提交命令
EventLoop 线程顺序执行
Netty 的核心结构也可以沿着这条链路理解:
BossGroup
↓ accept
WorkerGroup
↓ choose EventLoop
NioEventLoop
↓ Selector + TaskQueue
Channel
↓ Pipeline
最后,可以用一句话概括整篇文章:
Boss 只负责接收连接,Worker 负责管理连接;Boss 不直接修改 Worker 的 Selector,而是通过任务队列和 wakeup,把注册操作交给 Worker 自己完成。
评论区