Java 虚拟线程如何与网络 I/O 协作:Carrier、Continuation 与 Poller
虚拟线程执行下面这段代码时,表面上和普通线程没有区别:
int count = socket.getInputStream().read(buffer);
process(buffer, count);
当 Socket 暂时没有数据,当前虚拟线程会暂停。承载它的 Carrier Thread 随后可以运行其他虚拟线程。网络数据到达后,原虚拟线程重新得到调度,read() 返回,代码继续执行 process()。
这里容易产生几个疑问:
- Carrier Thread 只能线性执行代码,当前方法还没有执行完,它怎么去运行另一个任务?
- 虚拟线程暂停后保存在哪里?
- 网络事件到达时,Poller 如何找到对应的虚拟线程?
- 调度队列中保存的是虚拟线程、回调函数,还是完整调用栈?
- 网络数据怎样回到原来的 buffer?
- Continuation.yield()、Thread.yield() 和 LockSupport.park() 分别负责什么?
本文以 JDK 21 的实现为基线,沿着一次 Socket 读取的调用过程分析这些问题。JDK 21 的虚拟线程采用 M:N 调度,大量虚拟线程由 JDK 调度到少量平台线程上运行;默认调度器是一个独立的 ForkJoinPool。(OpenJDK)
Carrier Thread 没有任务切换指令
先去掉虚拟线程,只看一个普通的工作线程:
final class Worker extends Thread {
private final BlockingQueue<Runnable> queue;
Worker(BlockingQueue<Runnable> queue) {
this.queue = queue;
}
@Override
public void run() {
while (!isInterrupted()) {
try {
Runnable task = queue.take();
task.run();
} catch (InterruptedException e) {
interrupt();
}
}
}
}
这条线程始终线性执行。它从队列取出一个 Runnable,调用 run(),等待该方法返回,然后进入下一轮循环。
Runnable taskA = queue.take();
taskA.run();
Runnable taskB = queue.take();
taskB.run();
Carrier Thread 的工作方式也是如此。它不会在 taskA.run() 尚未返回时,突然跳到 taskB.run()。
虚拟线程能够切换,是因为 taskA.run() 可以在业务调用尚未完成时暂时返回。返回之前,JVM 会保存虚拟线程的调用栈。以后再次执行这个任务时,JVM 恢复原调用栈,代码继续运行。
因此,一条 Carrier Thread 上实际发生的是:
virtualThreadAResumeTask.run(); // 保存调用栈后暂时返回
virtualThreadBResumeTask.run(); // Carrier 回到调度循环后执行
Carrier Thread 仍然遵守普通方法调用规则。特殊能力位于 Continuation,它允许一段尚未完成的调用栈暂停和恢复。
调度队列中保存了什么
JDK 21 的 VirtualThread 内部有几个关键字段。省略诊断、线程容器和中断处理后,可以整理成下面的结构:
final class VirtualThread extends Thread {
private final Executor scheduler;
private final Continuation cont;
private final Runnable runContinuation;
private volatile int state;
private volatile boolean parkPermit;
private volatile Thread carrierThread;
}
其中:
- scheduler 指向虚拟线程调度器。
- cont 保存可暂停和恢复的执行状态。
- runContinuation 是提交给调度器的任务。
- state 记录虚拟线程当前处于运行、挂起或可调度状态。
- parkPermit 保存一次 unpark 许可。
- carrierThread 指向当前承载该虚拟线程的平台线程。
构造虚拟线程时,会创建 Continuation,并生成一个绑定到当前虚拟线程的方法引用:
this.cont = new VThreadContinuation(this, task);
this.runContinuation = this::runContinuation;
方法引用可以展开成一个普通的 Runnable:
final class ResumeTask implements Runnable {
private final VirtualThread virtualThread;
ResumeTask(VirtualThread virtualThread) {
this.virtualThread = virtualThread;
}
@Override
public void run() {
virtualThread.runContinuation();
}
}
调度队列中保存的核心内容是这样一个恢复任务。它持有 VirtualThread 引用,而 VirtualThread 再持有 Continuation。
ForkJoinPool.WorkQueue
└── ForkJoinTask
└── runContinuation
└── VirtualThread
└── Continuation
VirtualThread 不直接操作 ForkJoinPool 中的某个工作队列。它只调用:
scheduler.execute(runContinuation);
默认调度器收到普通 Runnable 后,会将其包装为 ForkJoinTask,再放入内部队列。OpenJDK 的注释说明,在调度器 Worker 上提交时,任务进入本地队列;其他线程提交时,任务进入外部提交队列。(GitHub)
虚拟线程第一次运行
启动虚拟线程时,内部会把状态由 NEW 修改为 STARTED,随后提交 runContinuation:
void start() {
if (!compareAndSetState(NEW, STARTED)) {
throw new IllegalThreadStateException();
}
submitRunContinuation();
}
private void submitRunContinuation() {
scheduler.execute(runContinuation);
}
某条 Carrier Thread 从 ForkJoinPool 队列中取出任务,最终调用:
virtualThread.runContinuation();
runContinuation() 的主要结构如下:
private void runContinuation() {
int initialState = state();
if (initialState != STARTED
&& initialState != UNPARKED
&& initialState != YIELDED) {
return;
}
if (!compareAndSetState(initialState, RUNNING)) {
return;
}
mount();
try {
cont.run();
} finally {
unmount();
if (cont.isDone()) {
afterDone();
} else {
afterYield();
}
}
}
mount() 将当前平台线程记录为 Carrier,并修改 JVM 中的当前线程身份:
private void mount() {
Thread carrier = Thread.currentCarrierThread();
setCarrierThread(carrier);
carrier.setCurrentThread(this);
}
因此,虚拟线程中的业务代码调用:
Thread.currentThread()
得到的是 VirtualThread 对象,而不是底层 Carrier Thread。
随后执行:
cont.run();
第一次调用时,Continuation 从用户任务入口开始执行。OpenJDK 的 runContinuation() 确实是在 mount() 后调用 cont.run(),并在 finally 中完成 unmount() 和后续状态处理。(GitHub)
Continuation.yield() 如何释放 Carrier
假设业务代码的调用关系如下:
void handleRequest(Socket socket) throws IOException {
User user = loadUser();
byte[] buffer = new byte[1024];
int count = socket.getInputStream().read(buffer);
process(user, buffer, count);
}
虚拟线程执行到网络等待位置时,Carrier Thread 上的逻辑调用栈可能是:
ForkJoinPool.runWorker()
VirtualThread.runContinuation()
Continuation.run()
handleRequest()
InputStream.read()
NioSocketImpl.implRead()
Poller.poll()
LockSupport.park()
VirtualThread.park()
Continuation.yield()
执行 Continuation.yield() 后,JVM 会冻结属于当前 Continuation 的栈帧。局部变量、对象引用、调用关系和程序执行位置都会被保留。
例如,handleRequest() 对应的栈帧中可能存在:
socket -> Socket@100
user -> User@200
buffer -> byte[]@300
count -> 尚未赋值
这些数据属于虚拟线程的执行状态,不再依赖当前 Carrier Thread 的平台栈。
冻结成功后,控制流会离开 Continuation,回到调用 cont.run() 的位置。此时属于虚拟线程的深层调用栈已经从 Carrier Thread 上移走,Carrier Thread 上只剩调度器相关栈帧:
ForkJoinPool.runWorker()
VirtualThread.runContinuation()
cont.run() 随后返回,runContinuation() 进入 finally,执行 unmount():
private void unmount() {
Thread carrier = this.carrierThread;
carrier.setCurrentThread(carrier);
setCarrierThread(null);
}
runContinuation() 返回后,Carrier Thread 回到 ForkJoinPool 的工作循环,从队列中获取其他任务。
后续某条 Carrier Thread 再次调用 cont.run() 时,JVM 恢复此前冻结的栈帧,原来的 Continuation.yield() 开始返回,调用链依次继续:
Continuation.yield() 返回
VirtualThread.park() 返回
LockSupport.park() 返回
Poller.poll() 返回
NioSocketImpl.implRead() 继续
这里没有重新调用 handleRequest(),也没有从方法开头重新执行。恢复点位于此前暂停的位置。
park() 为什么不会立即重新入队
Continuation.yield() 只负责暂停执行,它并不决定虚拟线程什么时候继续。调度策略取决于是谁调用了它。
虚拟线程调用 Thread.yield() 时,当前线程依然具备运行条件,只是暂时让出 Carrier。执行栈冻结后,afterYield() 会立即重新提交 runContinuation:
RUNNING
-> YIELDING
-> YIELDED
-> 重新提交
-> RUNNING
LockSupport.park() 表达的是等待条件尚未满足。例如等待网络数据、等待锁释放或者等待队列元素。执行栈冻结后,虚拟线程不能立即重新进入运行队列,否则会不断恢复、检查条件、再次挂起,造成 CPU 空转。
它的状态变化为:
RUNNING
-> PARKING
-> PARKED
只有其他线程调用 unpark() 后,它才重新具备运行条件:
PARKED
-> UNPARKED
-> 重新提交
-> RUNNING
OpenJDK 的 afterYield() 会检查当前状态。PARKING 会转为 PARKED,默认不重新提交;YIELDING 会转为 YIELDED 并立即提交。(GitHub)
对应代码可以简化为:
private void afterYield() {
int s = state();
if (s == PARKING) {
setState(PARKED);
if (parkPermit
&& compareAndSetState(PARKED, UNPARKED)) {
submitRunContinuation();
}
return;
}
if (s == YIELDING) {
setState(YIELDED);
submitRunContinuation();
}
}
Socket 读取如何进入 Poller
下面进入网络 I/O 路径。
int count = socket.getInputStream().read(buffer);
以 JDK 21 在 Linux 上的实现为例,JDK 会先尝试读取 Socket。内核接收缓冲区已有数据时,读取直接完成,不需要挂起虚拟线程。
如果当前没有数据,非阻塞读取会得到 EAGAIN 或等价状态。JDK 随后将文件描述符注册给 Poller,并挂起当前虚拟线程。JEP 444 描述了这一处理:JDK 中的阻塞网络操作无法立即完成时,虚拟线程会卸载;I/O 可以完成后,再把该虚拟线程提交回调度器。(OpenJDK)
简化后的读取逻辑如下:
int read(int fd, byte[] buffer) throws IOException {
while (true) {
int result = nonBlockingRead(fd, buffer);
if (result >= 0) {
return result;
}
Poller.poll(fd);
}
}
Poller.poll(fd) 内部可以抽象为:
void poll(int fd) {
Thread thread = Thread.currentThread();
waiters.put(fd, thread);
registerWithEpoll(fd);
LockSupport.park();
}
由于当前执行的是虚拟线程,Thread.currentThread() 返回对应的 VirtualThread。
此时存在两份关联。
Linux 内核中的 epoll 记录:
监听 fd 37 的可读事件
JDK 的 Poller 记录:
fd 37 -> VirtualThread@500
操作系统只认识文件描述符,不认识 Java 虚拟线程。Poller 中的映射负责将内核事件重新关联到 Java 线程对象。JDK 21 的 Poller 就是网络事件通知和虚拟线程唤醒之间的中间层。(GitHub)
注册完成后,当前虚拟线程执行:
LockSupport.park();
park() 进入 VirtualThread.park(),设置 PARKING 状态并调用 Continuation.yield()。执行栈冻结后,虚拟线程变成 PARKED,Carrier Thread 被释放。
等待期间的对象关系如下:
Linux epoll
└── fd 37
Poller
└── fd 37 -> VirtualThread@500
VirtualThread@500
├── state = PARKED
├── carrierThread = null
└── Continuation
└── 保存 read() 及其上层调用栈
ForkJoinPool
└── 当前没有该虚拟线程的恢复任务
虚拟线程此时没有占用 Carrier Thread,也不在运行队列中。Poller 保存它的等待关系,Continuation 保存它的执行状态。
网络就绪后如何重新调度
网络数据到达后,Linux 协议栈将字节放入 Socket 的内核接收缓冲区,并将 fd 标记为可读。
Linux 上的 Poller 线程通常阻塞在 epoll_wait()。事件返回后,它获得就绪的文件描述符:
int[] readyFds = epollWait();
for (int fd : readyFds) {
Thread thread = waiters.remove(fd);
if (thread != null) {
LockSupport.unpark(thread);
}
}
Poller 不执行用户的 handleRequest(),也不会在 Poller Thread 中调用原来的 InputStream.read()。它只找到等待该 fd 的虚拟线程,然后调用 unpark()。
VirtualThread.unpark() 的关键逻辑如下:
void unpark() {
Thread currentThread = Thread.currentThread();
if (!getAndSetParkPermit(true)
&& currentThread != this) {
int s = state();
if (s == PARKED
&& compareAndSetState(PARKED, UNPARKED)) {
submitRunContinuation();
}
}
}
submitRunContinuation() 最终执行:
scheduler.execute(runContinuation);
这一步把虚拟线程重新交给调度器。OpenJDK 的 unpark() 会先设置 parkPermit,再通过状态 CAS 将已挂起的虚拟线程转为 UNPARKED,随后提交其恢复任务。(GitHub)
由于调用 unpark() 的通常是 Poller Thread,它不属于默认调度器的 Worker,因此这次任务一般走 ForkJoinPool 的外部提交路径。
某条 Carrier Thread 之后取出该任务,再次调用:
virtualThread.runContinuation();
状态由 UNPARKED 修改为 RUNNING,虚拟线程挂载到新的 Carrier Thread,然后执行:
cont.run();
此前冻结的调用栈被恢复,LockSupport.park() 返回,Socket 读取逻辑继续。
网络数据没有进入调度队列
调度队列中没有网络数据,也没有下面这种对象:
new ResumeTask(virtualThread, networkData);
Poller 的职责是通知“fd 现在可能可读”,调度器负责安排虚拟线程继续执行。网络字节仍然保存在 Socket 的内核接收缓冲区中。
虚拟线程恢复后,读取循环再次调用操作系统:
int read(int fd, byte[] buffer) throws IOException {
while (true) {
int result = nonBlockingRead(fd, buffer);
if (result >= 0) {
return result;
}
Poller.poll(fd);
// park 返回后,再次进入循环
}
}
第一次 nonBlockingRead() 返回 EAGAIN,虚拟线程进入等待。
网络就绪后,Poller 调用 unpark(),虚拟线程恢复,随后再次执行:
nonBlockingRead(fd, buffer);
这一次内核接收缓冲区已有数据,系统调用将字节复制到原来的 buffer 中,并返回字节数。
之所以仍然能够访问原来的 buffer,是因为 Continuation 保存的调用栈中仍持有:
buffer -> byte[]@300
数据传递路径为:
网卡
-> Linux 网络协议栈
-> Socket 内核接收缓冲区
-> 虚拟线程被唤醒
-> 恢复 read() 调用栈
-> 再次执行 read()
-> 数据复制到原 buffer
调度队列只传递执行资格。Poller 只传递就绪通知。网络数据由内核 Socket 缓冲区保存。
parkPermit 解决丢失唤醒
虚拟线程准备挂起时,网络事件可能提前到达。
考虑下面的时序:
虚拟线程发现当前没有数据
Poller 注册 fd
网络数据立即到达
Poller 调用 unpark()
虚拟线程随后才执行 park()
如果 unpark() 只能唤醒已经处于 PARKED 状态的线程,这次通知就会丢失。虚拟线程之后进入 park(),可能一直无法恢复。
parkPermit 相当于一个容量为 1 的许可:
private volatile boolean parkPermit;
unpark() 首先设置:
parkPermit = true;
park() 开始时先尝试消费许可:
if (getAndSetParkPermit(false)) {
return;
}
如果网络事件已经到达,park() 直接返回,不再冻结调用栈。
连续调用多次 unpark() 也只会保存一个许可:
unpark(thread);
unpark(thread);
unpark(thread);
最终仍然只是:
parkPermit = true
下一次 park() 消费该许可后恢复为 false。
这也是 JUC 同步器通常使用循环检查条件的原因:
while (!conditionSatisfied()) {
LockSupport.park();
}
park() 返回只说明线程重新获得执行机会,业务条件仍需重新判断。
PARKING 解决重复调度
RUNNING 和 PARKED 两个状态还不够,因为调用栈冻结需要一段执行过程。
虚拟线程会先执行:
state = PARKING;
Continuation.yield();
在设置 PARKING 后、Continuation 完成冻结前,Poller 可能已经调用 unpark()。
此时原 Carrier Thread 仍可能执行当前虚拟线程,不能立刻把恢复任务提交给调度器。否则另一条 Carrier Thread 可能同时取出该虚拟线程。
因此,在 PARKING 状态下,unpark() 只设置 parkPermit,暂不提交任务。
等 Continuation 冻结完成,原 Carrier 执行 afterYield():
state = PARKED;
if (parkPermit
&& compareAndSetState(PARKED, UNPARKED)) {
submitRunContinuation();
}
这样既保留了提前到达的唤醒通知,也避免两条 Carrier Thread 同时运行同一个虚拟线程。
另外,runContinuation() 在恢复前还会对状态执行 CAS。即使调度队列中意外出现两个指向同一虚拟线程的恢复任务,也只有一条 Carrier Thread 能成功将状态改为 RUNNING。
三条线程如何协作
完整实现虽然涉及 ForkJoinPool、Continuation、Poller 和操作系统,但运行过程仍然可以拆成三条普通线程循环。
Carrier Thread 的循环:
while (true) {
Runnable task = scheduler.takeTask();
task.run();
}
Poller Thread 的循环:
while (true) {
int[] readyFds = epollWait();
for (int fd : readyFds) {
VirtualThread thread = waiters.remove(fd);
LockSupport.unpark(thread);
}
}
虚拟线程中的 Socket 读取循环:
while (true) {
int count = tryRead(fd, buffer);
if (count >= 0) {
return count;
}
registerPoller(fd, Thread.currentThread());
LockSupport.park();
}
三条执行流都保持线性。它们通过以下对象建立联系:
ForkJoinPool WorkQueue
保存可运行的 runContinuation
VirtualThread
持有 scheduler、Continuation 和线程状态
Poller
保存 fd 到 VirtualThread 的等待关系
Socket 内核接收缓冲区
保存网络数据
Carrier Thread 调用恢复任务,Poller 根据 fd 找到等待线程,虚拟线程通过 scheduler 重新提交自己,Continuation 保存暂停位置和局部对象引用。
完整调用路径
将整个过程连接起来,可以得到下面的关键路径:
Thread.startVirtualThread
-> VirtualThread.start
-> submitRunContinuation
-> ForkJoinPool 工作队列
-> Carrier 取出任务
-> VirtualThread.runContinuation
-> mount
-> Continuation.run
-> 用户代码
-> Socket.read
-> 非阻塞 read 返回 EAGAIN
-> Poller 注册 fd
-> LockSupport.park
-> VirtualThread.park
-> Continuation.yield
-> 调用栈被冻结
-> cont.run 暂时返回
-> unmount
-> VirtualThread 进入 PARKED
-> Carrier 返回调度循环
网络数据到达
-> Socket 内核接收缓冲区
-> epoll_wait 返回 fd
-> Poller 找到 VirtualThread
-> LockSupport.unpark
-> VirtualThread.unpark
-> submitRunContinuation
-> ForkJoinPool 工作队列
-> Carrier 取出恢复任务
-> VirtualThread.runContinuation
-> mount
-> Continuation.run
-> 恢复被冻结的调用栈
-> LockSupport.park 返回
-> 再次执行 read
-> 数据复制到原 buffer
-> 用户代码继续
理解这条路径后,虚拟线程的调度就不再神秘。Carrier Thread 始终执行普通任务循环;Continuation 允许当前任务在保留调用栈的情况下暂时返回;Poller 将网络 fd 与等待中的虚拟线程关联起来;unpark() 在事件到达后重新提交恢复任务。
JDK 21 与 JDK 24 的 Pinning 差异
JDK 21 中,如果虚拟线程在持有 synchronized Monitor 时进入某些阻塞操作,Continuation 可能无法卸载,虚拟线程和 Carrier Thread 会一起阻塞,这种情况称为 Pinning。
因此,JDK 21 阶段经常建议避免在长时间 synchronized 临界区中执行网络 I/O,必要时使用 ReentrantLock。
JDK 24 交付了 JEP 491,修改了 JVM Monitor 与虚拟线程的协作方式。虚拟线程在 synchronized 方法或代码块中阻塞时,绝大部分情况下也能卸载并释放 Carrier Thread。JDK 24 之后,选择 synchronized 还是 java.util.concurrent.locks,应更多依据可中断获取、公平性、超时和条件变量等语义,而不是单纯为了规避 Monitor Pinning。(OpenJDK)
持锁期间执行长时间网络 I/O 仍然需要谨慎。即便 Carrier 能够释放,锁本身仍然处于持有状态,其他需要进入同一临界区的线程仍会等待。
参考源码
- JEP 444:Virtual Threads。(OpenJDK)
- JDK 21 java.lang.VirtualThread。(GitHub)
- JDK 21 jdk.internal.vm.Continuation。(GitHub)
- JDK 21 sun.nio.ch.Poller。(GitHub)
- JEP 491:Synchronize Virtual Threads without Pinning。(OpenJDK)