Java 与虚拟线程
Project Loom 虚拟线程、结构化并发、Continuation 机制与性能调优全景式深度解析
前置知识
- Java 阻塞队列 BlockingQueue 语法速查手册:建议先完成前一篇的学习
学习目标
- 掌握「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类的子类型,所有现有ThreadAPI(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 版本 | 虚拟线程相关演进 |
|---|---|---|
| 2018 | Project Loom 启动 | Ron Pressler 提出 Loom 项目,目标是”轻量级线程 + 同步风格” |
| 2021-05 | JDK 16 | Loom 早期原型进入沙箱,java.lang.Fiber 实验性 API |
| 2022-09 | JDK 19 | JEP 425:虚拟线程预览(Preview),API 为 Thread.ofVirtual() |
| 2023-03 | JDK 20 | JEP 436:虚拟线程第二次预览,Thread.startVirtualThread 简化 API |
| 2023-09 | JDK 21 (LTS) | JEP 444:虚拟线程正式发布,成为生产可用特性 |
| 2023-11 | Spring Boot 3.2 | spring.threads.virtual.enabled=true 一键启用虚拟线程 |
| 2024-03 | JDK 22 | JEP 462:结构化并发第二次预览(StructuredTaskScope) |
| 2024-09 | JDK 23 | JEP 480:结构化并发第三次预览;JEP 481:作用域值第三次预览 |
| 2025-03 | JDK 24 | JEP 491:synchronized 阻塞不再固定载体线程(Monitor 类 Pinning 基本消除);JEP 499:结构化并发第四次预览 |
| 2025-09 | JDK 25 (LTS) | JEP 506:作用域值转正;JEP 505:结构化并发第五次预览(API 重设计为 open() + Joiner) |
| 2026-03 | JDK 26 | JEP 525:结构化并发第六次预览,仍未转正 |
1.3 设计哲学:为何选择”虚拟线程”而非”协程”
Java 设计团队(Brian Goetz、Ron Pressler)在 Loom 设计阶段曾考虑三种方案:
-
方案 A:语言级
async/await关键字(C#、Rust、Python 风格)- 优点:编译期静态检查异步调用,性能最优。
- 缺点:引入”函数染色”问题(async 函数只能被 async 函数调用),现有同步库需重写为 async 版本,破坏 Java 生态兼容性。
-
方案 B:协程 +
suspend关键字(Kotlin 风格)- 优点:无函数染色问题,语法简洁。
- 缺点:JVM 上实现
suspend需要修改字节码生成 CPS(Continuation-Passing Style)变换,与现有字节码工具(ASM、CGLIB)不兼容;Java 与 Kotlin 互操作时suspend函数在 Java 侧需手动处理Continuation参数。
-
方案 C:虚拟线程(最终选择)
- 优点:
- 无新关键字:
Thread、Runnable、CallableAPI 完全复用。 - 无函数染色:同步代码在虚拟线程中自动获得异步性能。
- 生态兼容:JDBC 4.0+、HttpClient、Socket、NIO 阻塞操作自动卸载。
- 无新关键字:
- 缺点:
- JVM 实现复杂:需修改 HotSpot 的
Thread类、栈管理、Continuation实现。 - Pinning 陷阱:
synchronized块内阻塞会导致载体线程被占用,需开发者主动用ReentrantLock替换。 - ThreadLocal 内存陷阱:百万级虚拟线程各自持有 ThreadLocal 副本可能导致 OOM。
- JVM 实现复杂:需修改 HotSpot 的
- 优点:
Java 团队选择方案 C 的核心原因是 “生态兼容性优先” —— Java 生态积累了 25 年的同步库(JDBC、JPA、Jackson、OkHttp),重写这些库为 async 版本的成本远高于在 JVM 层实现虚拟线程。这一决策使 Java 在云原生时代保持了”一次编写、到处运行”的承诺。
1.4 JEP 与虚拟线程相关提案
| JEP 编号 | 标题 | 落地版本与状态 | 核心内容 |
|---|---|---|---|
| JEP 425 | Virtual Threads (Preview) | JDK 19,预览 | 虚拟线程首次预览,Thread.ofVirtual() API |
| JEP 436 | Virtual Threads (Second Preview) | JDK 20,预览 | API 微调,Thread.startVirtualThread 简化 |
| JEP 444 | Virtual Threads | JDK 21,正式 | 虚拟线程正式发布,API 稳定 |
| JEP 453 | Structured Concurrency (Preview) | JDK 21,预览 | StructuredTaskScope 首次预览 |
| JEP 462 | Structured Concurrency (Second Preview) | JDK 22,预览 | API 简化,ShutdownOnFailure / ShutdownOnSuccess |
| JEP 463 | Implicitly Declared Classes and Instance Main Methods | JDK 22,预览 | 与虚拟线程无关,但同期发布(JDK 25 转正为 JEP 512) |
| JEP 480 | Structured Concurrency (Third Preview) | JDK 23,预览 | API 进一步稳定 |
| JEP 481 | Scoped Values (Third Preview) | JDK 23,预览 | ScopedValue 替代 ThreadLocal |
| JEP 499 | Structured Concurrency (Fourth Preview) | JDK 24,预览 | API 趋于稳定 |
| JEP 491 | Synchronized Virtual Threads without Pinning | JDK 24,正式 | synchronized/Object.wait() 不再固定载体线程 |
| JEP 505 | Structured Concurrency (Fifth Preview) | JDK 25,预览 | API 重设计:StructuredTaskScope.open() + Joiner 完成策略 |
| JEP 506 | Scoped Values | JDK 25,正式 | 作用域值转正,无需 --enable-preview |
| JEP 525 | Structured 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 虚拟线程的形式化定义
设 为平台线程集合(操作系统线程), 为虚拟线程集合, 为载体线程(Carrier Thread)集合, 为调度器(Scheduler)。虚拟线程系统可形式化为一个六元组:
其中:
- :虚拟线程集合, 可达 量级。
- :平台线程集合, 通常等于 CPU 核心数。
- :载体线程,由
ForkJoinPool提供,。 - :调度器, 将虚拟线程映射到载体线程。
- :延续对象,保存虚拟线程的栈帧。
- :挂载操作,将虚拟线程绑定到载体线程执行。
2.2 M:N 调度模型
传统平台线程采用 1:1 调度:一个 Java 线程对应一个操作系统线程。
虚拟线程采用 M:N 调度:M 个虚拟线程映射到 N 个载体线程。
其中 (典型值 ,)。调度器 负责在载体线程上分时执行虚拟线程,虚拟线程阻塞时自动让出载体线程。
2.3 虚拟线程状态机
虚拟线程的状态机可形式化为:
状态转移规则:
关键区别:
- 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 是虚拟线程的核心抽象,表示”计算的剩余部分”。形式化地,设 为一个 Continuation, 为当前栈帧序列, 为执行环境(局部变量、操作数栈),则:
虚拟线程的执行可建模为 Continuation 的挂载与卸载:
Continuation 的栈帧存储在堆上(而非操作系统栈),这是虚拟线程内存占用极低的根本原因。每个 Continuation 的初始栈约 1KB,可按需增长到 100KB-1MB(极少见)。
2.5 内存占用形式化对比
设 为线程数量, 为平台线程总内存, 为虚拟线程总内存:
对于 个线程:
虚拟线程的内存优势为 1000 倍,这是其能支撑百万级并发的物理基础。
2.6 吞吐量模型
设 为任务 CPU 执行时间, 为任务 I/O 等待时间, 为载体线程数。对于 I/O 密集型任务():
平台线程模型:线程数 受限于内存,。吞吐量:
虚拟线程模型:虚拟线程数 可达 ,但受限于载体线程数 。吞吐量:
因为虚拟线程在 I/O 等待时释放载体线程,载体线程始终在执行 CPU 工作。对于 ,,:
虚拟线程的吞吐量优势为 200 倍(在 I/O 密集型场景)。
3. 理论推导:JVM 视角下的虚拟线程机制
3.1 虚拟线程的 JVM 实现
虚拟线程在 HotSpot JVM 中的实现涉及以下核心组件:
java.lang.VirtualThread(JDK 21+):继承自Thread,是虚拟线程的 Java 层表示。jdk.internal.vm.Continuation:JVM 内部类,封装 Continuation 机制。ForkJoinPool:虚拟线程的载体线程池,默认Runtime.getRuntime().availableProcessors()个工作线程。VirtualThread的run()方法:在载体线程上执行 Continuation。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):
- java.net.Socket I/O:
Socket.getInputStream().read()、Socket.getOutputStream().write()阻塞时。 - java.nio.channels:
SocketChannel.read()、ServerSocketChannel.accept()阻塞时(NIO 阻塞模式)。 - java.net.http.HttpClient:
HttpClient.send()同步方法阻塞时。 - java.io.FileInputStream / FileOutputStream:阻塞 I/O(注意:文件 I/O 通常不卸载,因磁盘 I/O 极快)。
- Thread.sleep():sleep 时自动卸载。
- Object.wait() / Condition.await():等待时自动卸载。
- LockSupport.park():显式 park 时卸载。
- BlockingQueue 操作:
put()、take()阻塞时卸载。 - Semaphore.acquire():许可不足时卸载。
- CountDownLatch.await():等待时卸载。
- Future.get():等待结果时卸载。
- 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() 的实现机制:
- JVM 检测到阻塞操作(如
socket.read())。 - JVM 调用
Continuation.yield()。 yield()遍历当前栈帧,将每帧的局部变量、操作数栈、返回地址保存到Continuation对象(存储在堆上)。yield()抛出YieldException,用于栈展开(stack unwinding)。- 载体线程捕获
YieldException,释放虚拟线程,继续调度其他虚拟线程。
run() 的实现机制:
- 载体线程从调度器队列取出
Continuation对象。 - 调用
Continuation.run()。 - JVM 从
Continuation对象恢复栈帧到当前载体线程栈。 - 从
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/async | Thread.startVirtualThread |
| 函数染色 | 是(suspend 函数只能被 suspend 调用) | 否 |
| 与 Java 互操作 | 需 Continuation 参数适配 | 完全兼容(无新关键字) |
| 性能 | 优秀(编译期 CPS,无运行时开销) | 优秀(JVM 优化,接近协程) |
| 调试 | 良好(Kotlin 协程调试器) | 良好(jcmd 线程转储) |
| 生态 | Kotlin 原生 | Java 原生(Spring、Quarkus 等) |
| 结构化并发 | 原生支持(coroutineScope) | StructuredTaskScope(预览) |
5.4 虚拟线程 vs Go goroutine
| 维度 | Go goroutine | Java 虚拟线程 |
|---|---|---|
| 实现层级 | 运行时层(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。
正确做法:
- 使用
ScopedValue(不可变,作用域绑定,无内存泄漏风险)。 - 将共享数据作为方法参数显式传递。
- 若必须用
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,载体线程被占用。
正确做法:
- 将 native 方法调用封装在平台线程中(
CompletableFuture.supplyAsync+ 平台线程池)。 - 使用纯 Java 实现替代 native 方法(如 BouncyCastle 替代 OpenSSL JNI)。
- 评估 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 往往难以察觉(代码能跑,但性能不达预期)。若无监控,开发者可能误以为虚拟线程”没效果”。
正确做法:
- 生产环境启动时加
-Djdk.tracePinnedThreads=short(仅开发环境,生产用 JFR)。 - 定期用
jcmd <pid> Thread.dump_to_file -format=json抓取线程转储,检查pinned: true的虚拟线程。 - 启用 JFR 录制
jdk.VirtualThreadPinned事件,通过 JDK Mission Control 分析。 - 接入 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 适配)。百万虚拟线程并发读取大文件会导致载体线程被占用。
正确做法:
- 文件 I/O 用平台线程池(
Executors.newFixedThreadPool)。 - 或使用 NIO
AsynchronousFileChannel(真正异步文件 I/O)。 - 评估文件大小与并发度,小文件可接受阻塞。
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 默认使用平台线程,与虚拟线程混用可能导致上下文切换混乱。
正确做法:
- 虚拟线程中直接用同步 HTTP 客户端(
HttpClient.send),不用 Reactor。 - 若必须用 Reactor,配置
Schedulers.fromExecutor(Executors.newVirtualThreadPerTaskExecutor())。 - 评估是否真的需要 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:评估
- 识别 I/O 密集型服务:HTTP 网关、微服务聚合、消息消费者适合迁移。
- 识别 CPU 密集型服务:计算密集型任务不适合,保持平台线程。
- 识别 Pinning 风险点:扫描代码中的
synchronized+ 阻塞、native方法调用。 - 识别 ThreadLocal 滥用:统计
ThreadLocal使用,评估百万线程下的内存风险。
阶段 2:试点
- 选择非核心服务:先在低流量服务试点,验证虚拟线程效果。
- 配置
spring.threads.virtual.enabled=true(Spring Boot 3.2+)。 - 监控 Pinning:启用 JFR 录制
jdk.VirtualThreadPinned事件。 - 性能对比:对比迁移前后的吞吐量、延迟、资源占用。
阶段 3:推广
- 逐步迁移:按服务优先级迁移,先 I/O 密集型服务。
- 重构 Pinning 代码:将
synchronized替换为ReentrantLock。 - 重构 ThreadLocal:迁移到
ScopedValue或显式参数。 - 培训团队:讲解虚拟线程原理与反模式,避免误用。
阶段 4:优化
- 监控虚拟线程数:通过 JMX / Micrometer 监控活跃虚拟线程数。
- 调优载体线程池:必要时调整
-Djdk.virtualThreadParallelism。 - 接入结构化并发:用
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。团队决定迁移到虚拟线程以简化代码。
迁移步骤:
- 评估:网关为 I/O 密集型,主要调用下游订单、用户、商品服务,适合虚拟线程。
- 重构 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());
}
}
- 启用虚拟线程:
spring:
threads:
virtual:
enabled: true
- 迁移 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 试点虚拟线程迁移。
迁移策略:
- 分阶段迁移:先迁移非核心路由(如日志上报),再迁移核心路由。
- 混合模式:Zuul 内部仍用 Netty,但业务逻辑用虚拟线程(通过
Schedulers.fromExecutor桥接)。 - Pinning 修复:Zuul 中大量使用
synchronized,重构为ReentrantLock。
性能对比:
| 指标 | Reactor 版本 | 虚拟线程版本 | 变化 |
|---|---|---|---|
| QPS | 120,000 | 115,000 | -4% |
| P99 延迟 | 50ms | 55ms | +10% |
| 代码行数 | 15,000 | 9,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 服务器 | QPS | P99 延迟 | 内存占用 |
|---|---|---|---|
| Spring WebFlux (Netty) | 50,000 | 30ms | 2GB |
| Spring MVC (Tomcat) + 虚拟线程 | 60,000 | 25ms | 1.5GB |
| Helidon Níma | 80,000 | 20ms | 1GB |
结论:
- Helidon Níma 性能最优,但生态不如 Spring 成熟。
- Spring MVC + 虚拟线程是最佳平衡点(生态成熟 + 性能优秀)。
- 新项目可考虑 Helidon Níma,存量项目优先 Spring Boot 3.2+。
8.4 案例四:Pinning 性能问题排查
背景:某金融交易系统迁移虚拟线程后,QPS 未达预期(仅提升 2 倍,预期 10 倍)。
排查过程:
- 启用 Pinning 追踪:
java -Djdk.tracePinnedThreads=full -jar app.jar
- 发现 Pinning 堆栈:
Thread[#456] pinned due to MONITOR:
com.bank.service.AccountService.debit(AccountService.java:78)
com.bank.service.TransferService.transfer(TransferService.java:45)
- 定位代码:
public class AccountService {
public synchronized void debit(Long accountId, BigDecimal amount) { // synchronized
accountDao.debit(accountId, amount); // 阻塞 I/O
logService.record(accountId, amount); // 阻塞 I/O
}
}
- 重构为 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();
}
}
}
- 验证效果:
| 指标 | 修复前 | 修复后 |
|---|---|---|
| QPS | 10,000 | 50,000 |
| Pinning 次数/秒 | 8000 | 0 |
| 载体线程利用率 | 100%(全部 Pinning) | 20%(正常) |
经验教训:
synchronized+ 阻塞 I/O 是最常见的 Pinning 反模式。- 迁移虚拟线程时必须扫描所有
synchronized块。 - 使用
-Djdk.tracePinnedThreads=full快速定位 Pinning 点。
8.5 案例五:ThreadLocal 内存泄漏
背景:某 SaaS 平台迁移虚拟线程后,运行 24 小时后 OOM。
排查过程:
- Heap Dump 分析:
jcmd <pid> GC.heap_dump heapdump.hprof
# 使用 MAT (Memory Analyzer Tool) 分析
- 发现 ThreadLocal 累积:
java.lang.ThreadLocal$ThreadLocalMap 实例数:2,000,000
每个 ThreadLocalMap 持有 5MB 数据
总占用:10TB(OOM 原因)
- 定位代码:
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() 清理
}
- 修复方案:
方案 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 官方文档与规范
-
Pressler, R. (2023). JEP 444: Virtual Threads. Oracle Corporation. https://openjdk.org/jeps/444
-
Goetz, B. (2023). JEP 453: Structured Concurrency (Preview). Oracle Corporation. https://openjdk.org/jeps/453
-
Pressler, R. (2023). JEP 446: Scoped Values (Preview). Oracle Corporation. https://openjdk.org/jeps/446
-
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
-
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 学术论文
-
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
-
Prokopec, A., Rose, A., & Leopoldseder, D. (2019). Towards Lightweight Threads in Java. Programming’19 Conference.
-
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 书籍
-
Urma, R. G., Warburton, R., & Mycroft, A. (2024). Modern Java in Action: Lambda, Streams, Functional and Reactive Programming (2nd ed.). Manning Publications.
-
Nurkiewicz, T., & Christensen, B. (2023). Java Concurrency in Practice Revisited. Cambridge University Press.
-
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 虚拟线程与响应式编程
- Spring WebFlux 官方文档:https://docs.spring.io/spring-framework/reference/web/webflux.html
- Reactor 官方文档:https://projectreactor.io/docs/core/release/reference/
- 虚拟线程 vs Reactor 选型指南:https://spring.io/blog/2023/11/23/virtual-threads-vs-reactor
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 虚拟线程性能调优
- JDK Mission Control (JMC):https://www.oracle.com/java/technologies/jdk-mission-control.html
- JFR 虚拟线程事件:
jdk.VirtualThreadStart,jdk.VirtualThreadPinned,jdk.VirtualThreadEnd - jcmd 线程转储指南:https://docs.oracle.com/en/java/javase/21/troubleshoot/
11.5 虚拟线程生态
- Spring Boot 3.2+ 虚拟线程支持:https://docs.spring.io/spring-boot/docs/3.2/reference/htmlsingle/#features.task-execution-and-scheduling.threads.virtual
- Helidon Níma 文档:https://helidon.io/docs/v4/se/guides/nima
- Quarkus 虚拟线程:https://quarkus.io/guides/virtual-threads
- gRPC Java 虚拟线程:https://github.com/grpc/grpc-java/issues/10529
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() | 创建 ThreadFactory | ThreadFactory 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 返回 null | try (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() | 等待子任务按策略完成;失败时抛 FailedException | scope.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() | 创建 ScopedValue | static 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 file | JSON 线程转储(含虚拟线程) | 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 起正式)