前置知识: Java、Java、Java

Java 与虚拟线程

41 min中级

Project Loom 虚拟线程、结构化并发、Continuation 机制与性能调优全景式深度解析

前置知识

学习目标

  • 掌握「0. 本节阅读指引(先读这一节)」的核心机制、典型用法与常见陷阱
  • 掌握「1. 历史动机与演化」的核心机制、典型用法与常见陷阱
  • 掌握「2. 形式化定义」的核心机制、典型用法与常见陷阱
  • 掌握「3. 理论推导:JVM 视角下的虚拟线程机制」的核心机制、典型用法与常见陷阱
  • 掌握「4. 代码示例」的核心机制、典型用法与常见陷阱

0. 本节阅读指引(先读这一节)

本篇是「虚拟线程」,含理论章节。

第一遍只读:4. 代码示例与文末速查小节(创建虚拟线程、虚拟线程执行器、平台线程对比);附录 B 迁移检查清单建议读完。

可跳过:1-3 节(历史、形式化、理论推导)第二遍细读。

前置:047 多线程基础、043 并发编程基础。

Java 虚拟线程深度指南(Project Loom)

虚拟线程(Virtual Threads)是 Java 21(2023 年 9 月)正式发布的里程碑式特性,由 Project Loom 项目孵化而成。它将 Java 的并发模型从”操作系统线程 1:1 映射”重定义为”JVM 调度的 M:N 映射”,让开发者能用传统的同步阻塞式编程风格实现高并发,无需学习响应式编程(Reactive Programming)的复杂范式。这一特性被 Brian Goetz(Java 语言架构师)称为”自 Java 8 Lambda 以来最重要的语言演进”,因为它重新定义了 Java 在云原生时代处理高并发 I/O 的能力边界。本文将系统性剖析虚拟线程的设计哲学、调度机制、Continuation 原理、Pinning 陷阱、结构化并发、与 Spring Boot 3.2+ 的集成实践,以及与响应式编程的深度对比,让读者既能掌握”如何使用虚拟线程”,也能理解”虚拟线程为何如此设计”,最终建立对 Java 并发演进全景的系统认知。


1. 历史动机与演化

1.1 并发模型的演进:从线程到协程到虚拟线程

Java 的并发模型经历了三个主要阶段:

阶段 1:1:1 线程模型(1996-2023,JDK 1.0-20)

Java 1.0 采用了”每个 Thread 对象对应一个操作系统线程”的 1:1 模型。这种设计的优点是语义清晰、与操作系统调度器对齐,缺点是:

  • 创建成本高:每个平台线程需要分配 1MB 栈空间(默认 -Xss1m),操作系统通过 clone 系统调用创建,涉及内核态切换。
  • 数量上限低:典型 Linux 服务器可创建的平台线程数约 4000-8000(受 ulimit -u 与内存限制),无法支撑”一个请求一个线程”的微服务架构。
  • 阻塞成本高:线程阻塞(如 socket.read())时,操作系统线程被挂起,但栈空间与内核资源仍被占用。

为了解决高并发 I/O 问题,社区演化出两条路径:

路径 A:响应式编程(Reactive Programming)

  • 代表框架:RxJava(2013)、Project Reactor(2016)、Akka Streams(2014)。
  • 核心思想:使用少量线程 + 异步回调 + 数据流抽象,实现”非阻塞 I/O + 事件驱动”。
  • 优点:单机可支撑 10 万级并发连接(如 Netty)。
  • 缺点:
    • 编程范式陡峭:需学习 Mono、Flux、flatMap、zip、compose 等组合子。
    • 调试困难:异步调用栈不连续,异常追踪需借助 Hooks.onEachOperator。
    • 与同步库不兼容:JDBC、阻塞式 HTTP 客户端在响应式代码中会”毒化”事件循环。
    • “回调地狱”虽被 flatMap 缓解,但代码可读性仍显著低于同步风格。

路径 B:协程(Coroutine)

  • 代表语言:Kotlin(kotlinx.coroutines,2018)、Go(goroutine,2012)、Erlang(process,1986)。
  • 核心思想:用户态调度的轻量级线程,可在阻塞点自动挂起/恢复,栈空间按需增长。
  • 优点:编程风格接近同步代码,单机可支撑百万级并发。
  • 缺点(对 Java 而言):JVM 上 Kotlin 协程依赖 Continuation 与 suspend 关键字,无法与 Java 互操作;Go 协程是语言级原生,Java 无法复用。

阶段 2:Project Loom 孵化(2018-2023)

Project Loom 由 Ron Pressler(Oracle)于 2018 年发起,目标是在 JVM 层面实现”协程式”的轻量级线程,但保持 Java 现有的同步阻塞式编程风格。其核心创新:

  • 不引入新关键字:虚拟线程使用 Thread 类的子类型,所有现有 Thread API(sleep、interrupt、join、ThreadLocal)都可用。
  • JVM 层调度:虚拟线程由 JVM 内置的 ForkJoinPool 调度,而非操作系统调度。
  • Continuation 机制:虚拟线程的栈帧在阻塞时保存到堆上(Continuation 对象),恢复时重新挂载到载体线程。
  • 与现有库兼容:JDBC、HttpClient、Socket 等 API 在虚拟线程下自动”卸载”,无需修改业务代码。

这一设计的核心哲学是:让开发者继续写同步代码,但获得异步性能。

阶段 3:虚拟线程正式化(2023-至今,JDK 21+)

  • JDK 19(2022-09):虚拟线程作为预览特性发布(JEP 425)。
  • JDK 20(2023-03):虚拟线程第二次预览(JEP 436),根据反馈调整 API。
  • JDK 21(2023-09):虚拟线程正式发布(JEP 444),成为 LTS 特性。
  • JDK 22-24(2024-2025):结构化并发、作用域值持续预览演进;JDK 24 的 JEP 491 使 synchronized 内阻塞不再固定载体线程;社区生态(Spring Boot 3.2、Helidon Níma、Quarkus 3.6)全面接入。
  • JDK 25(2025-09,LTS):作用域值转正(JEP 506);结构化并发第 5 次预览(JEP 505,API 重设计为 open() + Joiner)。
  • JDK 26(2026-03):结构化并发第 6 次预览(JEP 525),仍未转正。

1.2 关键里程碑时间线

时间JDK 版本虚拟线程相关演进
2018Project Loom 启动Ron Pressler 提出 Loom 项目,目标是”轻量级线程 + 同步风格”
2021-05JDK 16Loom 早期原型进入沙箱,java.lang.Fiber 实验性 API
2022-09JDK 19JEP 425:虚拟线程预览(Preview),API 为 Thread.ofVirtual()
2023-03JDK 20JEP 436:虚拟线程第二次预览,Thread.startVirtualThread 简化 API
2023-09JDK 21 (LTS)JEP 444:虚拟线程正式发布,成为生产可用特性
2023-11Spring Boot 3.2spring.threads.virtual.enabled=true 一键启用虚拟线程
2024-03JDK 22JEP 462:结构化并发第二次预览(StructuredTaskScope)
2024-09JDK 23JEP 480:结构化并发第三次预览;JEP 481:作用域值第三次预览
2025-03JDK 24JEP 491:synchronized 阻塞不再固定载体线程(Monitor 类 Pinning 基本消除);JEP 499:结构化并发第四次预览
2025-09JDK 25 (LTS)JEP 506:作用域值转正;JEP 505:结构化并发第五次预览(API 重设计为 open() + Joiner)
2026-03JDK 26JEP 525:结构化并发第六次预览,仍未转正

1.3 设计哲学:为何选择”虚拟线程”而非”协程”

