Virtual Thread重投递至ForkJoinPool任务队列过程解析 Virtual Thread重投递至ForkJoinPool任务队列过程解析前言Virtual Thread重投递至ForkJoinPool任务队列过程Virtual Thread Wake-Up Architecture一、 Linux 内核态从网卡中断到 epoll_wait 唤醒1. 套接字数据就绪与回调触发 (net/core/sock.c)2. epoll 唤醒回调逻辑 (fs/eventpoll.c)二、 JVM 用户态 Poller 线程事件接收 (C / JNI / Java)1. Native 层 epoll_wait 绑定 (EPoll.c)2. Java Poller 线程循环与 Request 映射 (Poller.java)三、 VirtualThread 状态机变更与 Continuation 恢复提交1. 虚拟线程原子状态机 (VirtualThread.java)四、 ForkJoinPool 外部任务入队与 Carrier 线程唤醒1. 外部入队与环形缓冲区操作 (ForkJoinPool.java)2. 唤醒空闲 Carrier 线程 (signalWork)五、 Carrier 线程挂载 (Mount) 与 Native 栈帧解冻 (Thaw)1. 唤醒后重新执行非阻塞 Syscall (SocketChannelImpl.java)六、 完整唤醒链路时序图汇总前言本文旨在记录近期研读Java源码的学习心得与疑难问题。由于个人理解水平有限文中内容难免存在疏漏恳请读者不吝指正。Virtual Thread重投递至ForkJoinPool任务队列过程Project Loom 中虚拟线程Virtual Thread将“阻塞式编程”转换为底层“异步非阻塞 I/O”。当虚拟线程因等待 Socket 读写而挂起时底层 Carrier 线程JavaThread会被释放。当硬件网卡接收到数据后系统将沿着Linux Kernel 内核中断→ \to→JVM Native Poller 线程→ \to→Java 态VirtualThread.unpark()→ \to→ForkJoinPool任务队列的全链路唤醒虚拟线程。Virtual Thread Wake-Up Architecture Loom Virtual Thread Wake-Up Architecture [ Linux Kernel Space ] Hardware IRQ - softirq (NET_RX_SOFTIRQ) - TCP/IP Stack - sk_receive_queue - sock_def_readable() - ep_poll_callback() - Wakeup Poller Thread in ep_poll() │ ▼ (Syscall return) [ OpenJDK JVM Poller (C / JNI) ] Java_sun_nio_ch_EPoll_wait() - epoll_wait() unblocks │ ▼ [ OpenJDK Poller (Java Space) ] Poller$PollerThread.run() - Poller.polled() - Lookup Poller.Request - LockSupport.unpark(vthread) │ ▼ [ Virtual Thread State Machine ] VirtualThread.unpark() - CAS State: PARKED - RUNNABLE - VirtualThread.submitRunContinuation() │ ▼ [ ForkJoinPool Scheduler ] ForkJoinPool.externalPush(runContinuation) - Submission WorkQueue [ring buffer] - signalWork() - futex_wake() / Unsafe.unpark(carrierThread) │ ▼ (Carrier Thread Woken Up) [ Carrier Thread Continuation Thaw ] ForkJoinWorkerThread.run() - Pop runContinuation - Continuation.run() - Continuation::thaw() [C Assembly Stub] - Restore stack frames from stackChunkOop - Return to Java SocketChannelImpl.read() - Re-issue sys_read() [SUCCESS] 一、 Linux 内核态从网卡中断到epoll_wait唤醒数据包到达网卡后Linux 内核通过中断与 socket 唤醒链表解锁阻塞在epoll_wait的 JVM Poller 线程。------------------- ---------------------- ----------------------- | NIC HW Interrupt | --- | NET_RX_SOFTIRQ | --- | tcp_v4_rcv() | | (Data Packet) | | (NAPI Poll Processing| | (TCP State Machine) | ------------------- ---------------------- ----------------------- | v ------------------- ---------------------- ----------------------- | ep_poll_callback | --- | wake_up_interruptible| --- | sock_def_readable() | | (Add to rdllist) | | _sync_poll() | | (sk_data_ready) | ------------------- ---------------------- -----------------------1. 套接字数据就绪与回调触发 (net/core/sock.c)当 TCP 协议栈完成数据包重组并将sk_buff放入 socket 接收队列sk_receive_queue时内核调用 socket 的数据就绪回调函数sk_data_ready默认指向sock_def_readable// Linux Kernel: net/core/sock.cvoidsock_def_readable(structsock*sk){structsocket_alloc*si;rcu_read_lock();// 获取挂载在 socket 上的等待队列头 (wait_queue_head_t)structsocket_wq*wqrcu_dereference(sk-sk_wq);if(skwq_has_sleeper(wq)){// 触发等待队列上的回调函数唤醒关联的 epoll 项wake_up_interruptible_sync_poll(wq-wait,EPOLLIN|EPOLLPRI|EPOLLRDNORM|EPOLLRDBAND);}sk_wake_async(sk,SOCK_WAKE_WAITD,POLL_IN);rcu_read_unlock();}2.epoll唤醒回调逻辑 (fs/eventpoll.c)当epoll_ctl(EPOLL_CTL_ADD)将 socket 注册到epoll句柄时内核为该 socket 挂载了一个epitem并将ep_poll_callback注册为唤醒回调。// Linux Kernel: fs/eventpoll.cstaticintep_poll_callback(wait_queue_entry_t*wait,unsignedmode,intsync,void*key){intpwake0;structepitem*epiep_item_from_wait(wait);// 获取包含此 socket 的 epitemstructeventpoll*epepi-ep;uintptr_tpollflags(uintptr_t)key;spin_lock_irqsave(ep-lock,flags);// 1. 检查事件掩码是否匹配 (例如是否触发了 EPOLLIN)if(pollflags!(pollflagsepi-event.events))gotoout_unlock;// 2. 将就绪的 epitem 追加到 epoll 的就绪链表 rdllist 中if(!ep_is_linked(epi)){list_add_tail(epi-rdllnk,ep-rdllist);}// 3. 如果 epoll_wait 线程正处于 SLEEPING 状态 (等待在 ep-wq 队列上)if(waitqueue_active(ep-wq)){// 唤醒阻塞在 sys_epoll_wait 的 JVM Poller 线程wake_up_locked(ep-wq);}out_unlock:spin_unlock_irqrestore(ep-lock,flags);return1;}内核中的wake_up_locked(ep-wq)将处于TASK_INTERRUPTIBLE状态的 JVM Poller 线程的task_struct重新变更为TASK_RUNNING并将其放入 CFS / EEVDF 运行队列。调度器选中该线程后系统调用sys_epoll_wait返回就绪的事件数量。二、 JVM 用户态 Poller 线程事件接收 (C / JNI / Java)JVM Poller 是一个独立的 C / Java 平台线程非 Virtual Thread专门负责轮询epoll_fd。[ Linux Kernel ] --- epoll_wait() 返回 | v [ EPoll.c (JNI) ] --- Java_sun_nio_ch_EPoll_wait() 填入 address 内存数组 | v [ Poller.java ] --- Poller$PollerThread.run() | v Poller.polled(fd, events) | v LockSupport.unpark(vthread)1. Native 层epoll_wait绑定 (EPoll.c)Poller 线程通过 JNI 循环调用本地epoll_wait系统调用// OpenJDK: src/java.base/unix/native/libnio/ch/EPoll.cJNIEXPORT jint JNICALLJava_sun_nio_ch_EPoll_wait(JNIEnv*env,jclass clazz,jint epfd,jlong address,jint numfds,jint timeout){// 将传入的 Java DirectBuffer 地址强转为 struct epoll_event 数组structepoll_event*events(structepoll_event*)jlong_to_ptr(address);intres;RESTARTABLE(epoll_wait(epfd,events,numfds,timeout),res);if(res0){JNU_ThrowIOExceptionWithLastError(env,epoll_wait failed);return-1;}returnres;// 返回捕获的就绪文件描述符数量}2. Java Poller 线程循环与 Request 映射 (Poller.java)sun.nio.ch.Poller位于java.base/sun/nio/ch/Poller.java提取底层返回的就绪事件数组通过fd索引找出对应的等待请求Poller.Request// OpenJDK: src/java.base/share/classes/sun/nio/ch/Poller.javapackagesun.nio.ch;publicabstractclassPoller{// Poller 专属后台线程privateclassPollerThreadextendsThread{publicvoidrun(){for(;;){// 1. 调用 JNI 的 epoll_wait() 阻塞等待事件结果存入 address 缓冲区intnumeventsEPoll.wait(epfd,address,MAX_EVENTS,timeout);for(inti0;inumevents;i){longeventAddrEPoll.getEvent(address,i);intfdEPoll.getDescriptor(eventAddr);intevEPoll.getEvents(eventAddr);// 2. 根据 fd 响应就绪事件解绑事件并处理polled(fd,ev);}}}}// 响应事件的核心处理逻辑protectedfinalvoidpolled(intfd,intev){// 根据 fd 从哈希表/数组中找出之前 parked 的 Poller.RequestRequestrequestmap.remove(fd);if(request!null){// 获取绑定在该请求上的 VirtualThread 引用Threadthreadrequest.thread();// 3. 唤醒被挂起的 VirtualThreadLockSupport.unpark(thread);}}}三、VirtualThread状态机变更与 Continuation 恢复提交LockSupport.unpark(thread)将引发VirtualThread内部原子状态机的转换并将该虚拟线程的runContinuation提交至 ForkJoinPool。[ LockSupport.unpark() ] | v [ VirtualThread.unpark() ] | --- CAS(State): PARKED - RUNNABLE | v [ submitRunContinuation() ] | v [ scheduler.execute(runContinuation) ] (提交至 ForkJoinPool)1. 虚拟线程原子状态机 (VirtualThread.java)虚拟线程在挂起与唤醒过程中使用volatile int state结合UnsafeCAS 保证线程安全// OpenJDK: src/java.base/share/classes/java/lang/VirtualThread.javapackagejava.lang;publicclassVirtualThreadextendsThread{// 虚拟线程状态常量privatestaticfinalintNEW0;privatestaticfinalintSTARTED1;privatestaticfinalintRUNNING2;// 正在 Carrier 线程上运行privatestaticfinalintPARKING3;// 正在 Unmount/Freeze 过程中privatestaticfinalintPARKED4;// 已挂起在堆中休眠privatestaticfinalintPINNED5;// 钉住状态privatestaticfinalintTIMED_PARK6;privatevolatileintstate;privatefinalRunnablerunContinuation;// 用于包装 Continuation.run() 的任务privatefinalExecutorscheduler;// 通常为默认的 ForkJoinPoolOverridevoidunpark(){if(!Thread.currentThread().equals(this)){// 外部线程如 Poller 线程调用唤醒intsstate;// 1. 尝试通过 CAS 将状态从 PARKED 变更为 RUNNABLEwhile(sPARKED||sTIMED_PARK){if(STATE.compareAndSet(this,s,RUNNABLE)){// 2. CAS 变更成功将任务提交至调度器submitRunContinuation();break;}sstate;// 重新读取最新状态 retry}}}privatevoidsubmitRunContinuation(){try{// 3. 将 runContinuation 重新提交给 ForkJoinPoolscheduler.execute(runContinuation);}catch(RejectedExecutionExceptionree){// 异常兜底逻辑}}}四、ForkJoinPool外部任务入队与 Carrier 线程唤醒由于PollerThread属于非 ForkJoinWorker 的外部平台线程External Thread它必须通过ForkJoinPool.externalPush将runContinuation压入无锁的 Submission Queue提交队列。PollerThread (External non-FJP thread) | v [ ForkJoinPool.externalPush(runContinuation) ] | v [ Find/Create WorkQueue (Submission Queue) ] | v [ Push into WorkQueue.array (Ring Buffer via CAS) ] | v [ ForkJoinPool.signalWork() ] | v [ Unsafe.unpark(Carrier Thread) / futex_wake() ]1. 外部入队与环形缓冲区操作 (ForkJoinPool.java)// OpenJDK: src/java.base/share/classes/java/util/concurrent/ForkJoinPool.javapackagejava.util.concurrent;publicclassForkJoinPoolextendsAbstractExecutorService{// 外部线程提交任务入口finalvoidexternalPush(ForkJoinTask?task){WorkQueueq;// 获取外部提交线程对应的随机 WorkQueue 槽位intprobeThreadLocalRandom.getProbe();if(probe0){ThreadLocalRandom.localInit();probeThreadLocalRandom.getProbe();}// 1. 找到对应的 submission WorkQueue (偶数索引槽位)if((qfindSubmissionQueue(probe))!null){// 2. 将 task 压入该队列的 array 数组中q.push(task,this);return;}// 慢速初始化队列并入队externalSubmit(task);}// WorkQueue 内部的压栈逻辑 (非锁 Ring Buffer)staticfinalclassWorkQueue{volatileinttop;// 栈顶指针 (仅 Push 线程可更新)intbase;// 栈底指针ForkJoinTask?[]array;// 任务环形数组finalvoidpush(ForkJoinTask?task,ForkJoinPoolpool){intstop;ForkJoinTask?[]aarray;if(a!null){intsizea.length;intcapsize-1;// 1. 计算元素在数组中的内存偏移量 Offsetintindexscap;// 2. 将 runContinuation 写入 Ring Buffer 对应槽位QA.setRelease(a,index,task);// 3. 递增 top 指针 (Release 语义保障可见性)TOP.setRelease(this,s1);// 4. 唤醒或创建空闲的 Carrier Worker 线程去 Steal 任务pool.signalWork();}}}}2. 唤醒空闲 Carrier 线程 (signalWork)signalWork()通过解析ForkJoinPool内部的 64 位 atomicctl状态字段寻找处于休眠状态AC_MASK/TC_MASK控制的 Worker 线程并对其执行LockSupport.unpark()底层在 Linux 上对应sys_futex(FUTEX_WAKE)// OpenJDK: src/java.base/share/classes/java/util/concurrent/ForkJoinPool.javafinalvoidsignalWork(){longc;// 检查是否有空闲线程 (ctl 包含平摊活跃线程数与空闲链表 Head 指针)while((cctl)0L){intsp(int)c;// 获取休眠栈顶 Worker 的 IDif(sp0){// 没有处于休眠状态的 Worker检查是否需要创建新的 Carrier 线程if((cADD_WORKER)!0L)tryAddWorker(c);break;}// 尝试 CAS 弹出现有的休眠 WorkerWorkQueue[]qesqueues;intispSMASK;if(qes!nulliqes.lengthqes[i]!null){WorkQueuevqes[i];longnc(v.stackPredSP_MASK)|((cAC_UNIT)~SP_MASK);if(CTL.compareAndSet(this,c,nc)){v.phasesp;// 获取绑定的 Carrier 平台线程对象 (ForkJoinWorkerThread)Threadpv.owner;if(p!null){// 唤醒该 Carrier 线程LockSupport.unpark(p);break;}}}}}五、 Carrier 线程挂载 (Mount) 与 Native 栈帧解冻 (Thaw)Carrier 线程被唤醒后从WorkQueue弹出runContinuation并运行最终进入 C 层的Continuation::thaw将stackChunkOop内的 Java 栈帧重新写回当前 Carrier 线程的物理 Native Stack。[ Carrier Thread (ForkJoinWorkerThread) ] | v Executes runContinuation.run() | v jdk.internal.vm.Continuation.run() | v Continuation::thaw() [C Native Code / Assembly Stub] | --- 1. 从 stackChunkOop 反序列化栈帧 --- 2. 将 JIT / Interpreter 栈帧写回 Carrier Native Stack --- 3. 调整硬件 RSP / RBP 寄存器指针 --- 4. 设置 Thread.currentThread() VirtualThread --- 5. JMP 至挂起时的 _pc (RIP) 地址 | v [ Java Land: SocketChannelImpl.read() ] | v Re-issue sys_read(fd, buf, len) ---【返回读到的真正数据 SUCCESS】1. 唤醒后重新执行非阻塞 Syscall (SocketChannelImpl.java)以网络读操作为例当 Virtual Thread 恢复执行后它会从上次 Unmount 离开的指令地址_pc继续向下运行即跳出park()循环并再次调用read()// OpenJDK: src/java.base/share/classes/sun/nio/ch/SocketChannelImpl.javapublicintread(ByteBufferbuf)throwsIOException{for(;;){// 1. 尝试发起 sys_read 系统调用intnTokens.read(fd,buf);if(nIOStatus.UNAVAILABLE){// 如果是在 Mount 前会在这里将当前线程注册到 Poller 并 yieldPoller.register(fd,Net.POLLIN);Continuation.yield(SCOPE);// -------------------------------------------------------------// 【当重新 Thaw 恢复执行后PC 指向此处进入下一次循环再次发起 read】// -------------------------------------------------------------continue;}// 2. 再次执行 read 时Linux Socket 接收缓冲区已有数据sys_read 返回读取的字节数 nreturnIOStatus.normalize(n);}}六、 完整唤醒链路时序图汇总下图归纳了从内核收到数据包到 Java 逻辑执行恢复的端到端调用序列序号步骤执行主体运行空间核心 API / 函数关键动作1网卡中断Hardware / SoftIRQ内核态tcp_v4_rcv协议栈重组数据存入sk_receive_queue2链表唤醒Socket Subsystem内核态sock_def_readable→ \to→ep_poll_callback将epitem加入rdllist唤醒阻塞在epoll_wait的线程3Syscall 返回JVM Poller 线程用户态Java_sun_nio_ch_EPoll_waitepoll_wait接收到就绪的fd数量4映射 ThreadPoller Java 逻辑用户态Poller.polled()从fd找到对应的Poller.Request和VirtualThread5状态 CASVirtualThread用户态VirtualThread.unpark()CAS 将虚拟线程状态从PARKED转换为RUNNABLE6压入队列Poller Java 逻辑用户态ForkJoinPool.externalPush()将runContinuation压入ForkJoinPool无锁 Submission 队列7唤醒 CarrierForkJoinPool用户态 / 内核态LockSupport.unpark(Carrier)执行sys_futex(FUTEX_WAKE)唤醒休眠的 ForkJoin Worker 线程8Mount/ThawCarrier 线程用户态 (C)Continuation::thaw()从stackChunkOop恢复栈帧写回 Carrier 物理栈重置 RSP/RBP9续放执行VirtualThread用户态 (Java)SocketChannelImpl.read()跳转至挂起时的PC继续运行重新发起sys_read成功拿到数据