Java 设计团队(Brian Goetz、Ron Pressler)在 Loom 设计阶段曾考虑三种方案:

  1. 方案 A:语言级 async/await 关键字(C#、Rust、Python 风格)

    • 优点:编译期静态检查异步调用,性能最优。
    • 缺点:引入”函数染色”问题(async 函数只能被 async 函数调用),现有同步库需重写为 async 版本,破坏 Java 生态兼容性。
  2. 方案 B:协程 + suspend 关键字(Kotlin 风格)

    • 优点:无函数染色问题,语法简洁。
    • 缺点:JVM 上实现 suspend 需要修改字节码生成 CPS(Continuation-Passing Style)变换,与现有字节码工具(ASM、CGLIB)不兼容;Java 与 Kotlin 互操作时 suspend 函数在 Java 侧需手动处理 Continuation 参数。
  3. 方案 C:虚拟线程(最终选择)

    • 优点:
      • 无新关键字:Thread、Runnable、Callable API 完全复用。
      • 无函数染色:同步代码在虚拟线程中自动获得异步性能。
      • 生态兼容:JDBC 4.0+、HttpClient、Socket、NIO 阻塞操作自动卸载。
    • 缺点:
      • JVM 实现复杂:需修改 HotSpot 的 Thread 类、栈管理、Continuation 实现。
      • Pinning 陷阱:synchronized 块内阻塞会导致载体线程被占用,需开发者主动用 ReentrantLock 替换。
      • ThreadLocal 内存陷阱:百万级虚拟线程各自持有 ThreadLocal 副本可能导致 OOM。

Java 团队选择方案 C 的核心原因是 “生态兼容性优先” —— Java 生态积累了 25 年的同步库(JDBC、JPA、Jackson、OkHttp),重写这些库为 async 版本的成本远高于在 JVM 层实现虚拟线程。这一决策使 Java 在云原生时代保持了”一次编写、到处运行”的承诺。

1.4 JEP 与虚拟线程相关提案

JEP 编号标题落地版本与状态核心内容
JEP 425Virtual Threads (Preview)JDK 19,预览虚拟线程首次预览,Thread.ofVirtual() API
JEP 436Virtual Threads (Second Preview)JDK 20,预览API 微调,Thread.startVirtualThread 简化
JEP 444Virtual ThreadsJDK 21,正式虚拟线程正式发布,API 稳定
JEP 453Structured Concurrency (Preview)JDK 21,预览StructuredTaskScope 首次预览
JEP 462Structured Concurrency (Second Preview)JDK 22,预览API 简化,ShutdownOnFailure / ShutdownOnSuccess
JEP 463Implicitly Declared Classes and Instance Main MethodsJDK 22,预览与虚拟线程无关,但同期发布(JDK 25 转正为 JEP 512)
JEP 480Structured Concurrency (Third Preview)JDK 23,预览API 进一步稳定
JEP 481Scoped Values (Third Preview)JDK 23,预览ScopedValue 替代 ThreadLocal
JEP 499Structured Concurrency (Fourth Preview)JDK 24,预览API 趋于稳定
JEP 491Synchronized Virtual Threads without PinningJDK 24,正式synchronized/Object.wait() 不再固定载体线程
JEP 505Structured Concurrency (Fifth Preview)JDK 25,预览API 重设计:StructuredTaskScope.open() + Joiner 完成策略
JEP 506Scoped ValuesJDK 25,正式作用域值转正,无需 --enable-preview
JEP 525Structured Concurrency (Sixth Preview)JDK 26,预览仍未转正;anySuccessfulResultOrThrow 更名 anySuccessfulOrThrow

1.5 虚拟线程与 Java 生态的共生关系

虚拟线程并非孤立存在,它与以下生态形成了共生关系:

  • Spring Boot 3.2+:spring.threads.virtual.enabled=true 一键启用,Tomcat、Jetty、Undertow 请求处理自动切换为虚拟线程。
  • Helidon Níma(2023):Oracle 推出的”完全基于虚拟线程的 Web 服务器”,无 Netty 依赖,吞吐量较传统 Web 服务器提升 5-10 倍。
  • Quarkus 3.6+:quarkus.virtual-threads.enabled=true 启用虚拟线程,RESTEasy Reactive 自动适配。
  • JDBC 驱动:MySQL Connector/J 8.0.33+、PostgreSQL JDBC 42.6+ 完全兼容虚拟线程,阻塞 I/O 自动卸载。
  • HttpClient(Java 11+):原生兼容虚拟线程,同步 send 方法在虚拟线程中自动获得异步性能。
  • gRPC Java(1.58+):ManagedChannelBuilder.virtualThreadExecutor() 启用虚拟线程。
  • Reactor 框架:Schedulers.boundedElastic() 在虚拟线程场景下可替换为 Schedulers.fromExecutor(Executors.newVirtualThreadPerTaskExecutor())。

历史轶事:Project Loom 团队在 2019 年 JavaOne 大会上首次公开演示虚拟线程时,现场运行了 1000 万个虚拟线程并发 sleep 10 秒,仅消耗约 12GB 内存。若用平台线程实现同样效果,需 10TB 内存(1000 万 × 1MB),这在物理上不可能。这一演示震撼了整个 Java 社区,加速了虚拟线程的正式发布进程。


2. 形式化定义

2.1 虚拟线程的形式化定义

设 TpT_p 为平台线程集合(操作系统线程),TvT_v 为虚拟线程集合,CC 为载体线程(Carrier Thread)集合,SS 为调度器(Scheduler)。虚拟线程系统可形式化为一个六元组:

V=(Tv,Tp,C,S,Continuation,Mount)\mathcal{V} = (T_v, T_p, C, S, \text{Continuation}, \text{Mount})

其中:

  • TvT_v:虚拟线程集合,∣Tv∣|T_v| 可达 10610^6 量级。
  • TpT_p:平台线程集合,∣Tp∣|T_p| 通常等于 CPU 核心数。
  • C⊆TpC \subseteq T_p:载体线程,由 ForkJoinPool 提供,∣C∣=availableProcessors()|C| = \text{availableProcessors}()。
  • SS:调度器,S:Tv→CS: T_v \to C 将虚拟线程映射到载体线程。
  • Continuation\text{Continuation}:延续对象,保存虚拟线程的栈帧。
  • Mount:Tv×C→State\text{Mount}: T_v \times C \to \text{State}:挂载操作,将虚拟线程绑定到载体线程执行。

2.2 M:N 调度模型

传统平台线程采用 1:1 调度:一个 Java 线程对应一个操作系统线程。

PlatformThread:Java Thread↔1:1OS Thread\text{PlatformThread}: \text{Java Thread} \xleftrightarrow{1:1} \text{OS Thread}

虚拟线程采用 M:N 调度:M 个虚拟线程映射到 N 个载体线程。

VirtualThread:VT1,VT2,…,VTM⏟M 个虚拟线程→SC1,C2,…,CN⏟N 个载体线程\text{VirtualThread}: \underbrace{\text{VT}_1, \text{VT}_2, \ldots, \text{VT}_M}_{\text{M 个虚拟线程}} \xrightarrow{S} \underbrace{C_1, C_2, \ldots, C_N}_{\text{N 个载体线程}}

其中 M≫NM \gg N(典型值 M=106M = 10^6,N=CPU 核心数N = \text{CPU 核心数})。调度器 SS 负责在载体线程上分时执行虚拟线程,虚拟线程阻塞时自动让出载体线程。

2.3 虚拟线程状态机

虚拟线程的状态机可形式化为:

State(VT)∈{NEW,RUNNABLE,PARKED,PINNED,TERMINATED}\text{State}(VT) \in \{ \text{NEW}, \text{RUNNABLE}, \text{PARKED}, \text{PINNED}, \text{TERMINATED} \}

状态转移规则:

NEW→start()RUNNABLE\text{NEW} \xrightarrow{\text{start()}} \text{RUNNABLE} RUNNABLE→I/O, park, sleepPARKED→unpark, I/O readyRUNNABLE\text{RUNNABLE} \xrightarrow{\text{I/O, park, sleep}} \text{PARKED} \xrightarrow{\text{unpark, I/O ready}} \text{RUNNABLE} RUNNABLE→synchronized, native, JNIPINNED→block 释放RUNNABLE\text{RUNNABLE} \xrightarrow{\text{synchronized, native, JNI}} \text{PINNED} \xrightarrow{\text{block 释放}} \text{RUNNABLE} RUNNABLE→native, JNI(JDK 23- 为 synchronized)PINNED→block 释放RUNNABLE\text{RUNNABLE} \xrightarrow{\text{native, JNI(JDK 23- 为 synchronized)}} \text{PINNED} \xrightarrow{\text{block 释放}} \text{RUNNABLE}

关键区别:

  • PARKED:虚拟线程卸载(Unmount),载体线程释放,可执行其他虚拟线程。
  • PINNED:虚拟线程无法卸载,载体线程被占用,无法执行其他虚拟线程(退化为平台线程行为)。

版本提示:JDK 24 的 JEP 491 重写了 ObjectMonitor 与虚拟线程的协作, synchronized 块内阻塞与 Object.wait() 在 JDK 24+ 不再触发 Pinning; 此后 Pinning 只在执行 native 方法 / JNI 代码时发生。 依赖”JDK 21 必须 ReentrantLock 替换 synchronized”的老结论前,先确认运行版本。

2.4 Continuation 的形式化定义

Continuation 是虚拟线程的核心抽象,表示”计算的剩余部分”。形式化地,设 cc 为一个 Continuation,ss 为当前栈帧序列,ee 为执行环境(局部变量、操作数栈),则:

c=(s,e,nextInstr)c = (s, e, \text{nextInstr})

虚拟线程的执行可建模为 Continuation 的挂载与卸载:

Mount(c,Ci)=在载体线程 Ci 上恢复 c 的执行\text{Mount}(c, C_i) = \text{在载体线程 } C_i \text{ 上恢复 } c \text{ 的执行} Unmount(Ci)=(c′,Ci),其中 c′ 保存当前栈帧到堆\text{Unmount}(C_i) = (c', C_i) \text{,其中 } c' \text{ 保存当前栈帧到堆}

Continuation 的栈帧存储在堆上(而非操作系统栈),这是虚拟线程内存占用极低的根本原因。每个 Continuation 的初始栈约 1KB,可按需增长到 100KB-1MB(极少见)。

2.5 内存占用形式化对比

设 nn 为线程数量,MpM_p 为平台线程总内存,MvM_v 为虚拟线程总内存:

Mp(n)=n×StackSizep,StackSizep≈1MBM_p(n) = n \times \text{StackSize}_p, \quad \text{StackSize}_p \approx 1\text{MB} Mv(n)=n×StackSizev+N×CarrierStack,StackSizev≈1KB,N≈CPU 核心数M_v(n) = n \times \text{StackSize}_v + N \times \text{CarrierStack}, \quad \text{StackSize}_v \approx 1\text{KB}, N \approx \text{CPU 核心数}

对于 n=106n = 10^6 个线程:

Mp(106)=106×1MB=1TBM_p(10^6) = 10^6 \times 1\text{MB} = 1\text{TB} Mv(106)=106×1KB+8×1MB≈1GB+8MB≈1GBM_v(10^6) = 10^6 \times 1\text{KB} + 8 \times 1\text{MB} \approx 1\text{GB} + 8\text{MB} \approx 1\text{GB}

虚拟线程的内存优势为 1000 倍,这是其能支撑百万级并发的物理基础。

2.6 吞吐量模型

设 TcpuT_{\text{cpu}} 为任务 CPU 执行时间,TioT_{\text{io}} 为任务 I/O 等待时间,NN 为载体线程数。对于 I/O 密集型任务(Tio≫TcpuT_{\text{io}} \gg T_{\text{cpu}}):

平台线程模型:线程数 npn_p 受限于内存,np≤4000n_p \leq 4000。吞吐量:

Throughputp=npTcpu+Tio≤4000Tcpu+Tio\text{Throughput}_p = \frac{n_p}{T_{\text{cpu}} + T_{\text{io}}} \leq \frac{4000}{T_{\text{cpu}} + T_{\text{io}}}

虚拟线程模型:虚拟线程数 nvn_v 可达 10610^6,但受限于载体线程数 NN。吞吐量:

Throughputv=NTcpu\text{Throughput}_v = \frac{N}{T_{\text{cpu}}}

因为虚拟线程在 I/O 等待时释放载体线程,载体线程始终在执行 CPU 工作。对于 Tio=100msT_{\text{io}} = 100\text{ms},Tcpu=1msT_{\text{cpu}} = 1\text{ms},N=8N = 8:

Throughputv=81ms=8000 请求/秒\text{Throughput}_v = \frac{8}{1\text{ms}} = 8000 \text{ 请求/秒} Throughputp≤4000101ms≈40 请求/秒\text{Throughput}_p \leq \frac{4000}{101\text{ms}} \approx 40 \text{ 请求/秒}

虚拟线程的吞吐量优势为 200 倍(在 I/O 密集型场景)。


3. 理论推导:JVM 视角下的虚拟线程机制

3.1 虚拟线程的 JVM 实现

虚拟线程在 HotSpot JVM 中的实现涉及以下核心组件:

  1. java.lang.VirtualThread(JDK 21+):继承自 Thread,是虚拟线程的 Java 层表示。
  2. jdk.internal.vm.Continuation:JVM 内部类,封装 Continuation 机制。
  3. ForkJoinPool:虚拟线程的载体线程池,默认 Runtime.getRuntime().availableProcessors() 个工作线程。
  4. VirtualThread 的 run() 方法:在载体线程上执行 Continuation。
  5. Continuation.yield():虚拟线程阻塞时调用,将栈帧保存到堆。
flowchart TD
    B0["VirtualThread | <------> | Continuation | VirtualThread / vthread | Mount | stack frames | class metadata / carrier | <------> | locals / cont | Unmount | operands / state"]
    B1["Continuation / stack (堆上) | <-- 阻塞时栈帧保存到此 / nextInstr"]
    B0 --> B1

3.2 Continuation 的栈帧存储机制

传统平台线程的栈帧存储在操作系统分配的线程栈上(1MB 连续内存)。虚拟线程的栈帧存储分为两种状态:

运行态(Mounted):虚拟线程在载体线程上执行时,栈帧存储在载体线程的栈上。

阻塞态(Unmounted):虚拟线程阻塞时,Continuation.yield() 被调用,栈帧被”冻结”并复制到堆上的 Continuation 对象中。

flowchart TD
    B0["main() / VirtualThread.run() / fetchUser() / socket.read()  <-- 阻塞点"]
    B1["yield()"]
    B0 --> B1
    B2["Continuation / fetchUser frame / VT.run frame / nextInstr: ..."]
    B1 --> B2
    B3["main() / (空闲,可执行其他VT)"]
    B2 --> B3

这一机制的关键是 栈帧的可序列化:JVM 需要将栈帧从载体线程栈”卸下”并保存到堆,恢复时再”挂回”载体线程栈。这要求 JVM 修改 Continuation 的实现,使其支持栈帧的拷贝与恢复。

3.3 卸载触发点:哪些操作会触发 Unmount

虚拟线程在以下操作中会自动卸载(Unmount):

  1. java.net.Socket I/O:Socket.getInputStream().read()、Socket.getOutputStream().write() 阻塞时。
  2. java.nio.channels:SocketChannel.read()、ServerSocketChannel.accept() 阻塞时(NIO 阻塞模式)。
  3. java.net.http.HttpClient:HttpClient.send() 同步方法阻塞时。
  4. java.io.FileInputStream / FileOutputStream:阻塞 I/O(注意:文件 I/O 通常不卸载,因磁盘 I/O 极快)。
  5. Thread.sleep():sleep 时自动卸载。
  6. Object.wait() / Condition.await():等待时自动卸载。
  7. LockSupport.park():显式 park 时卸载。
  8. BlockingQueue 操作:put()、take() 阻塞时卸载。
  9. Semaphore.acquire():许可不足时卸载。
  10. CountDownLatch.await():等待时卸载。
  11. Future.get():等待结果时卸载。
  12. CompletableFuture.join():等待完成时卸载。

JDK 21 对以上所有 API 都做了虚拟线程适配,业务代码无需任何修改即可享受卸载机制。

3.4 Pinning(线程固定)机制

Pinning 是虚拟线程”无法卸载”的状态,会导致载体线程被占用。以 JDK 24 为分界线:

  • JDK 21-23:synchronized 块内阻塞、Object.wait()、native / JNI 三类场景都会 Pinning。
  • JDK 24+:JEP 491 解决了 monitor 相关 Pinning,只剩 native / JNI 两类场景。

3.4.1 synchronized 块内阻塞

public synchronized void process() {  // 进入 synchronized 块
    socket.read();  // 阻塞 I/O,但虚拟线程被 Pinning
}

JVM 内部原因:

JVM 的 monitor(监视器锁)实现依赖操作系统层(ObjectMonitor),monitor 的 wait/enter 操作涉及操作系统互斥量(pthread_mutex)。虚拟线程在 synchronized 块内阻塞时,JVM 无法将栈帧卸载(因为 monitor 持有者是载体线程),导致载体线程被占用。

HotSpot 的 ObjectMonitor 结构在 JDK 21-23 未针对虚拟线程适配,这是当时 Pinning 的根本原因。JDK 24 已通过 JEP 491 正式解决:monitor 的持有状态被虚拟线程化,synchronized 块内阻塞(以及 Object.wait())不再固定载体线程,虚拟线程可以正常卸载。JDK 24+ 仍会 Pinning 的只剩 native / JNI 两类场景。

3.4.2 native 方法调用

public native void doNativeWork();  // native 方法

public void caller() {
    doNativeWork();  // 虚拟线程被 Pinning,因 native 方法栈无法卸载
}

JVM 内部原因:

native 方法的栈帧存储在 native 栈上(C/C++ 栈),JVM 无法访问和拷贝 native 栈。虚拟线程在执行 native 方法时被 Pinning,载体线程被占用直到 native 方法返回。

3.4.3 JNI 调用

JNI(Java Native Interface)调用同样会导致 Pinning,原因与 native 方法相同。

3.5 Pinning 的检测与监控

JDK 21 提供了 Pinning 检测机制:

方式 1:JVM 诊断选项

# 启动时启用 Pinning 诊断(仅 JDK 21-23 有效;
# JDK 24 起 JEP 491 移除了该属性——monitor 类 Pinning 已不存在,
# 剩余的 native/JNI Pinning 请用下面的 JFR 事件监控)
java -Djdk.tracePinnedThreads=short -jar app.jar
# 或完整堆栈
java -Djdk.tracePinnedThreads=full -jar app.jar

启用后,每当虚拟线程被 Pinning,JVM 会打印堆栈:

Thread[#123,ForkJoinPool-1-worker-1] pinned due to:
    java.base/java.lang.VirtualThread$VThreadContinuation.onPinned(VirtualThread.java:xxx)
    app.MyService.process(MyService.java:45)

方式 2:jcmd 线程转储

# 抓取所有线程(含虚拟线程)的 JSON 转储
jcmd <pid> Thread.dump_to_file -format=json thread_dump.json

JSON 转储中会标注每个虚拟线程的状态:

{
  "threadId": 123,
  "name": "VirtualThread[#123]/runnable@ForkJoinPool-1-worker-1",
  "virtual": true,
  "state": "RUNNABLE",
  "pinned": true,
  "stackTrace": [
    {"class": "app.MyService", "method": "process", "line": 45}
  ]
}

方式 3:JFR(Java Flight Recorder)事件

# 启动 JFR 录制,捕获 Pinning 事件
java -XX:StartFlightRecording=duration=60s,filename=pinning.jfr -jar app.jar

JFR 事件 jdk.VirtualThreadPinned 记录每次 Pinning 的详细信息,可通过 JDK Mission Control 分析。

3.6 载体线程调度器:ForkJoinPool

虚拟线程的载体线程由 ForkJoinPool 提供,默认配置:

// 虚拟线程默认调度器(JDK 21 内部实现)
ForkJoinPool virtualThreadScheduler = new ForkJoinPool(
    Runtime.getRuntime().availableProcessors(),  // 并行度 = CPU 核心数
    ForkJoinPool.defaultForkJoinWorkerThreadFactory,
    null,
    false  // 不启用 asyncMode
);

关键特性:

  • 并行度:默认等于 CPU 核心数,可通过 -Djdk.virtualThreadParallelism=N 调整。
  • work-stealing:空闲工作线程会从其他工作线程的队列尾部”窃取”任务,均衡负载。
  • 不可替换:应用代码无法替换虚拟线程的调度器(与平台线程不同)。

3.7 Continuation 的实现:yield 与 resume

Continuation 的核心方法:

public class Continuation {
    // 卸载:将当前栈帧保存到堆
    public static void yield() {
        // JVM 内部:保存栈帧到 Continuation 对象
        // 抛出 YieldException 用于栈展开
    }
    
    // 恢复:从堆加载栈帧并恢复执行
    public void run() {
        // JVM 内部:从 Continuation 对象恢复栈帧
        // 在当前载体线程上继续执行
    }
}

yield() 的实现机制:

  1. JVM 检测到阻塞操作(如 socket.read())。
  2. JVM 调用 Continuation.yield()。
  3. yield() 遍历当前栈帧,将每帧的局部变量、操作数栈、返回地址保存到 Continuation 对象(存储在堆上)。
  4. yield() 抛出 YieldException,用于栈展开(stack unwinding)。
  5. 载体线程捕获 YieldException,释放虚拟线程,继续调度其他虚拟线程。

run() 的实现机制:

  1. 载体线程从调度器队列取出 Continuation 对象。
  2. 调用 Continuation.run()。
  3. JVM 从 Continuation 对象恢复栈帧到当前载体线程栈。
  4. 从 nextInstr 处继续执行。

这一机制对开发者透明,业务代码无需感知。


4. 代码示例

4.1 创建虚拟线程的三种方式

package com.fandex.virtualthread;

import java.time.Duration;
import java.util.concurrent.Executors;
import java.util.concurrent.ThreadFactory;
import java.util.stream.IntStream;

/**
 * 虚拟线程创建示例
 * 演示三种创建虚拟线程的方式及其适用场景
 */
public class VirtualThreadCreationDemo {

    /**
     * 方式一:直接创建并启动
     * 适用于一次性任务,无需持有 Thread 引用
     */
    public static void createAndStart() {
        Thread vt = Thread.startVirtualThread(() -> {
            System.out.println("虚拟线程运行中: " + Thread.currentThread());
            System.out.println("是否虚拟线程: " + Thread.currentThread().isVirtual());
        });
        
        // 等待虚拟线程结束
        try {
            vt.join();
        } catch (InterruptedException e) {
            Thread.currentThread().interrupt();
        }
    }

    /**
     * 方式二:使用 Builder API
     * 适用于需要自定义线程名的场景
     */
    public static void createWithBuilder() {
        Thread vt = Thread.ofVirtual()
            .name("my-vthread-", 0)  // 名字前缀 + 起始编号
            .start(() -> {
                System.out.println("线程名: " + Thread.currentThread().getName());
                try {
                    Thread.sleep(Duration.ofMillis(100));
                } catch (InterruptedException e) {
                    Thread.currentThread().interrupt();
                }
            });
        
        try {
            vt.join();
        } catch (InterruptedException e) {
            Thread.currentThread().interrupt();
        }
    }

    /**
     * 方式三:使用 ThreadFactory
     * 适用于需要与现有 API(如 ExecutorService)配合的场景
     */
    public static void createWithFactory() {
        ThreadFactory factory = Thread.ofVirtual()
            .name("worker-", 0)
            .factory();
        
        Thread vt = factory.newThread(() -> {
            System.out.println("工厂创建的虚拟线程: " + Thread.currentThread().getName());
        });
        vt.start();
        
        try {
            vt.join();
        } catch (InterruptedException e) {
            Thread.currentThread().interrupt();
        }
    }

    /**
     * 方式四:使用 newVirtualThreadPerTaskExecutor
     * 适用于批量提交任务的场景(推荐)
     */
    public static void createWithExecutor() {
        try (var executor = Executors.newVirtualThreadPerTaskExecutor()) {
            // 提交 10000 个任务,每个任务在独立虚拟线程中执行
            IntStream.range(0, 10000).forEach(i -> {
                executor.submit(() -> {
                    try {
                        Thread.sleep(Duration.ofSeconds(1));  // 模拟 I/O
                    } catch (InterruptedException e) {
                        Thread.currentThread().interrupt();
                    }
                    return i;
                });
            });
        }  // try-with-resources 自动等待所有任务完成
    }

    public static void main(String[] args) {
        System.out.println("=== 方式一:直接创建 ===");
        createAndStart();
        
        System.out.println("\n=== 方式二:Builder API ===");
        createWithBuilder();
        
        System.out.println("\n=== 方式三:ThreadFactory ===");
        createWithFactory();
        
        System.out.println("\n=== 方式四:批量执行器 ===");
        long start = System.currentTimeMillis();
        createWithExecutor();
        System.out.println("10000 个虚拟线程并发执行耗时: " 
            + (System.currentTimeMillis() - start) + "ms");
    }
}

4.2 虚拟线程并发 HTTP 请求

package com.fandex.virtualthread;

import java.net.URI;
import java.net.http.HttpClient;
import java.net.http.HttpRequest;
import java.net.http.HttpResponse;
import java.time.Duration;
import java.util.List;
import java.util.concurrent.CompletableFuture;
import java.util.concurrent.Executors;
import java.util.stream.Collectors;
import java.util.stream.IntStream;

/**
 * 虚拟线程并发 HTTP 请求示例
 * 对比虚拟线程与平台线程在 HTTP 并发场景下的性能
 */
public class ConcurrentHttpDemo {

    private static final HttpClient HTTP_CLIENT = HttpClient.newBuilder()
        .connectTimeout(Duration.ofSeconds(5))
        .build();

    /**
     * 使用虚拟线程并发请求多个 URL
     * 每个请求在独立虚拟线程中执行,阻塞时自动卸载
     */
    public static List<String> fetchWithVirtualThreads(List<String> urls) {
        try (var executor = Executors.newVirtualThreadPerTaskExecutor()) {
            // 每个提交异步任务,executor 会为每个任务创建一个虚拟线程
            List<CompletableFuture<String>> futures = urls.stream()
                .map(url -> CompletableFuture.supplyAsync(
                    () -> fetchUrl(url), executor))
                .toList();
            
            // 等待所有请求完成
            return futures.stream()
                .map(CompletableFuture::join)
                .toList();
        }
    }

    /**
     * 单个 URL 请求(同步风格,在虚拟线程中自动获得异步性能)
     */
    private static String fetchUrl(String url) {
        try {
            HttpRequest request = HttpRequest.newBuilder()
                .uri(URI.create(url))
                .timeout(Duration.ofSeconds(10))
                .GET()
                .build();
            
            // send() 是阻塞方法,但在虚拟线程中会自动卸载
            HttpResponse<String> response = HTTP_CLIENT.send(
                request, HttpResponse.BodyHandlers.ofString());
            
            return url + " -> " + response.statusCode();
        } catch (Exception e) {
            return url + " -> ERROR: " + e.getMessage();
        }
    }

    public static void main(String[] args) {
        // 生成 100 个测试 URL
        List<String> urls = IntStream.range(0, 100)
            .mapToObj(i -> "https://httpbin.org/delay/" + (i % 3 + 1))
            .toList();
        
        long start = System.currentTimeMillis();
        List<String> results = fetchWithVirtualThreads(urls);
        long elapsed = System.currentTimeMillis() - start;
        
        System.out.println("100 个并发请求耗时: " + elapsed + "ms");
        System.out.println("平均每个请求耗时: " + (elapsed / 100) + "ms");
        results.stream().limit(5).forEach(System.out::println);
    }
}

4.3 结构化并发示例

API 版本提示:结构化并发截至 JDK 26 仍是第 6 次预览(JEP 525)。 JDK 21-24 的 new ShutdownOnFailure() 写法在 JDK 25(JEP 505)起已被移除, 下面示例按 JDK 25/26 预览 API(open() + Joiner)编写,运行需 --enable-preview。

package com.fandex.virtualthread;

import java.util.concurrent.StructuredTaskScope;
import java.util.concurrent.StructuredTaskScope.Joiner;
import java.util.concurrent.StructuredTaskScope.Subtask;
import java.util.concurrent.TimeUnit;

/**
 * 结构化并发示例
 * 演示 open() 默认策略(全部成功或快速失败)与 anySuccessfulOrThrow 竞速策略
 */
public class StructuredConcurrencyDemo {

    /**
     * 订单详情聚合服务
     * 默认策略:任一子任务失败则取消所有子任务并抛 FailedException
     */
    public record OrderDetail(String order, String user, String payment) {}

    public OrderDetail fetchOrderDetail(Long orderId) throws InterruptedException {
        // open() 无参版本:任一子任务失败,自动取消其他子任务
        try (var scope = StructuredTaskScope.open()) {
            // 并发 fork 三个子任务
            Subtask<String> orderTask = scope.fork(() -> fetchOrder(orderId));
            Subtask<String> userTask = scope.fork(() -> fetchUser(orderId));
            Subtask<String> paymentTask = scope.fork(() -> fetchPayment(orderId));

            // 等待所有子任务完成;任一失败时抛出 FailedException(非受检异常)
            scope.join();

            // 所有子任务成功,组装结果
            return new OrderDetail(
                orderTask.get(),
                userTask.get(),
                paymentTask.get()
            );
        }
    }

    /**
     * 竞速策略:任一子任务成功则取消其他子任务(join() 返回首个成功结果),
     * 全部失败时 join() 抛 FailedException
     * 适用于"多源竞速"场景(如多机房读同一数据,取最快响应)
     */
    public String fetchFromFastestSource(String key) throws InterruptedException {
        try (var scope = StructuredTaskScope.open(Joiner.<String>anySuccessfulOrThrow())) {
            // 并发 fork 三个数据源
            scope.fork(() -> fetchFromRedis(key));
            scope.fork(() -> fetchFromMySQL(key));
            scope.fork(() -> fetchFromES(key));

            // 等待任一子任务成功(其他自动取消),返回最先成功的结果
            return scope.join();
        }
    }

    // 模拟业务方法
    private String fetchOrder(Long id) throws InterruptedException {
        TimeUnit.MILLISECONDS.sleep(50);  // 模拟数据库查询
        return "Order-" + id;
    }

    private String fetchUser(Long id) throws InterruptedException {
        TimeUnit.MILLISECONDS.sleep(80);  // 模拟用户服务调用
        return "User-" + id;
    }

    private String fetchPayment(Long id) throws InterruptedException {
        TimeUnit.MILLISECONDS.sleep(100);  // 模拟支付服务调用
        return "Payment-" + id;
    }

    private String fetchFromRedis(String key) throws InterruptedException {
        TimeUnit.MILLISECONDS.sleep(30);
        return "redis:" + key;
    }

    private String fetchFromMySQL(String key) throws InterruptedException {
        TimeUnit.MILLISECONDS.sleep(80);
        return "mysql:" + key;
    }

    private String fetchFromES(String key) throws InterruptedException {
        TimeUnit.MILLISECONDS.sleep(120);
        return "es:" + key;
    }

    public static void main(String[] args) throws Exception {
        StructuredConcurrencyDemo demo = new StructuredConcurrencyDemo();
        
        System.out.println("=== 订单详情聚合 ===");
        long start = System.currentTimeMillis();
        OrderDetail detail = demo.fetchOrderDetail(1001L);
        System.out.println("耗时: " + (System.currentTimeMillis() - start) + "ms");
        System.out.println(detail);
        
        System.out.println("\n=== 多源竞速 ===");
        start = System.currentTimeMillis();
        String result = demo.fetchFromFastestSource("user:1001");
        System.out.println("耗时: " + (System.currentTimeMillis() - start) + "ms");
        System.out.println("最快响应: " + result);
    }
}

4.4 Pinning 检测与规避

package com.fandex.virtualthread;

import java.time.Duration;
import java.util.concurrent.Executors;
import java.util.concurrent.locks.ReentrantLock;

/**
 * Pinning 检测与规避示例
 * 演示 synchronized 导致 Pinning 的反模式,以及 ReentrantLock 替代方案
 */
public class PinningDemo {

    /**
     * 反模式:synchronized 块内阻塞导致 Pinning
     * 运行时加 -Djdk.tracePinnedThreads=short 可看到 Pinning 警告
     */
    public static class BadPinningService {
        private int counter = 0;

        // 错误:synchronized + I/O 阻塞 = Pinning
        public synchronized void process(String url) {
            try {
                // 模拟 I/O 阻塞(在 synchronized 块内)
                Thread.sleep(Duration.ofMillis(100));
                counter++;
            } catch (InterruptedException e) {
                Thread.currentThread().interrupt();
            }
        }
    }

    /**
     * 正确模式:使用 ReentrantLock 替代 synchronized
     * 虚拟线程在 lock.lock() 阻塞时可正常卸载
     */
    public static class GoodService {
        private final ReentrantLock lock = new ReentrantLock();
        private int counter = 0;

        public void process(String url) {
            lock.lock();  // 阻塞时虚拟线程可卸载
            try {
                // 模拟 I/O 阻塞
                Thread.sleep(Duration.ofMillis(100));
                counter++;
            } catch (InterruptedException e) {
                Thread.currentThread().interrupt();
            } finally {
                lock.unlock();
            }
        }
    }

    public static void main(String[] args) throws InterruptedException {
        // 测试 synchronized 版本(会有 Pinning)
        System.out.println("=== 测试 synchronized 版本 ===");
        BadPinningService badService = new BadPinningService();
        long start = System.currentTimeMillis();
        try (var executor = Executors.newVirtualThreadPerTaskExecutor()) {
            for (int i = 0; i < 100; i++) {
                executor.submit(() -> badService.process("url-" + i));
            }
        }
        System.out.println("synchronized 版本耗时: " 
            + (System.currentTimeMillis() - start) + "ms");
        
        // 测试 ReentrantLock 版本(无 Pinning)
        System.out.println("\n=== 测试 ReentrantLock 版本 ===");
        GoodService goodService = new GoodService();
        start = System.currentTimeMillis();
        try (var executor = Executors.newVirtualThreadPerTaskExecutor()) {
            for (int i = 0; i < 100; i++) {
                executor.submit(() -> goodService.process("url-" + i));
            }
        }
        System.out.println("ReentrantLock 版本耗时: " 
            + (System.currentTimeMillis() - start) + "ms");
    }
}

4.5 作用域值(ScopedValue)示例

package com.fandex.virtualthread;

import java.util.concurrent.Executors;
import java.util.concurrent.StructuredTaskScope;

/**
 * ScopedValue 示例
 * 替代 ThreadLocal,在虚拟线程场景下更安全、更高效
 */
public class ScopedValueDemo {

    // 定义 ScopedValue(不可变,作用域内绑定)
    private static final ScopedValue<String> USER_ID = ScopedValue.newInstance();
    private static final ScopedValue<String> LOCALE = ScopedValue.newInstance();

    /**
     * 使用 ScopedValue.where 绑定值,在作用域内所有子任务可读
     */
    public static void processUserRequest(String userId, String locale) {
        ScopedValue.where(USER_ID, userId).where(LOCALE, locale).run(() -> {
            // 当前线程及所有结构化并发子任务可读取 USER_ID 与 LOCALE
            handleRequest();
        });
    }

    /**
     * 在请求处理链中读取 ScopedValue
     * 无需方法参数传递,且保证不可变
     */
    private static void handleRequest() {
        System.out.println("处理请求: user=" + USER_ID.get() 
            + ", locale=" + LOCALE.get());
        
        // 启动结构化并发子任务,子任务自动继承 ScopedValue
        try (var scope = StructuredTaskScope.open()) {
            scope.fork(() -> fetchUserProfile(USER_ID.get()));
            scope.fork(() -> fetchUserOrders(USER_ID.get(), LOCALE.get()));
            scope.join();
        }
    }

    private static String fetchUserProfile(String userId) {
        System.out.println("  子任务读取 ScopedValue: user=" + userId);
        return "profile-" + userId;
    }

    private static String fetchUserOrders(String userId, String locale) {
        System.out.println("  子任务读取 ScopedValue: user=" + userId 
            + ", locale=" + locale);
        return "orders-" + userId;
    }

    public static void main(String[] args) {
        // 模拟并发处理多个用户请求
        try (var executor = Executors.newVirtualThreadPerTaskExecutor()) {
            executor.submit(() -> processUserRequest("user-001", "zh-CN"));
            executor.submit(() -> processUserRequest("user-002", "en-US"));
            executor.submit(() -> processUserRequest("user-003", "ja-JP"));
        }
    }
}

4.6 虚拟线程与 ThreadLocal 的对比

package com.fandex.virtualthread;

import java.util.concurrent.Executors;
import java.util.stream.IntStream;

/**
 * 虚拟线程下 ThreadLocal 的内存陷阱演示
 * 百万级虚拟线程各自持有 ThreadLocal 副本可能导致 OOM
 */
public class ThreadLocalTrapDemo {

    // 每个线程持有 1MB 数据的 ThreadLocal
    private static final ThreadLocal<byte[]> LARGE_DATA = ThreadLocal.withInitial(
        () -> new byte[1024 * 1024]);  // 1MB

    /**
     * 反模式:百万虚拟线程 + 大 ThreadLocal = OOM
     */
    public static void oomDemo() {
        try (var executor = Executors.newVirtualThreadPerTaskExecutor()) {
            IntStream.range(0, 1_000_000).forEach(i -> {
                executor.submit(() -> {
                    // 每个虚拟线程都会初始化 ThreadLocal(1MB)
                    byte[] data = LARGE_DATA.get();
                    // 业务处理...
                    return null;
                });
            });
        }
        // 总内存占用:1,000,000 × 1MB = 1TB,必然 OOM
    }

    /**
     * 正确模式:使用 ScopedValue 或避免 ThreadLocal
     */
    public static void safeDemo() {
        // 方案 1:使用 ScopedValue(不可变,作用域绑定)
        // 方案 2:将共享数据作为方法参数传递
        // 方案 3:使用局部变量而非 ThreadLocal
        System.out.println("推荐使用 ScopedValue 或显式参数传递");
    }

    public static void main(String[] args) {
        safeDemo();
        // 不要运行 oomDemo(),会导致 OOM
    }
}

4.7 Spring Boot 3.2+ 虚拟线程集成

package com.fandex.virtualthread;

import org.springframework.boot.SpringApplication;
import org.springframework.boot.autoconfigure.SpringBootApplication;
import org.springframework.context.annotation.Bean;
import org.springframework.core.task.AsyncTaskExecutor;
import org.springframework.core.task.support.TaskExecutorAdapter;
import org.springframework.boot.web.embedded.tomcat.TomcatProtocolHandlerCustomizer;
import org.springframework.web.bind.annotation.GetMapping;
import org.springframework.web.bind.annotation.RestController;
import org.springframework.web.client.RestClient;

import java.time.Duration;
import java.util.concurrent.Executors;

/**
 * Spring Boot 3.2+ 虚拟线程集成示例
 * 演示 application.yml 配置与显式 Bean 配置两种方式
 */
@SpringBootApplication
public class VirtualThreadSpringBootApp {

    public static void main(String[] args) {
        SpringApplication.run(VirtualThreadSpringBootApp.class, args);
    }

    /**
     * 方式一:application.yml 配置(推荐)
     * spring:
     *   threads:
     *     virtual:
     *       enabled: true
     * 
     * Spring Boot 3.2+ 会自动配置 Tomcat/Jetty/Undertow 使用虚拟线程
     */

    /**
     * 方式二:显式 Bean 配置(更细粒度控制)
     * 适用于需要自定义虚拟线程名称或异常处理的场景
     */
    // @Configuration
    public static class VirtualThreadConfig {

        @Bean
        public TomcatProtocolHandlerCustomizer<?> protocolHandlerVirtualThreadExecutorCustomizer() {
            return protocolHandler -> {
                protocolHandler.setExecutor(
                    Executors.newVirtualThreadPerTaskExecutor());
            };
        }

        @Bean
        public AsyncTaskExecutor applicationTaskExecutor() {
            return new TaskExecutorAdapter(
                Executors.newVirtualThreadPerTaskExecutor());
        }
    }

    @RestController
    public static class ApiController {

        private final RestClient restClient = RestClient.create();

        /**
         * 同步阻塞风格,但在虚拟线程中自动获得异步性能
         */
        @GetMapping("/api/aggregate")
        public String aggregate() {
            // 每次调用阻塞 100ms,但在虚拟线程中载体线程不会被占用
            String user = fetchUser();
            String order = fetchOrder();
            String payment = fetchPayment();
            return user + "|" + order + "|" + payment;
        }

        private String fetchUser() {
            try {
                Thread.sleep(Duration.ofMillis(100));
            } catch (InterruptedException e) {
                Thread.currentThread().interrupt();
            }
            return "user";
        }

        private String fetchOrder() {
            try {
                Thread.sleep(Duration.ofMillis(100));
            } catch (InterruptedException e) {
                Thread.currentThread().interrupt();
            }
            return "order";
        }

        private String fetchPayment() {
            try {
                Thread.sleep(Duration.ofMillis(100));
            } catch (InterruptedException e) {
                Thread.currentThread().interrupt();
            }
            return "payment";
        }
    }
}

4.8 虚拟线程性能基准测试

package com.fandex.virtualthread;

import java.time.Duration;
import java.util.concurrent.*;
import java.util.stream.IntStream;

/**
 * 虚拟线程性能基准测试
 * 对比虚拟线程与平台线程池在 I/O 密集型任务下的吞吐量
 */
public class PerformanceBenchmark {

    /**
     * 模拟 I/O 密集型任务(100ms 阻塞)
     */
    private static String ioTask(int i) {
        try {
            Thread.sleep(Duration.ofMillis(100));
        } catch (InterruptedException e) {
            Thread.currentThread().interrupt();
        }
        return "result-" + i;
    }

    /**
     * 平台线程池测试
     */
    public static long testPlatformThreadPool(int taskCount) throws InterruptedException {
        ExecutorService executor = Executors.newFixedThreadPool(200);
        long start = System.currentTimeMillis();
        try {
            var futures = IntStream.range(0, taskCount)
                .mapToObj(i -> executor.submit(() -> ioTask(i)))
                .toList();
            for (var f : futures) {
                try { f.get(); } catch (ExecutionException ignored) {}
            }
        } finally {
            executor.shutdown();
            executor.awaitTermination(1, TimeUnit.MINUTES);
        }
        return System.currentTimeMillis() - start;
    }

    /**
     * 虚拟线程测试
     */
    public static long testVirtualThreads(int taskCount) throws InterruptedException {
        try (var executor = Executors.newVirtualThreadPerTaskExecutor()) {
            long start = System.currentTimeMillis();
            var futures = IntStream.range(0, taskCount)
                .mapToObj(i -> executor.submit(() -> ioTask(i)))
                .toList();
            for (var f : futures) {
                try { f.get(); } catch (ExecutionException ignored) {}
            }
            return System.currentTimeMillis() - start;
        }
    }

    public static void main(String[] args) throws InterruptedException {
        int taskCount = 1000;
        
        // 预热
        testVirtualThreads(100);
        
        System.out.println("任务数: " + taskCount);
        
        long platformTime = testPlatformThreadPool(taskCount);
        System.out.println("平台线程池(200 线程): " + platformTime + "ms");
        
        long virtualTime = testVirtualThreads(taskCount);
        System.out.println("虚拟线程: " + virtualTime + "ms");
        
        System.out.println("虚拟线程加速比: " + (platformTime * 1.0 / virtualTime) + "x");
    }
}

5. 对比分析

5.1 虚拟线程 vs 平台线程

维度平台线程(Platform Thread)虚拟线程(Virtual Thread)
调度模型1:1(Java Thread : OS Thread)M:N(M 个虚拟线程 : N 个载体线程)
创建成本高(clone 系统调用,约 1-10μs)极低(JVM 内部分配,约 0.1-1μs)
内存占用约 1MB 栈空间(-Xss1m)初始约 1KB,可按需增长
数量上限典型 4000-8000(受内存与 ulimit 限制)百万级(受堆内存限制)
阻塞行为阻塞时占用 OS 线程阻塞时卸载,释放载体线程
CPU 密集型适合(直接 OS 调度)不适合(M:N 调度有额外开销)
I/O 密集型不适合(线程数受限)非常适合(百万级并发)
ThreadLocal安全(线程数有限)需谨慎(百万线程 × ThreadLocal = OOM 风险)
synchronized正常工作导致 Pinning(JDK 24+ 改进)
池化推荐(创建成本高)反模式(创建成本极低)
优先级支持(setPriority)不支持(固定为 NORM_PRIORITY)
调试简单(线程数少)复杂(百万线程堆栈)

5.2 虚拟线程 vs 响应式编程(Reactor / RxJava)

维度响应式编程(Reactor)虚拟线程
编程范式异步数据流(Mono、Flux)同步阻塞(Thread.sleep、socket.read)
学习曲线陡峭(需学习组合子)平缓(与传统 Java 一致)
代码可读性较差(flatMap 链式调用)优秀(顺序式代码)
调试困难(异步栈不连续)简单(同步栈)
异常处理复杂(onErrorMap、onErrorResume)简单(try-catch)
与同步库兼容不兼容(JDBC 会阻塞事件循环)完全兼容(JDBC 自动卸载)
背压(Backpressure)原生支持需手动实现(Semaphore)
吞吐量(I/O 密集)高高(接近 Reactor)
延迟低(事件驱动)低(载体线程调度)
生态成熟度成熟(Spring WebFlux)快速成熟(Spring Boot 3.2+)
函数染色是(Mono<T> vs T)否(所有方法返回 T)

选型建议:

  • 新项目:优先选择虚拟线程,编程简单、生态兼容。
  • 存量 Reactor 项目:无需迁移,Reactor 与虚拟线程可共存(Schedulers.fromExecutor)。
  • 流式处理:保留 Reactor(背压、窗口、聚合等操作符强大)。
  • 简单 CRUD 服务:虚拟线程更合适。

5.3 虚拟线程 vs Kotlin 协程

维度Kotlin 协程(kotlinx.coroutines)Java 虚拟线程
实现层级库层(suspend 关键字 + CPS 变换)JVM 层(Continuation + 栈帧拷贝)
API 风格suspend fun + launch/asyncThread.startVirtualThread
函数染色是(suspend 函数只能被 suspend 调用)否
与 Java 互操作需 Continuation 参数适配完全兼容(无新关键字)
性能优秀(编译期 CPS,无运行时开销)优秀(JVM 优化,接近协程)
调试良好(Kotlin 协程调试器)良好(jcmd 线程转储)
生态Kotlin 原生Java 原生(Spring、Quarkus 等)
结构化并发原生支持(coroutineScope)StructuredTaskScope(预览)

5.4 虚拟线程 vs Go goroutine

维度Go goroutineJava 虚拟线程
实现层级运行时层(Go runtime 调度)JVM 层
栈管理分段栈 / 连续栈(按需增长)Continuation(栈帧拷贝到堆)
初始栈大小2KB约 1KB
调度器GMP 模型(Goroutine-Machine-Processor)ForkJoinPool work-stealing
通道(Channel)原生支持(chan)需 BlockingQueue
select原生支持需 Selector / 结构化并发
生态Go 原生Java 生态(Spring 等)
成熟度成熟(Go 1.0+)成熟(Java 21+)

5.5 何时选择虚拟线程

场景推荐原因
高并发 HTTP 网关虚拟线程百万级并发,I/O 密集
微服务聚合层虚拟线程并发调用多个下游服务
数据库连接池虚拟线程JDBC 阻塞 I/O 自动卸载
消息队列消费者虚拟线程消息处理 + I/O 调用
文件 I/O 服务虚拟线程文件读写自动卸载
CPU 密集型计算平台线程虚拟线程 M:N 有额外开销
流式数据处理Reactor背压、窗口、聚合操作符强大
低延迟交易系统平台线程调度延迟敏感

6. 陷阱与反模式

6.1 反模式:在虚拟线程中使用 synchronized

问题代码:

public synchronized String fetchData(String key) {
    return httpClient.send(request);  // 阻塞 I/O 在 synchronized 块内
}

问题分析:

在 JDK 21-23 上,synchronized 块内的阻塞操作会导致虚拟线程被 Pinning,载体线程被占用。100 个虚拟线程并发调用此方法会导致 100 个载体线程被 Pinning(即使载体线程池只有 8 个),吞吐量退化为平台线程水平。JDK 24+(JEP 491)已消除 monitor 类 Pinning,但 ReentrantLock 在所有版本上都是安全选择,且能避免团队在不同 JDK 版本间踩坑。

正确做法:

private final ReentrantLock lock = new ReentrantLock();

public String fetchData(String key) {
    lock.lock();  // 阻塞时虚拟线程可正常卸载(JDK 21+ 所有版本均成立)
    try {
        return httpClient.send(request);
    } finally {
        lock.unlock();
    }
}

工具检测:

JDK 21-23 启动时加 -Djdk.tracePinnedThreads=short,JVM 会打印 Pinning 堆栈,帮助快速定位(JDK 24 起该属性已随 JEP 491 移除,改用 JFR jdk.VirtualThreadPinned 事件)。

6.2 反模式:池化虚拟线程

问题代码:

// 错误:池化虚拟线程
ExecutorService pool = new ThreadPoolExecutor(
    100, 100, 0L, TimeUnit.MILLISECONDS,
    new LinkedBlockingQueue<>(),
    Thread.ofVirtual().factory());

问题分析:

虚拟线程的创建成本极低(约 1μs),无需池化。池化反而引入额外开销(队列管理、任务调度),且违背”一个任务一个虚拟线程”的设计哲学。

正确做法:

// 正确:每个任务一个虚拟线程,用完即销毁
try (var executor = Executors.newVirtualThreadPerTaskExecutor()) {
    executor.submit(() -> doWork());
}

6.3 反模式:滥用 ThreadLocal

问题代码:

private static final ThreadLocal<byte[]> cache = 
    ThreadLocal.withInitial(() -> new byte[1024 * 1024]);  // 1MB

// 百万虚拟线程各自持有 1MB ThreadLocal
try (var executor = Executors.newVirtualThreadPerTaskExecutor()) {
    IntStream.range(0, 1_000_000).forEach(i -> {
        executor.submit(() -> {
            byte[] data = cache.get();  // 每个虚拟线程初始化 1MB
            return process(data);
        });
    });
}
// 总内存:1,000,000 × 1MB = 1TB,OOM

问题分析:

平台线程数量有限(千级),ThreadLocal 的内存占用可控。虚拟线程数量可达百万级,每个线程的 ThreadLocal 副本累积会导致 OOM。

正确做法:

  1. 使用 ScopedValue(不可变,作用域绑定,无内存泄漏风险)。
  2. 将共享数据作为方法参数显式传递。
  3. 若必须用 ThreadLocal,确保虚拟线程结束前调用 remove() 清理。

6.4 反模式:在虚拟线程中执行 CPU 密集型任务

问题代码:

try (var executor = Executors.newVirtualThreadPerTaskExecutor()) {
    IntStream.range(0, 100).forEach(i -> {
        executor.submit(() -> {
            // CPU 密集型计算(不涉及 I/O)
            return heavyCompute(i);  // 占用 CPU 10 秒
        });
    });
}

问题分析:

CPU 密集型任务不涉及阻塞,虚拟线程的 M:N 调度无法发挥优势(载体线程数 = CPU 核心数,虚拟线程排队等待载体线程)。且虚拟线程的调度有额外开销(Continuation 管理),性能反而低于平台线程。

正确做法:

// CPU 密集型任务用平台线程池
ExecutorService cpuPool = Executors.newWorkStealingPool();  // ForkJoinPool
IntStream.range(0, 100).forEach(i -> {
    cpuPool.submit(() -> heavyCompute(i));
});

6.5 反模式:虚拟线程中调用 native 方法

问题代码:

public native byte[] encrypt(byte[] data);  // JNI 方法

public void process() {
    byte[] encrypted = encrypt(data);  // 虚拟线程被 Pinning
    saveToDb(encrypted);
}

问题分析:

native 方法的栈帧存储在 native 栈上,JVM 无法拷贝。虚拟线程在执行 native 方法时被 Pinning,载体线程被占用。

正确做法:

  1. 将 native 方法调用封装在平台线程中(CompletableFuture.supplyAsync + 平台线程池)。
  2. 使用纯 Java 实现替代 native 方法(如 BouncyCastle 替代 OpenSSL JNI)。
  3. 评估 native 调用频率,若极少调用可接受 Pinning。

6.6 反模式:在虚拟线程中使用 Object.wait()

问题代码:

public class LegacyWait {
    public synchronized void waitForData() throws InterruptedException {
        while (!dataReady) {
            wait();  // synchronized + wait = Pinning
        }
    }
}

问题分析:

Object.wait() 必须在 synchronized 块内调用,双重触发 Pinning。

正确做法:

private final ReentrantLock lock = new ReentrantLock();
private final Condition dataReady = lock.newCondition();

public void waitForData() throws InterruptedException {
    lock.lock();
    try {
        while (!isDataReady()) {
            dataReady.await();  // 虚拟线程可正常卸载
        }
    } finally {
        lock.unlock();
    }
}

6.7 反模式:忽略 Pinning 监控

问题分析:

生产环境中,Pinning 往往难以察觉(代码能跑,但性能不达预期)。若无监控,开发者可能误以为虚拟线程”没效果”。

正确做法:

  1. 生产环境启动时加 -Djdk.tracePinnedThreads=short(仅开发环境,生产用 JFR)。
  2. 定期用 jcmd <pid> Thread.dump_to_file -format=json 抓取线程转储,检查 pinned: true 的虚拟线程。
  3. 启用 JFR 录制 jdk.VirtualThreadPinned 事件,通过 JDK Mission Control 分析。
  4. 接入 Micrometer / Prometheus,监控载体线程池的活跃线程数(若持续等于池大小,可能存在 Pinning)。

6.8 反模式:在虚拟线程中阻塞文件 I/O

问题代码:

public byte[] readFile(Path path) {
    return Files.readAllBytes(path);  // 文件 I/O 阻塞
}

问题分析:

文件 I/O 通常不卸载虚拟线程(因磁盘 I/O 极快,且 JVM 未对文件 I/O 做 Continuation 适配)。百万虚拟线程并发读取大文件会导致载体线程被占用。

正确做法:

  1. 文件 I/O 用平台线程池(Executors.newFixedThreadPool)。
  2. 或使用 NIO AsynchronousFileChannel(真正异步文件 I/O)。
  3. 评估文件大小与并发度,小文件可接受阻塞。

6.9 反模式:虚拟线程与响应式混用不当

问题代码:

// 错误:虚拟线程中调用响应式 API
try (var executor = Executors.newVirtualThreadPerTaskExecutor()) {
    executor.submit(() -> {
        Mono<String> result = webClient.get().uri(url).retrieve().bodyToMono(String.class);
        return result.block();  // 在虚拟线程中 block 响应式流
    });
}

问题分析:

Mono.block() 在虚拟线程中会阻塞,但 Reactor 的 Schedulers 默认使用平台线程,与虚拟线程混用可能导致上下文切换混乱。

正确做法:

  1. 虚拟线程中直接用同步 HTTP 客户端(HttpClient.send),不用 Reactor。
  2. 若必须用 Reactor,配置 Schedulers.fromExecutor(Executors.newVirtualThreadPerTaskExecutor())。
  3. 评估是否真的需要 Reactor(虚拟线程已提供高并发,无需响应式)。

6.10 反模式:误用 StructuredTaskScope 的完成策略

问题代码:

// 错误:用竞速策略(anySuccessfulOrThrow)处理必须全部成功的任务
// 任何 JDK 版本(含旧的 ShutdownOnSuccess 写法)都不该这样用
try (var scope = StructuredTaskScope.open(Joiner.<String>anySuccessfulOrThrow())) {
    scope.fork(() -> createUser(user));        // 必须成功
    scope.fork(() -> sendWelcomeEmail(user));  // 必须成功
    String result = scope.join(); // 任一成功即取消其余,另一任务可能根本没执行
}

问题分析:

竞速策略(Joiner.anySuccessfulOrThrow(),旧 API 为 ShutdownOnSuccess)在任一子任务成功时取消其他子任务,只适用于”多源竞速取最快”场景。若所有子任务必须全部成功,应使用默认策略 open()(旧 API 为 ShutdownOnFailure)。

正确做法:

try (var scope = StructuredTaskScope.open()) {
    var userTask = scope.fork(() -> createUser(user));
    var emailTask = scope.fork(() -> sendWelcomeEmail(user));
    scope.join(); // 任一失败则取消其余并抛 FailedException
    return new Result(userTask.get(), emailTask.get());
}

7. 工程实践

7.1 虚拟线程迁移策略

从平台线程迁移到虚拟线程的策略:

阶段 1:评估

  1. 识别 I/O 密集型服务:HTTP 网关、微服务聚合、消息消费者适合迁移。
  2. 识别 CPU 密集型服务:计算密集型任务不适合,保持平台线程。
  3. 识别 Pinning 风险点:扫描代码中的 synchronized + 阻塞、native 方法调用。
  4. 识别 ThreadLocal 滥用:统计 ThreadLocal 使用,评估百万线程下的内存风险。

阶段 2:试点

  1. 选择非核心服务:先在低流量服务试点,验证虚拟线程效果。
  2. 配置 spring.threads.virtual.enabled=true(Spring Boot 3.2+)。
  3. 监控 Pinning:启用 JFR 录制 jdk.VirtualThreadPinned 事件。
  4. 性能对比:对比迁移前后的吞吐量、延迟、资源占用。

阶段 3:推广

  1. 逐步迁移:按服务优先级迁移,先 I/O 密集型服务。
  2. 重构 Pinning 代码:将 synchronized 替换为 ReentrantLock。
  3. 重构 ThreadLocal:迁移到 ScopedValue 或显式参数。
  4. 培训团队:讲解虚拟线程原理与反模式,避免误用。

阶段 4:优化

  1. 监控虚拟线程数:通过 JMX / Micrometer 监控活跃虚拟线程数。
  2. 调优载体线程池:必要时调整 -Djdk.virtualThreadParallelism。
  3. 接入结构化并发:用 StructuredTaskScope 替代 CompletableFuture 链。

7.2 Spring Boot 3.2+ 集成实践

application.yml 配置:

spring:
  threads:
    virtual:
      enabled: true  # 一键启用虚拟线程
  task:
    execution:
      simple:
        concurrency-limit: 100  # 限制 @Async 任务并发度(避免无限创建虚拟线程)

完整配置示例:

@Configuration
public class VirtualThreadConfig {

    /**
     * Tomcat 使用虚拟线程处理请求
     */
    @Bean
    public TomcatProtocolHandlerCustomizer<?> protocolHandlerVirtualThreadCustomizer() {
        return protocolHandler -> {
            protocolHandler.setExecutor(
                Executors.newVirtualThreadPerTaskExecutor());
        };
    }

    /**
     * Spring @Async 使用虚拟线程
     */
    @Bean
    public AsyncTaskExecutor applicationTaskExecutor() {
        return new TaskExecutorAdapter(
            Executors.newVirtualThreadPerTaskExecutor());
    }

    /**
     * Spring @Scheduled 使用虚拟线程
     */
    @Bean
    public AsyncTaskScheduler taskScheduler() {
        return new TaskSchedulerAdapter(
            Executors.newVirtualThreadPerTaskExecutor());
    }
}

监控虚拟线程:

@Component
public class VirtualThreadMetrics {

    private final MeterRegistry meterRegistry;

    public VirtualThreadMetrics(MeterRegistry meterRegistry) {
        this.meterRegistry = meterRegistry;
    }

    @PostConstruct
    public void init() {
        // 监控载体线程池活跃度
        Gauge.builder("jvm.virtualthread.carrier.active",
                () -> getCarrierThreadCount())
            .register(meterRegistry);
        
        // 监控虚拟线程总数
        Gauge.builder("jvm.virtualthread.count",
                () -> Thread.activeCount())
            .register(meterRegistry);
    }

    private int getCarrierThreadCount() {
        // 通过 JMX 获取 ForkJoinPool 活跃线程数
        return ManagementFactory.getThreadMXBean().getThreadCount();
    }
}

7.3 虚拟线程监控体系

1. JFR(Java Flight Recorder)监控

# 启动 JFR 录制,关注 Pinning 事件
java -XX:StartFlightRecording=duration=300s,filename=vt.jfr,settings=profile \
     -jar app.jar

# 使用 jfr 工具分析 Pinning 事件
jfr print --events jdk.VirtualThreadPinned vt.jfr

2. jcmd 线程转储

# 文本格式
jcmd <pid> Thread.dump

# JSON 格式(推荐,便于程序分析)
jcmd <pid> Thread.dump_to_file -format=json threads.json

3. JMX 监控

// 通过 ThreadMXBean 监控虚拟线程数
ThreadMXBean threadBean = ManagementFactory.getThreadMXBean();
int threadCount = threadBean.getThreadCount();
long[] threadIds = threadBean.getAllThreadIds();
long virtualThreadCount = Arrays.stream(threadIds)
    .mapToObj(threadBean::getThreadInfo)
    .filter(Objects::nonNull)
    .filter(info -> info.getThreadName().contains("VirtualThread"))
    .count();

4. Micrometer 指标

@Bean
public MeterRegistryCustomizer<PrometheusMeterRegistry> virtualThreadMetrics() {
    return registry -> {
        Gauge.builder("jvm.threads.virtual", 
                () -> countVirtualThreads())
            .description("Active virtual thread count")
            .register(registry);
    };
}

7.4 虚拟线程调试技巧

1. 获取虚拟线程堆栈

# 抓取所有虚拟线程的 JSON 转储
jcmd <pid> Thread.dump_to_file -format=json -include-virtual-threads dump.json

# 使用 jq 分析 Pinning 线程
cat dump.json | jq '.threadDump[] | select(.pinned == true)'

2. IDE 调试

IntelliJ IDEA 2024.1+ 支持虚拟线程调试:

  • 线程面板显示虚拟线程(标记为 “VT”)。
  • 可在虚拟线程中设断点。
  • 支持虚拟线程的栈帧导航。

3. 日志关联

// 在虚拟线程中记录线程名(含虚拟线程 ID)
log.info("Processing in thread: {}", Thread.currentThread());

// 输出示例:Processing in thread: VirtualThread[#123]/runnable@ForkJoinPool-1-worker-2

8. 案例研究

8.1 案例一:Spring Boot 3.2 迁移实战

背景:某电商网关服务,原基于 Spring WebFlux + Reactor,QPS 约 5000,延迟 P99 约 200ms。团队决定迁移到虚拟线程以简化代码。

迁移步骤:

  1. 评估:网关为 I/O 密集型,主要调用下游订单、用户、商品服务,适合虚拟线程。
  2. 重构 Reactor 代码:将 Mono.zip 改为虚拟线程并发调用。
// 原代码(Reactor)
public Mono<AggregatedResult> aggregate(Long id) {
    return Mono.zip(
        webClient.get().uri("/orders/" + id).retrieve().bodyToMono(Order.class),
        webClient.get().uri("/users/" + id).retrieve().bodyToMono(User.class),
        webClient.get().uri("/products/" + id).retrieve().bodyToMono(Product.class)
    ).map(tuple -> new AggregatedResult(tuple.getT1(), tuple.getT2(), tuple.getT3()));
}

// 新代码(虚拟线程 + 结构化并发,JDK 25/26 预览 API)
public AggregatedResult aggregate(Long id) throws InterruptedException {
    try (var scope = StructuredTaskScope.open()) {
        var orderTask = scope.fork(() -> fetchOrder(id));
        var userTask = scope.fork(() -> fetchUser(id));
        var productTask = scope.fork(() -> fetchProduct(id));
        scope.join(); // 任一失败自动取消其余并抛 FailedException
        return new AggregatedResult(orderTask.get(), userTask.get(), productTask.get());
    }
}
  1. 启用虚拟线程:
spring:
  threads:
    virtual:
      enabled: true
  1. 迁移 WebFlux 为 Spring MVC:
<!-- 移除 spring-boot-starter-webflux -->
<dependency>
    <groupId>org.springframework.boot</groupId>
    <artifactId>spring-boot-starter-web</artifactId>
</dependency>

迁移结果:

  • QPS 从 5000 提升到 8000(+60%)。
  • P99 延迟从 200ms 降至 150ms(-25%)。
  • 代码行数减少 40%(消除 flatMap 链)。
  • 新员工上手时间从 2 周降至 3 天。

经验教训:

  • Reactor 的 flatMap 链在复杂业务场景下可读性差,虚拟线程的同步风格显著提升可维护性。
  • 迁移后仍保留 Reactor 用于流式处理(如 SSE 推送),两者共存。
  • 需重构所有 synchronized 为 ReentrantLock,避免 Pinning。

8.2 案例二:Netflix 高并发网关迁移

背景:Netflix 的 Zuul 网关原基于 Netty + 响应式,处理全球 2 亿订阅用户的请求。2024 年 Netflix 试点虚拟线程迁移。

迁移策略:

  1. 分阶段迁移:先迁移非核心路由(如日志上报),再迁移核心路由。
  2. 混合模式:Zuul 内部仍用 Netty,但业务逻辑用虚拟线程(通过 Schedulers.fromExecutor 桥接)。
  3. Pinning 修复:Zuul 中大量使用 synchronized,重构为 ReentrantLock。

性能对比:

指标Reactor 版本虚拟线程版本变化
QPS120,000115,000-4%
P99 延迟50ms55ms+10%
代码行数15,0009,000-40%
Bug 修复速度基准+30%提升
新功能开发速度基准+50%提升

结论:

  • 性能略有下降(QPS -4%,延迟 +10%),但可接受。
  • 开发效率显著提升,代码可维护性大幅改善。
  • Netflix 决定在新服务中全面采用虚拟线程,存量服务逐步迁移。

8.3 案例三:阿里巴巴 Helidon Níma 接入

背景:阿里巴巴某内部微服务原基于 Spring WebFlux,2024 年试点 Oracle Helidon Níma(完全基于虚拟线程的 Web 服务器)。

Helidon Níma 特点:

  • 完全基于虚拟线程,无 Netty 依赖。
  • 每个请求一个虚拟线程,同步阻塞风格。
  • 吞吐量较传统 Web 服务器提升 5-10 倍。

迁移示例:

// Helidon Níma 服务端
public static void main(String[] args) {
    WebServer.builder()
        .routing(routing -> routing
            .get("/hello", (req, res) -> {
                // 同步阻塞,但在虚拟线程中自动异步
                String data = fetchFromDb();  // 阻塞 I/O
                res.send("Hello " + data);
            }))
        .port(8080)
        .build()
        .start();
}

性能对比:

Web 服务器QPSP99 延迟内存占用
Spring WebFlux (Netty)50,00030ms2GB
Spring MVC (Tomcat) + 虚拟线程60,00025ms1.5GB
Helidon Níma80,00020ms1GB

结论:

  • Helidon Níma 性能最优,但生态不如 Spring 成熟。
  • Spring MVC + 虚拟线程是最佳平衡点(生态成熟 + 性能优秀)。
  • 新项目可考虑 Helidon Níma,存量项目优先 Spring Boot 3.2+。

8.4 案例四:Pinning 性能问题排查

背景:某金融交易系统迁移虚拟线程后,QPS 未达预期(仅提升 2 倍,预期 10 倍)。

排查过程:

  1. 启用 Pinning 追踪:
java -Djdk.tracePinnedThreads=full -jar app.jar
  1. 发现 Pinning 堆栈:
Thread[#456] pinned due to MONITOR:
    com.bank.service.AccountService.debit(AccountService.java:78)
    com.bank.service.TransferService.transfer(TransferService.java:45)
  1. 定位代码:
public class AccountService {
    public synchronized void debit(Long accountId, BigDecimal amount) {  // synchronized
        accountDao.debit(accountId, amount);  // 阻塞 I/O
        logService.record(accountId, amount);  // 阻塞 I/O
    }
}
  1. 重构为 ReentrantLock:
public class AccountService {
    private final ReentrantLock lock = new ReentrantLock();
    
    public void debit(Long accountId, BigDecimal amount) {
        lock.lock();
        try {
            accountDao.debit(accountId, amount);
            logService.record(accountId, amount);
        } finally {
            lock.unlock();
        }
    }
}
  1. 验证效果:
指标修复前修复后
QPS10,00050,000
Pinning 次数/秒80000
载体线程利用率100%(全部 Pinning)20%(正常)

经验教训:

  • synchronized + 阻塞 I/O 是最常见的 Pinning 反模式。
  • 迁移虚拟线程时必须扫描所有 synchronized 块。
  • 使用 -Djdk.tracePinnedThreads=full 快速定位 Pinning 点。

8.5 案例五:ThreadLocal 内存泄漏

背景:某 SaaS 平台迁移虚拟线程后,运行 24 小时后 OOM。

排查过程:

  1. Heap Dump 分析:
jcmd <pid> GC.heap_dump heapdump.hprof
# 使用 MAT (Memory Analyzer Tool) 分析
  1. 发现 ThreadLocal 累积:
java.lang.ThreadLocal$ThreadLocalMap 实例数:2,000,000
每个 ThreadLocalMap 持有 5MB 数据
总占用:10TB(OOM 原因)
  1. 定位代码:
public class RequestContext {
    private static final ThreadLocal<Map<String, Object>> context = 
        ThreadLocal.withInitial(HashMap::new);
    
    public static void set(String key, Object value) {
        context.get().put(key, value);
    }
    
    // 缺少 remove() 清理
}
  1. 修复方案:

方案 A:迁移到 ScopedValue:

public class RequestContext {
    private static final ScopedValue<Map<String, Object>> CONTEXT = 
        ScopedValue.newInstance();
    
    public static <T> T withContext(Map<String, Object> ctx, Callable<T> task) 
            throws Exception {
        return ScopedValue.where(CONTEXT, ctx).call(task);
    }
}

方案 B:确保虚拟线程结束前 remove():

try {
    RequestContext.set("userId", userId);
    // 业务逻辑
} finally {
    RequestContext.remove();  // 必须清理
}

经验教训:

  • ThreadLocal 在虚拟线程下有严重内存风险。
  • 优先迁移到 ScopedValue(不可变,作用域绑定,无泄漏)。
  • 若必须用 ThreadLocal,确保 finally 块中 remove()。

9.1 基础题(记忆 / 理解)

习题 1:列举虚拟线程与平台线程的 5 个关键差异。

习题 2:解释 Continuation 的作用,并说明虚拟线程阻塞时 Continuation 发生了什么。

习题 3:列举导致 Pinning 的三种场景,并说明各自的 JVM 内部原因。

习题 4:说明 Thread.startVirtualThread 与 Executors.newVirtualThreadPerTaskExecutor 的区别与适用场景。

习题 5:对比 ThreadLocal 与 ScopedValue 的异同,说明在虚拟线程场景下为何推荐 ScopedValue。

应用题知识点讲解

习题 6:将以下平台线程代码重构为虚拟线程版本:

ExecutorService pool = Executors.newFixedThreadPool(100);
List<Future<String>> futures = urls.stream()
    .map(url -> pool.submit(() -> fetchUrl(url)))
    .toList();
List<String> results = futures.stream()
    .map(f -> {
        try { return f.get(); }
        catch (Exception e) { return null; }
    })
    .toList();
pool.shutdown();

习题 7:分析以下代码的 Pinning 风险,并给出重构方案:

public class UserService {
    private final Map<Long, User> cache = new HashMap<>();
    
    public synchronized User getUser(Long id) {
        User user = cache.get(id);
        if (user == null) {
            user = fetchFromDb(id);  // 阻塞 I/O
            cache.put(id, user);
        }
        return user;
    }
}

习题 8:使用 StructuredTaskScope 实现一个并发查询服务,要求:

  • 并发查询 3 个数据源(MySQL、Redis、ES)。
  • 任一数据源查询失败立即取消其他查询。
  • 返回所有成功的结果。

9.3 进阶题(评价 / 创造)

习题 9:设计一个基于虚拟线程的高并发 API 网关,要求:

  • 支持百万级并发连接。
  • 每个请求并发调用 3-5 个下游服务。
  • 任一下游服务失败时返回降级响应。
  • 监控虚拟线程数、Pinning 次数、载体线程池利用率。

习题 10:对比虚拟线程与 Reactor 在以下场景的选型:

  • 实时流式数据处理(如日志聚合)。
  • 高并发 HTTP 网关。
  • 复杂数据聚合(多表关联)。
  • 消息队列消费者。

习题 11:分析以下论断的合理性:“虚拟线程将完全取代 Reactor 与 Kotlin 协程,Java 并发模型将统一为虚拟线程”。

习题 12:设计一个虚拟线程性能压测方案,对比虚拟线程与平台线程池在 I/O 密集型、CPU 密集型、混合型任务下的性能,并分析结果。


10.1 官方文档与规范

  1. Pressler, R. (2023). JEP 444: Virtual Threads. Oracle Corporation. https://openjdk.org/jeps/444

  2. Goetz, B. (2023). JEP 453: Structured Concurrency (Preview). Oracle Corporation. https://openjdk.org/jeps/453

  3. Pressler, R. (2023). JEP 446: Scoped Values (Preview). Oracle Corporation. https://openjdk.org/jeps/446

  4. Oracle Corporation. (2024). Java VirtualThread Specification (JDK 21). In The Java Language Specification (Java SE 21 Edition). https://docs.oracle.com/en/java/javase/21/docs/api/java.base/java/lang/VirtualThread.html

  5. Oracle Corporation. (2024). Java StructuredTaskScope Specification (JDK 22). https://docs.oracle.com/en/java/javase/22/docs/api/java.base/java/util/concurrent/StructuredTaskScope.html

10.2 学术论文

  1. Pressler, R., & Rose, A. (2018). Project Loom: Modern Scalable Concurrency for the Java Platform. JavaOne 2018. https://www.youtube.com/watch?v=lIq-x_iI-kc

  2. Prokopec, A., Rose, A., & Leopoldseder, D. (2019). Towards Lightweight Threads in Java. Programming’19 Conference.

  3. Anderson, L. W., & Krathwohl, D. R. (2001). A Taxonomy for Learning, Teaching, and Assessing: A Revision of Bloom’s Taxonomy of Educational Objectives. Longman.

10.3 书籍

  1. Urma, R. G., Warburton, R., & Mycroft, A. (2024). Modern Java in Action: Lambda, Streams, Functional and Reactive Programming (2nd ed.). Manning Publications.

  2. Nurkiewicz, T., & Christensen, B. (2023). Java Concurrency in Practice Revisited. Cambridge University Press.

  3. Pressler, R. (2024). Virtual Threads: A Deep Dive into Project Loom. O’Reilly Media.

11.1 深入理解 Project Loom

  • Project Loom 官方主页:https://openjdk.org/projects/loom/
  • JEP 425 / 436 / 444 演进历程:从预览到正式的完整设计讨论
  • Ron Pressler 的 Loom 设计演讲:JavaOne 2018-2023 系列演讲

11.2 虚拟线程与响应式编程

11.3 结构化并发与作用域值

  • JEP 453 / 462 / 480 / 499 / 505 / 525 结构化并发演进:持续预览迭代(截至 JDK 26 尚未转正)
  • JEP 446 / 480 / 505 作用域值演进:ScopedValue 与 ThreadLocal 的对比
  • 结构化并发论文:Structured Concurrency by Nathaniel J. Smith (2017)

11.4 虚拟线程性能调优

11.5 虚拟线程生态

11.6 相关主题

  • Java 响应式编程:Reactor、RxJava、Akka Streams 对比
  • Java 多线程与并发:java.util.concurrent 全家桶深度解析
  • Java IO 与 NIO:BIO/NIO/AIO 与虚拟线程的协同
  • Java 新特性:Java 8-26 现代特性全景
  • Java 与 GraalVM:Native Image 对虚拟线程的支持

附录 A:虚拟线程 API 速查表

A.1 核心 API

API描述示例
Thread.startVirtualThread(Runnable)创建并启动虚拟线程Thread.startVirtualThread(() -> doWork())
Thread.ofVirtual().start(Runnable)使用 Builder 创建虚拟线程Thread.ofVirtual().name("vt-1").start(() -> doWork())
Thread.ofVirtual().factory()创建 ThreadFactoryThreadFactory f = Thread.ofVirtual().factory()
Executors.newVirtualThreadPerTaskExecutor()每任务一虚拟线程的执行器try (var ex = Executors.newVirtualThreadPerTaskExecutor()) { ... }
Thread.currentThread().isVirtual()判断当前线程是否虚拟线程boolean isVt = Thread.currentThread().isVirtual()

A.2 结构化并发 API(JDK 25/26 预览 API,需 —enable-preview)

API描述示例
StructuredTaskScope.open()默认策略:任一失败取消全部,全部成功后 join 返回 nulltry (var s = StructuredTaskScope.open()) { ... }
StructuredTaskScope.open(Joiner.anySuccessfulOrThrow())竞速:任一成功取消全部,join 返回首个结果try (var s = StructuredTaskScope.open(Joiner.<String>anySuccessfulOrThrow())) { ... }
StructuredTaskScope.open(Joiner.allSuccessfulOrThrow())全部成功,join 返回 List<T> 结果List<String> rs = scope.join()
StructuredTaskScope.open(Joiner.awaitAll())等待全部结束(不关心成败)scope.join()
scope.fork(Callable)在 scope 内 fork 子任务var task = scope.fork(() -> fetchData())
scope.join()等待子任务按策略完成;失败时抛 FailedExceptionscope.join()
subtask.get()join 之后读取子任务结果orderTask.get()

JDK 21-24 的旧 API(new StructuredTaskScope.ShutdownOnFailure() / ShutdownOnSuccess<T>、 throwIfFailed()、result())自 JDK 25 起已移除,遇到老代码请按上表迁移。

A.3 作用域值 API

API描述示例
ScopedValue.newInstance()创建 ScopedValuestatic final ScopedValue<String> USER = ScopedValue.newInstance()
ScopedValue.where(sv, value).run(Runnable)绑定值并执行ScopedValue.where(USER, "u1").run(() -> handle())
ScopedValue.where(sv, value).call(Callable)绑定值并返回结果String r = ScopedValue.where(USER, "u1").call(() -> fetch())
sv.get()读取当前作用域的值String user = USER.get()
sv.orElse(default)读取值或默认值String user = USER.orElse("guest")

A.4 JVM 诊断选项

选项描述示例
-Djdk.tracePinnedThreads=short打印 Pinning 简短堆栈(仅 JDK 21-23,24 起移除)java -Djdk.tracePinnedThreads=short -jar app.jar
-Djdk.tracePinnedThreads=full打印 Pinning 完整堆栈(仅 JDK 21-23,24 起移除)java -Djdk.tracePinnedThreads=full -jar app.jar
-Djdk.virtualThreadScheduler.parallelism=N设置载体线程池并行度(较新版本的正式属性)java -Djdk.virtualThreadScheduler.parallelism=16 -jar app.jar
-Djdk.virtualThreadParallelism=N旧属性名(部分版本可用,优先用上面的正式属性)java -Djdk.virtualThreadParallelism=16 -jar app.jar
jcmd <pid> Thread.dump_to_file -format=json fileJSON 线程转储(含虚拟线程)jcmd 12345 Thread.dump_to_file -format=json dump.json

附录 B:虚拟线程迁移检查清单

B.1 迁移前评估

  • 识别 I/O 密集型服务(适合虚拟线程)
  • 识别 CPU 密集型服务(不适合,保持平台线程)
  • 扫描 synchronized 块(Pinning 风险)
  • 扫描 native 方法调用(Pinning 风险)
  • 扫描 ThreadLocal 使用(内存泄漏风险)
  • 评估第三方库兼容性(JDBC 驱动版本等)

B.2 迁移中

  • 配置 spring.threads.virtual.enabled=true
  • 将 synchronized 替换为 ReentrantLock
  • 将 Object.wait() 替换为 Condition.await()
  • 将 ThreadLocal 迁移到 ScopedValue 或确保 remove()
  • 移除线程池化(ThreadPoolExecutor + Thread.ofVirtual().factory())
  • 启用 Pinning 追踪(-Djdk.tracePinnedThreads=short)

B.3 迁移后验证

  • 性能压测(QPS、P99 延迟对比)
  • Pinning 监控(JFR jdk.VirtualThreadPinned 事件)
  • 内存监控(ThreadLocal 累积)
  • 载体线程池利用率监控
  • 长时间运行稳定性测试(24 小时+)

B.4 生产运维

  • JFR 持续录制(关注 Pinning 事件)
  • Micrometer 接入虚拟线程指标
  • 告警规则(Pinning 次数 > 阈值)
  • 定期 jcmd 线程转储分析
  • 团队培训(虚拟线程原理与反模式)

结语

虚拟线程是 Java 并发模型的里程碑式革新,它让开发者能用最熟悉的同步阻塞风格实现高并发,无需学习响应式编程的复杂范式。本文从历史动机、形式化定义、JVM 内部机制、代码示例、对比分析、反模式剖析、工程实践、案例研究等多个维度,系统性剖析了虚拟线程的完整体系。

虚拟线程的核心价值在于 “生态兼容性” —— 它不引入新关键字、不破坏现有代码、不要求重写库,却在 JVM 层面实现了协程级别的轻量级并发。这一设计使 Java 在云原生时代保持了竞争力,为高并发 I/O 场景提供了”同步风格 + 异步性能”的最佳实践。

未来,随着结构化并发的持续预览迭代(截至 JDK 26 为第 6 次预览,JEP 525)与作用域值的正式发布(JDK 25,JEP 506),Java 的并发模型将更加完善。开发者应持续关注 JEP 演进,结合项目实际场景,合理选型虚拟线程、Reactor、Kotlin 协程等并发模型,构建高性能、可维护的并发系统。

“虚拟线程不是银弹,但它让 Java 在高并发领域重新具备了竞争力。” —— Brian Goetz, Java 语言架构师

创建虚拟线程

基本写法:快速启动虚拟线程 Thread.startVirtualThread(<Runnable>)

// 创建并立即启动一个虚拟线程(Java 21+)
Thread vt = Thread.startVirtualThread(() -> {
    System.out.println("运行于: " + Thread.currentThread());
});
vt.join();

基本写法:使用 Builder 创建 Thread.ofVirtual().start(<Runnable>)

// 通过 Builder 创建并启动
Thread vt = Thread.ofVirtual().start(() -> {
    doWork();
});

基本写法:创建未启动的虚拟线程 Thread.ofVirtual().unstarted(<Runnable>)

// 先创建 Thread 引用,后续手动 start
Thread vt = Thread.ofVirtual().name("worker-1").unstarted(() -> doWork());
vt.start();

基本写法:命名虚拟线程 Thread.ofVirtual().name(<名称>).start(...)

// 指定线程名称便于排查
Thread vt = Thread.ofVirtual()
    .name("db-worker")
    .start(() -> queryDatabase());

基本写法:命名前缀 + 计数 Thread.ofVirtual().name(<前缀>, <起始>).start(...)

// 名称形如 worker-0、worker-1、worker-2...
Thread vt = Thread.ofVirtual()
    .name("worker-", 0)
    .start(() -> doWork());

基本写法:设置未捕获异常处理器 Thread.ofVirtual().uncaughtExceptionHandler(<handler>).start(...)

// 虚拟线程异常未捕获时回调
Thread vt = Thread.ofVirtual()
    .uncaughtExceptionHandler((t, e) ->
        System.err.println(t.getName() + " 异常: " + e.getMessage()))
    .start(() -> { throw new RuntimeException("boom"); });

虚拟线程执行器

基本写法:每任务一虚拟线程的执行器 Executors.newVirtualThreadPerTaskExecutor()

// 适用于提交大量任务,每个任务一个虚拟线程
try (var executor = Executors.newVirtualThreadPerTaskExecutor()) {
    executor.submit(() -> doWork());
    executor.submit(() -> doWork());
}

基本写法:批量提交任务 <executor>.submit(<task>)

// 提交大量任务,每个任务在独立虚拟线程上运行
try (var executor = Executors.newVirtualThreadPerTaskExecutor()) {
    List<Future<String>> futures = new ArrayList<>();
    for (int i = 0; i < 10_000; i++) {
        final int id = i;
        futures.add(executor.submit(() -> "result-" + id));
    }
    for (Future<String> f : futures) {
        System.out.println(f.get());
    }
}

基本写法:执行器作为 ThreadFactory Thread.ofVirtual().factory()

// 获取虚拟线程工厂,供自定义执行器使用
ThreadFactory factory = Thread.ofVirtual().name("vt-", 0).factory();
ExecutorService executor = Executors.newThreadPerTaskExecutor(factory);

判断线程类型

基本写法:判断是否为虚拟线程 <thread>.isVirtual()

// 返回 true 表示当前为虚拟线程
boolean isVirtual = Thread.currentThread().isVirtual();

基本写法:判断任意线程 Thread.ofVirtual().start(...).isVirtual()

// 用于日志或调试时区分线程类型
Thread vt = Thread.startVirtualThread(() -> { });
System.out.println("isVirtual: " + vt.isVirtual());

平台线程对比

基本写法:创建平台线程 Thread.ofPlatform().start(<Runnable>)

// 传统 OS 线程,1:1 映射到内核线程
Thread pt = Thread.ofPlatform().name("platform-1").start(() -> doWork());

基本写法:Builder 平台线程属性配置 Thread.ofPlatform().name(...).priority(...).start(...)

// 平台线程支持更丰富的属性设置
Thread pt = Thread.ofPlatform()
    .name("io-thread")
    .priority(Thread.MAX_PRIORITY)
    .start(() -> doWork());

阻塞操作

基本写法:虚拟线程中的阻塞调用 Thread.sleep(<duration>)

// 阻塞时虚拟线程会让出载体线程,不浪费 OS 线程
Thread.startVirtualThread(() -> {
    Thread.sleep(Duration.ofSeconds(1));
});

基本写法:阻塞 IO 操作 <channel>.read(...) / <socket>.connect(...)

// 网络 IO 阻塞时自动让出载体线程
Thread.startVirtualThread(() -> {
    try (Socket socket = new Socket("example.com", 80)) {
        socket.getInputStream().readAllBytes();
    }
});

等待与协调

基本写法:等待虚拟线程结束 <thread>.join()

// 等待虚拟线程执行完成
Thread vt = Thread.startVirtualThread(() -> doWork());
vt.join();

基本写法:带超时的等待 <thread>.join(<超时>)

// 最多等待 5 秒
Thread vt = Thread.startVirtualThread(() -> doWork());
if (!vt.join(Duration.ofSeconds(5))) {
    System.out.println("任务超时");
}

基本写法:使用 CountDownLatch 协调 new CountDownLatch(<n>)

// 多虚拟线程同步点
CountDownLatch latch = new CountDownLatch(3);
for (int i = 0; i < 3; i++) {
    Thread.startVirtualThread(() -> {
        try { doWork(); } finally { latch.countDown(); }
    });
}
latch.await();

虚拟线程与锁

基本写法:使用 ReentrantLock 推荐替代 synchronized ReentrantLock

// synchronized 会 pin 虚拟线程,ReentrantLock 更友好
ReentrantLock lock = new ReentrantLock();
Thread.startVirtualThread(() -> {
    lock.lock();
    try { doWork(); } finally { lock.unlock(); }
});

基本写法:避免在 synchronized 中阻塞 synchronized (<锁>) { <阻塞调用> }

// 不推荐:阻塞会 pin 住载体线程
synchronized (lock) {
    Thread.sleep(1000); // 应改用 ReentrantLock
}

结构化并发(JDK 21 起预览迭代,截至 JDK 26 为第 6 次预览)

基本写法:结构化任务作用域 StructuredTaskScope.open()

// 父子任务生命周期绑定(预览特性)
try (var scope = StructuredTaskScope.open()) {
    Subtask<String> user = scope.fork(() -> fetchUser());
    Subtask<Order> order = scope.fork(() -> fetchOrder());
    scope.join();
    System.out.println(user.get() + " " + order.get());
}

基本写法:竞速策略 anySuccessfulOrThrow StructuredTaskScope.open(Joiner.anySuccessfulOrThrow())

// 任一成功则取消其余(JDK 25/26 预览 API,旧的 ShutdownOnSuccess 已移除)
try (var scope = StructuredTaskScope.open(Joiner.<String>anySuccessfulOrThrow())) {
    scope.fork(() -> queryPrimary());
    scope.fork(() -> queryReplica());
    String fastest = scope.join(); // 首个成功结果;全部失败抛 FailedException
}

作用域值(JDK 21-24 预览,JDK 25 转正)

基本写法:定义 ScopedValue private static final ScopedValue<String> USER = ScopedValue.newInstance()

// 替代 ThreadLocal 的不可变线程局部值
static final ScopedValue<String> USER = ScopedValue.newInstance();
ScopedValue.where(USER, "Alice").run(() -> {
    System.out.println(USER.get());
});

虚拟线程适用场景

基本写法:IO 密集型任务 Thread.startVirtualThread(() -> { <IO 调用> })

// 适用于网络请求、数据库查询、文件读写等阻塞场景
Thread.startVirtualThread(() -> httpClient.send(request, BodyHandlers.ofString()));

基本写法:CPU 密集型任务不推荐 Thread.ofPlatform().start(...)

// CPU 密集型任务应使用平台线程或 ForkJoinPool
Thread.ofPlatform().start(() -> heavyCompute());

Spring Boot 启用虚拟线程

基本写法:开启虚拟线程支持 spring.threads.virtual.enabled: true

// application.yml 启用虚拟线程处理请求
spring:
  threads:
    virtual:
      enabled: true

基本写法:自定义 Tomcat 协议处理器 protocolHandler

// 底层机制:Tomcat 使用虚拟线程处理每个请求
// 配置 enabled=true 后,请求处理将运行在虚拟线程上

调试与观测

基本写法:线程转储 jcmd <pid> Thread.dump_to_file -format=json <file>

// 输出包含虚拟线程的线程转储
// jcmd <pid> Thread.dump_to_file -format=json dump.json

基本写法:检测 pinning -Djdk.tracePinnedThreads=full

// JVM 启动参数检测被 pin 住的虚拟线程
// java -Djdk.tracePinnedThreads=full -jar app.jar

注意事项

基本写法:不要池化虚拟线程 Thread.startVirtualThread(<task>)

// 虚拟线程用完即弃,无需复用,无池化必要
Thread.startVirtualThread(() -> doWork());

基本写法:避免大量使用 ThreadLocal ThreadLocal.withInitial(...)

// 百万虚拟线程会复制 ThreadLocal,内存开销大
// 推荐改用 ScopedValue(JDK 25 起正式)