Java深入解析篇二十二之虚拟线程
·
Java 虚拟线程(Virtual Threads)详解
本文基于 JDK 21(JEP 444)正式特性,深入分析虚拟线程原理、调度机制、实战用法与迁移策略,配合完整可运行代码示例。
目录
- 虚拟线程概述
- 虚拟线程 vs 平台线程
- 虚拟线程调度原理
- 创建虚拟线程
- 虚拟线程与结构化并发
- synchronized 与 Pin 问题
- ReentrantLock 替代方案
- ThreadLocal 与 ScopedValue
- 虚拟线程池模式(每任务一线程)
- 虚拟线程与 IO 密集型任务
- 虚拟线程监控
- 从线程池迁移到虚拟线程
- 虚拟线程适用场景与反模式
- 性能对比测试
- 最佳实践
一、虚拟线程概述
1.1 什么是虚拟线程
虚拟线程(Virtual Thread)是 JDK 21 正式引入的轻量级线程实现(JEP 444),由 JVM 而非操作系统调度。它保留了 java.lang.Thread 的完整 API,但底层采用 M:N 调度模型,将大量虚拟线程映射到少量平台线程(载体线程)上执行。
核心设计目标:
- 保持 Java “每请求一线程” 的简单编程模型
- 获得接近异步/响应式编程的吞吐量
- 无需学习复杂的异步 API(CompletableFuture 链式调用等)
1.2 演进历程
| 版本 | JEP | 状态 | 关键变化 |
|---|---|---|---|
| JDK 19 | JEP 425 | 第一次预览 | 需 --enable-preview,基础 API |
| JDK 20 | JEP 436 | 第二次预览 | 小幅 API 调整,性能优化 |
| JDK 21 | JEP 444 | 正式发布 | 无需预览标志,API 稳定 |
| JDK 24 | — | 持续优化 | synchronized 内阻塞不再 Pin |
1.3 为什么需要虚拟线程
// 传统平台线程:创建成本高,数量受限
// 一台服务器通常只能创建 200~2000 个平台线程
ExecutorService pool = Executors.newFixedThreadPool(200);
// 当并发请求 > 200 时,请求排队等待 → 吞吐量受限
// 虚拟线程:创建成本极低,可轻松创建百万级
ExecutorService executor = Executors.newVirtualThreadPerTaskExecutor();
// 每个请求一个虚拟线程,10万并发无压力
根本矛盾:
- 平台线程 1:1 映射 OS 线程,每个占用 ~1MB 栈内存
- OS 线程上下文切换代价高(内核态切换)
- 服务器并发能力被线程数量硬性限制
- 异步编程(Reactor/CompletableFuture)虽然解决了吞吐量问题,但代码复杂度高
虚拟线程的解法:
- 虚拟线程是 JVM 管理的轻量对象,栈在堆上按需分配
- 阻塞时自动让出载体线程(用户态切换,无内核参与)
- 保持同步阻塞的编程风格,获得异步的吞吐量
1.4 启用与基本使用
// JDK 21+ 无需任何特殊标志,直接使用
public class VirtualThreadDemo {
public static void main(String[] args) throws Exception {
// 最简单的创建方式
Thread vThread = Thread.startVirtualThread(() -> {
System.out.println("Hello from virtual thread!");
System.out.println("Is virtual: " + Thread.currentThread().isVirtual());
System.out.println("Thread: " + Thread.currentThread());
});
vThread.join();
// 输出:
// Hello from virtual thread!
// Is virtual: true
// Thread: VirtualThread[#22]/runnable@ForkJoinPool-1-worker-1
}
}
二、虚拟线程 vs 平台线程
2.1 架构对比
// 平台线程:1:1 映射 OS 线程
Thread platformThread = Thread.ofPlatform()
.name("platform-worker")
.start(() -> {
// 此线程直接对应一个 OS 线程
// 阻塞时 OS 线程被占用,无法执行其他任务
System.out.println("Platform: " + Thread.currentThread());
// 输出: Platform: Thread[#21,platform-worker,5,main]
});
// 虚拟线程:M:N 映射到载体线程
Thread virtualThread = Thread.ofVirtual()
.name("virtual-worker")
.start(() -> {
// 此线程由 JVM 调度,运行在某个 Carrier 上
// 阻塞时自动卸载,Carrier 可执行其他虚拟线程
System.out.println("Virtual: " + Thread.currentThread());
// 输出: Virtual: VirtualThread[#22,virtual-worker]/runnable@ForkJoinPool-1-worker-1
});
2.2 核心差异详解
| 维度 | 平台线程 | 虚拟线程 |
|---|---|---|
| 底层实现 | OS 线程(pthread/WinThread) | JVM Continuation 对象 |
| 栈内存 | 固定 ~1MB(-Xss 配置) | 初始几KB,动态增长,存于堆 |
| 创建开销 | ~1ms(系统调用) | ~1μs(对象分配) |
| 最大数量 | 数千(OS 限制) | 数百万(堆内存限制) |
| 调度 | OS 内核抢占式调度 | JVM 协作式调度(ForkJoinPool) |
| 阻塞行为 | 占用 OS 线程 | 自动卸载,释放 Carrier |
| 上下文切换 | 内核态(~1-10μs) | 用户态(~几百ns) |
| 线程池 | 必须池化 | 不应池化 |
| 守护属性 | 可配置 | 始终为守护线程 |
| 优先级 | 1-10 可设置 | 固定 NORM_PRIORITY |
| stop/suspend/resume | 已废弃但存在 | 不支持(抛 UnsupportedOperationException) |
2.3 内存占用实测
public class MemoryComparison {
public static void main(String[] args) throws Exception {
// 测试:创建大量线程观察内存
int count = 100_000;
// 虚拟线程:轻松创建 10 万个
long startMem = Runtime.getRuntime().totalMemory() - Runtime.getRuntime().freeMemory();
CountDownLatch latch = new CountDownLatch(count);
try (var executor = Executors.newVirtualThreadPerTaskExecutor()) {
for (int i = 0; i < count; i++) {
executor.submit(() -> {
try {
Thread.sleep(Duration.ofSeconds(1));
} catch (InterruptedException e) {
Thread.currentThread().interrupt();
} finally {
latch.countDown();
}
});
}
latch.await();
}
long endMem = Runtime.getRuntime().totalMemory() - Runtime.getRuntime().freeMemory();
System.out.println("100,000 虚拟线程内存增量: " + (endMem - startMem) / 1024 / 1024 + " MB");
// 典型输出: 约 100~200 MB(每线程仅几KB)
// 对比:10万个平台线程 → 需要 ~100GB 栈内存 → 不可能创建
}
}
2.4 API 兼容性
public class ApiCompatibility {
public static void main(String[] args) throws Exception {
Thread vt = Thread.ofVirtual().start(() -> {
// 以下 API 在虚拟线程中正常工作
System.out.println("ID: " + Thread.currentThread().threadId());
System.out.println("Name: " + Thread.currentThread().getName());
System.out.println("State: " + Thread.currentThread().getState());
System.out.println("IsVirtual: " + Thread.currentThread().isVirtual());
System.out.println("IsDaemon: " + Thread.currentThread().isDaemon()); // 始终 true
// 以下操作在虚拟线程中无效或抛异常
// Thread.currentThread().setPriority(10); // 无效,忽略
// Thread.currentThread().setDaemon(false); // 抛 UnsupportedOperationException
});
vt.join();
// 虚拟线程支持的 Thread API
System.out.println("isAlive: " + vt.isAlive());
System.out.println("getState: " + vt.getState()); // TERMINATED
vt.interrupt(); // 支持中断
}
}
三、虚拟线程调度原理
3.1 M:N 调度模型
虚拟线程采用 M:N 调度模型:M 个虚拟线程映射到 N 个载体线程(Carrier Thread)上执行。
┌────────────────────────────────────────────────────────────┐
│ 虚拟线程层(M 个,由 JVM 管理) │
│ │
│ VT-1 VT-2 VT-3 VT-4 VT-5 ... VT-1000000 │
│ │ │ │ │ │ │ │
│ └──────┴──────┼──────┴──────┴──────────────┘ │
│ │ │
│ ForkJoinPool(调度器) │
│ 工作窃取 + 任务队列 │
│ │ │
│ ┌─────────────┼─────────────────┐ │
│ │ │ │ │
│ ▼ ▼ ▼ │
│ Carrier-1 Carrier-2 ... Carrier-N │
│ (平台线程) (平台线程) (平台线程) │
│ │
│ N = CPU 核心数(默认),可通过 JVM 参数调整 │
└────────────────────────────────────────────────────────────┘
3.2 Carrier Thread(载体线程)
载体线程本质是平台线程,由 ForkJoinPool 管理:
public class CarrierThreadDemo {
public static void main(String[] args) throws Exception {
// 观察虚拟线程运行在哪个 Carrier 上
for (int i = 0; i < 10; i++) {
Thread.startVirtualThread(() -> {
// toString() 中包含 Carrier 信息
System.out.println(Thread.currentThread());
// 输出: VirtualThread[#25]/runnable@ForkJoinPool-1-worker-3
// @后面就是 Carrier 线程名
}).join();
}
}
}
Carrier 线程池配置:
# 默认载体线程数 = Runtime.getRuntime().availableProcessors()
# 可通过 JVM 系统属性调整:
# 设置载体线程并行度(正常情况下的 Carrier 数量)
-Djdk.virtualThreadScheduler.parallelism=16
# 设置最大载体线程数(当发生 Pin 时可扩展到此数量)
-Djdk.virtualThreadScheduler.maxPoolSize=256
# 设置最小可运行载体线程数(低于此值时强制创建新 Carrier)
-Djdk.virtualThreadScheduler.minRunnable=1
3.3 Continuation(续体)机制
Continuation 是虚拟线程挂起/恢复的核心单元:
虚拟线程执行流程:
│
├── 1. 虚拟线程被提交到 ForkJoinPool
│
├── 2. 空闲 Carrier 从队列取出任务
│ └── mount():将虚拟线程的 Continuation 绑定到 Carrier
│
├── 3. 执行 Continuation.run()
│ └── 虚拟线程代码开始执行
│
├── 4. 遇到阻塞操作(如 IO、sleep、LockSupport.park)
│ ├── Continuation.yield():保存当前栈帧到堆内存
│ ├── unmount():虚拟线程与 Carrier 解绑
│ └── Carrier 回到 ForkJoinPool,可执行其他虚拟线程
│
├── 5. 阻塞操作完成(如 IO 就绪)
│ ├── 虚拟线程重新提交到 ForkJoinPool 队列
│ └── 等待下一个空闲 Carrier
│
└── 6. 新 Carrier 取出任务
├── mount():重新绑定
├── Continuation.run():从上次 yield 点恢复执行
└── 继续执行直到完成或再次阻塞
3.4 调度器内部实现
// 虚拟线程调度器的核心逻辑(简化版,基于 OpenJDK 源码)
// java.lang.VirtualThread 内部
class VirtualThread extends Thread {
// 调度器:全局共享的 ForkJoinPool
private static final ForkJoinPool DEFAULT_SCHEDULER = createScheduler();
private static ForkJoinPool createScheduler() {
int parallelism = Runtime.getRuntime().availableProcessors();
// 可通过系统属性覆盖
String prop = System.getProperty("jdk.virtualThreadScheduler.parallelism");
if (prop != null) parallelism = Integer.parseInt(prop);
return new ForkJoinPool(
parallelism,
factory, // 创建 Carrier 线程的工厂
null, // 无 UncaughtExceptionHandler
true, // asyncMode = true(FIFO,适合IO任务)
0, parallelism, // core/max pool size
1, null, // minimumRunnable, saturate
0L, TimeUnit.MILLISECONDS
);
}
// 虚拟线程启动时提交到调度器
@Override
public void start() {
// 将自身包装为任务提交到 ForkJoinPool
DEFAULT_SCHEDULER.execute(this::run);
}
// park(阻塞)时让出 Carrier
private void park() {
// 保存 Continuation 状态
Continuation.yield(VTHREAD_SCOPE);
// --- 此处虚拟线程被挂起 ---
// --- 恢复后从这里继续执行 ---
}
}
3.5 阻塞操作如何触发卸载
public class BlockingBehaviorDemo {
public static void main(String[] args) throws Exception {
// 以下操作都会触发虚拟线程卸载(yield)
Thread.startVirtualThread(() -> {
try {
// 1. Thread.sleep → 触发 yield
Thread.sleep(Duration.ofMillis(100));
// 2. IO 操作 → 触发 yield(底层使用 NIO)
var client = java.net.http.HttpClient.newHttpClient();
var request = java.net.http.HttpRequest.newBuilder()
.uri(java.net.URI.create("https://httpbin.org/delay/1"))
.build();
client.send(request, java.net.http.HttpResponse.BodyHandlers.ofString());
// 3. LockSupport.park → 触发 yield
java.util.concurrent.locks.LockSupport.parkNanos(1_000_000);
// 4. BlockingQueue.take → 触发 yield
// 5. CountDownLatch.await → 触发 yield
// 6. Future.get → 触发 yield
} catch (Exception e) {
e.printStackTrace();
}
}).join();
}
}
四、创建虚拟线程
4.1 Thread.ofVirtual()(Builder 模式)
public class VirtualThreadBuilder {
public static void main(String[] args) throws Exception {
// 基本创建
Thread vt1 = Thread.ofVirtual().start(() -> {
System.out.println("Basic virtual thread: " + Thread.currentThread());
});
vt1.join();
// 带名称的虚拟线程(便于调试)
Thread vt2 = Thread.ofVirtual()
.name("my-worker")
.start(() -> {
System.out.println("Named: " + Thread.currentThread().getName());
});
vt2.join();
// 带名称前缀和自增序号(批量创建时推荐)
Thread.Builder.OfVirtual builder = Thread.ofVirtual().name("task-", 0);
for (int i = 0; i < 5; i++) {
builder.start(() -> {
System.out.println("Running: " + Thread.currentThread().getName());
// 输出: task-0, task-1, task-2, task-3, task-4
});
}
// 创建但不启动
Thread vt3 = Thread.ofVirtual().unstarted(() -> {
System.out.println("Not started yet");
});
System.out.println("State: " + vt3.getState()); // NEW
vt3.start();
vt3.join();
}
}
4.2 Thread.startVirtualThread()(最简方式)
public class StartVirtualThread {
public static void main(String[] args) throws Exception {
// 一行代码创建并启动虚拟线程
Thread vt = Thread.startVirtualThread(() -> {
System.out.println("Simplest way to create virtual thread");
System.out.println("Thread: " + Thread.currentThread());
});
vt.join();
// 带返回值的场景:配合 Future
ExecutorService executor = Executors.newVirtualThreadPerTaskExecutor();
Future<String> future = executor.submit(() -> {
Thread.sleep(100);
return "computed result";
});
System.out.println(future.get()); // computed result
executor.close();
}
}
4.3 Executors.newVirtualThreadPerTaskExecutor()(批量任务推荐)
public class VirtualThreadExecutor {
public static void main(String[] args) throws Exception {
// try-with-resources:close() 时自动等待所有任务完成
try (var executor = Executors.newVirtualThreadPerTaskExecutor()) {
// 提交大量任务,每个任务独占一个虚拟线程
List<Future<String>> futures = new ArrayList<>();
for (int i = 0; i < 100_000; i++) {
final int taskId = i;
futures.add(executor.submit(() -> {
// 模拟 IO 操作
Thread.sleep(Duration.ofMillis(100));
return "Task-" + taskId + " done by " + Thread.currentThread();
}));
}
// 收集结果
int successCount = 0;
for (Future<String> f : futures) {
f.get(); // 阻塞等待(在虚拟线程中也会 yield)
successCount++;
}
System.out.println("Completed: " + successCount + " tasks");
}
// executor.close() 确保所有任务完成后才继续
System.out.println("All tasks finished");
}
}
4.4 ThreadFactory(兼容旧 API)
public class VirtualThreadFactory {
public static void main(String[] args) throws Exception {
// 创建虚拟线程工厂
ThreadFactory factory = Thread.ofVirtual()
.name("legacy-worker-", 0)
.factory();
// 用于需要 ThreadFactory 的旧 API
// 例如:某些框架接受 ThreadFactory 参数
Thread t1 = factory.newThread(() -> {
System.out.println("Created by factory: " + Thread.currentThread().getName());
});
t1.start();
t1.join();
// 配合 ScheduledExecutorService 等(注意:不建议用于虚拟线程池化)
// 这里仅演示 API 兼容性
ExecutorService legacyPool = Executors.newFixedThreadPool(4, factory);
legacyPool.submit(() -> System.out.println("In legacy pool: " + Thread.currentThread()));
legacyPool.shutdown();
}
}
4.5 创建方式对比
| 方式 | 适用场景 | 特点 |
|---|---|---|
Thread.ofVirtual().start() |
单个线程,需配置名称 | Builder 模式,灵活 |
Thread.startVirtualThread() |
单个线程,最简场景 | 一行代码,无配置 |
newVirtualThreadPerTaskExecutor() |
批量任务提交 | 自动管理生命周期 |
Thread.ofVirtual().factory() |
兼容需要 ThreadFactory 的 API | 适配器模式 |
五、虚拟线程与结构化并发
5.1 结构化并发概念
结构化并发(Structured Concurrency,JEP 453)将并发任务视为结构化的工作单元,确保子任务的生命周期不超过父任务,实现取消传播和错误传播。
// 非结构化并发(传统方式):子任务可能逃逸
public String handleRequest_NonStructured(String id) {
// 这两个异步任务的生命周期不受控
CompletableFuture<String> userFuture = CompletableFuture.supplyAsync(() -> findUser(id));
CompletableFuture<String> orderFuture = CompletableFuture.supplyAsync(() -> fetchOrder(id));
// 如果这里抛异常,上面的 Future 可能还在运行 → 资源泄漏
// 如果调用者取消请求,子任务不会自动取消
return userFuture.join() + orderFuture.join();
}
5.2 StructuredTaskScope 基本用法
// 需要 --enable-preview(JDK 21 中为预览特性)
import jdk.incubator.concurrent.StructuredTaskScope;
import jdk.incubator.concurrent.Subtask;
public class StructuredConcurrencyDemo {
record Response(String user, String order) {}
// 结构化并发:所有子任务在 scope 内完成或取消
public Response handleRequest(String id) throws Exception {
try (var scope = new StructuredTaskScope.ShutdownOnFailure()) {
// fork() 在新虚拟线程中执行子任务
Subtask<String> userTask = scope.fork(() -> findUser(id));
Subtask<String> orderTask = scope.fork(() -> fetchOrder(id));
// 等待所有子任务完成
scope.join();
// 如果任一子任务失败,抛出异常
scope.throwIfFailed();
// 所有子任务成功,获取结果
return new Response(userTask.get(), orderTask.get());
}
// scope.close() 确保:
// ① 所有子任务已完成或已取消
// ② 不会有子任务逃逸到 scope 外
}
private String findUser(String id) throws Exception {
Thread.sleep(100); // 模拟 IO
return "User-" + id;
}
private String fetchOrder(String id) throws Exception {
Thread.sleep(200); // 模拟 IO
return "Order-" + id;
}
}
5.3 ShutdownOnFailure 策略
public class ShutdownOnFailureDemo {
public static void main(String[] args) throws Exception {
try (var scope = new StructuredTaskScope.ShutdownOnFailure()) {
Subtask<String> task1 = scope.fork(() -> {
Thread.sleep(100);
return "result-1";
});
Subtask<String> task2 = scope.fork(() -> {
Thread.sleep(50);
throw new RuntimeException("task2 failed!");
// task2 失败 → scope 自动取消 task1
});
Subtask<String> task3 = scope.fork(() -> {
Thread.sleep(200);
return "result-3";
// task3 会被取消,不会执行完
});
scope.join(); // 等待(直到全部完成或任一失败)
try {
scope.throwIfFailed(); // 抛出 task2 的异常
} catch (RuntimeException e) {
System.out.println("Caught: " + e.getMessage());
// 输出: Caught: task2 failed!
}
}
}
}
5.4 ShutdownOnSuccess 策略
public class ShutdownOnSuccessDemo {
public static void main(String[] args) throws Exception {
// 竞速模式:取最快成功的结果
try (var scope = new StructuredTaskScope.ShutdownOnSuccess<String>()) {
// 向多个数据源同时查询,取最快返回的
scope.fork(() -> querySourceA()); // 耗时 300ms
scope.fork(() -> querySourceB()); // 耗时 100ms ← 最快
scope.fork(() -> querySourceC()); // 耗时 200ms
scope.join();
String fastest = scope.result();
System.out.println("Fastest result: " + fastest);
// 输出: Fastest result: source-B-data
// 其余任务自动取消
}
}
static String querySourceA() throws Exception {
Thread.sleep(300); return "source-A-data";
}
static String querySourceB() throws Exception {
Thread.sleep(100); return "source-B-data";
}
static String querySourceC() throws Exception {
Thread.sleep(200); return "source-C-data";
}
}
5.5 自定义 TaskScope
public class CollectingScope<T> extends StructuredTaskScope<T> {
private final List<T> results = new CopyOnWriteArrayList<>();
private final List<Throwable> errors = new CopyOnWriteArrayList<>();
@Override
protected void handleComplete(Subtask<? extends T> subtask) {
switch (subtask.state()) {
case SUCCESS -> results.add(subtask.get());
case FAILED -> errors.add(subtask.exception());
case UNAVAILABLE -> { /* 被取消 */ }
}
}
public List<T> results() { return List.copyOf(results); }
public List<Throwable> errors() { return List.copyOf(errors); }
}
// 使用:收集所有成功结果,忽略失败
public class CustomScopeDemo {
public static void main(String[] args) throws Exception {
try (var scope = new CollectingScope<String>()) {
for (int i = 0; i < 10; i++) {
final int id = i;
scope.fork(() -> {
if (id % 3 == 0) throw new RuntimeException("fail-" + id);
Thread.sleep(50);
return "data-" + id;
});
}
scope.join();
System.out.println("Success: " + scope.results());
System.out.println("Errors: " + scope.errors().size());
}
}
}
六、synchronized 与 Pin 问题
6.1 什么是 Pin(钉住)
当虚拟线程在 synchronized 块或 native 方法内执行阻塞操作时,JVM 无法将其从 Carrier 上卸载(因为 monitor 锁与 OS 线程绑定),虚拟线程被"钉住"(Pinned)在 Carrier 上。
public class PinProblemDemo {
private static final Object LOCK = new Object();
public static void main(String[] args) throws Exception {
// 模拟:大量虚拟线程在 synchronized 内执行阻塞 IO
try (var executor = Executors.newVirtualThreadPerTaskExecutor()) {
for (int i = 0; i < 1000; i++) {
executor.submit(() -> {
synchronized (LOCK) {
// 阻塞 IO 在 synchronized 内 → Pin!
// Carrier 被独占,无法服务其他虚拟线程
try {
Thread.sleep(Duration.ofMillis(100));
} catch (InterruptedException e) {
Thread.currentThread().interrupt();
}
}
});
}
}
// 如果 Carrier 数 = 8,则同时最多 8 个虚拟线程在执行
// 其余 992 个等待 → 严重降低并发度
}
}
6.2 Pin 的触发条件
public class PinTriggers {
// 触发条件 1:synchronized 方法内阻塞
public synchronized void pinnedMethod() throws Exception {
// 整个方法执行期间,虚拟线程被 Pin
Thread.sleep(1000); // 阻塞 → Carrier 被占用 1 秒
}
// 触发条件 2:synchronized 块内阻塞
public void pinnedBlock() throws Exception {
synchronized (this) {
// 块内阻塞 → Pin
var conn = DriverManager.getConnection("jdbc:mysql://...");
conn.createStatement().executeQuery("SELECT SLEEP(1)");
}
}
// 触发条件 3:native 方法内阻塞
public void nativePinning() {
// JNI 调用内部如果阻塞 → Pin
nativeBlockingCall();
}
private native void nativeBlockingCall();
// 不触发 Pin 的情况:
public void noPinning() throws Exception {
// synchronized 内只有 CPU 计算(不阻塞)→ 短暂 Pin,影响小
synchronized (this) {
int sum = 0;
for (int i = 0; i < 1000; i++) sum += i;
}
// 阻塞操作在 synchronized 外 → 不 Pin
Thread.sleep(1000); // 正常 yield,Carrier 释放
}
}
6.3 Pin 的危害分析
假设 Carrier 线程数 = 8(8核CPU)
场景:1000 个虚拟线程同时执行 synchronized + sleep(1s)
时间线:
T=0s: 8个虚拟线程获得 Carrier,进入 synchronized → Pin
其余 992 个等待
T=1s: 前 8 个完成,释放 Carrier
下 8 个获得 Carrier → Pin
...
T=125s: 全部完成(1000/8 × 1s)
对比(无 Pin,使用 ReentrantLock):
T=0s: 1000 个虚拟线程全部启动
进入 sleep → 全部 yield → Carrier 空闲
T=1s: 全部完成
性能差距:125 倍!
6.4 检测 Pin 事件
# 方法一:JVM 参数(开发/测试环境)
java -Djdk.tracePinnedThreads=short -jar app.jar
java -Djdk.tracePinnedThreads=full -jar app.jar
# 输出示例(short 模式):
# Thread[#42,ForkJoinPool-1-worker-1,5,CarrierThreads]
# java.base/java.lang.VirtualThread$VThreadContinuation.onPinned(VirtualThread.java:180)
# com.example.PinProblemDemo.pinnedMethod(PinProblemDemo.java:15) <== monitors:1
# 方法二:JFR(生产环境)
java -XX:StartFlightRecording=filename=recording.jfr,duration=60s -jar app.jar
# 然后用 JMC 分析 jdk.VirtualThreadPinned 事件
// 方法三:程序化检测(JDK 21+)
// 通过 JFR API 监听 Pin 事件
import jdk.jfr.consumer.RecordingStream;
import java.time.Duration;
public class PinDetector {
public static void startMonitoring() {
RecordingStream rs = new RecordingStream();
rs.enable("jdk.VirtualThreadPinned").withThreshold(Duration.ofMillis(20));
rs.onEvent("jdk.VirtualThreadPinned", event -> {
System.out.println("[PIN DETECTED] Duration: " + event.getDuration());
System.out.println(" Thread: " + event.getThread().getJavaName());
// 记录日志、触发告警
});
rs.startAsync();
}
}
6.5 JDK 24 的改进
JDK 24(JEP 491 相关改进):
│
├── synchronized 内的阻塞操作不再导致 Pin
│ └── JVM 内部将 synchronized 改为基于 ReentrantLock 的实现
│
├── 影响:
│ ├── 旧代码无需修改即可受益
│ ├── 性能大幅提升(不再 Carrier 饥饿)
│ └── -Djdk.tracePinnedThreads 不再报告 synchronized Pin
│
└── 注意:
├── 仍需 JDK 21 LTS 的项目应主动使用 ReentrantLock
└── native 方法内的阻塞仍会 Pin
七、ReentrantLock 替代方案
7.1 基本替换模式
public class LockMigration {
// ===== 替换前:synchronized(导致 Pin)=====
private final Object lock = new Object();
private int counter_Sync = 0;
public int incrementSync() throws Exception {
synchronized (lock) {
Thread.sleep(10); // 模拟 IO → Pin!
return ++counter_Sync;
}
}
// ===== 替换后:ReentrantLock(安全)=====
private final ReentrantLock reLock = new ReentrantLock();
private int counter_Lock = 0;
public int incrementLock() throws Exception {
reLock.lock();
try {
Thread.sleep(10); // 模拟 IO → 正常 yield,不 Pin
return ++counter_Lock;
} finally {
reLock.unlock(); // 必须在 finally 中释放!
}
}
}
7.2 可中断锁获取
public class InterruptibleLock {
private final ReentrantLock lock = new ReentrantLock();
public void doWork() throws InterruptedException {
// 可中断的锁获取(虚拟线程友好)
lock.lockInterruptibly();
try {
// 临界区
Thread.sleep(1000);
} finally {
lock.unlock();
}
}
public boolean tryDoWork() {
// 非阻塞尝试获取锁
if (lock.tryLock()) {
try {
// 临界区
return true;
} finally {
lock.unlock();
}
}
return false; // 获取失败,不阻塞
}
public boolean timedDoWork() throws InterruptedException {
// 超时获取锁
if (lock.tryLock(5, TimeUnit.SECONDS)) {
try {
// 临界区
return true;
} finally {
lock.unlock();
}
}
return false; // 超时
}
}
7.3 读写锁替代
public class ReadWriteLockMigration {
// ===== 替换前 =====
private final Map<String, String> cache = new HashMap<>();
public synchronized String getSync(String key) throws Exception {
Thread.sleep(10); // IO in synchronized → Pin
return cache.get(key);
}
public synchronized void putSync(String key, String value) throws Exception {
Thread.sleep(10); // Pin
cache.put(key, value);
}
// ===== 替换后:ReadWriteLock =====
private final ReadWriteLock rwLock = new ReentrantReadWriteLock();
private final Map<String, String> safeCache = new HashMap<>();
public String getLock(String key) throws Exception {
rwLock.readLock().lock();
try {
Thread.sleep(10); // 不 Pin,多个读者可并发
return safeCache.get(key);
} finally {
rwLock.readLock().unlock();
}
}
public void putLock(String key, String value) throws Exception {
rwLock.writeLock().lock();
try {
Thread.sleep(10); // 不 Pin
safeCache.put(key, value);
} finally {
rwLock.writeLock().unlock();
}
}
}
7.4 StampedLock(乐观读)
public class StampedLockDemo {
private final StampedLock sl = new StampedLock();
private double x, y;
// 乐观读(无锁,适合读多写少)
public double distanceFromOrigin() {
long stamp = sl.tryOptimisticRead(); // 获取乐观读戳
double currentX = x, currentY = y; // 读取共享变量
if (!sl.validate(stamp)) {
// 乐观读期间有写操作 → 升级为悲观读锁
stamp = sl.readLock();
try {
currentX = x;
currentY = y;
} finally {
sl.unlockRead(stamp);
}
}
return Math.sqrt(currentX * currentX + currentY * currentY);
}
// 写操作
public void move(double deltaX, double deltaY) {
long stamp = sl.writeLock();
try {
x += deltaX;
y += deltaY;
} finally {
sl.unlockWrite(stamp);
}
}
}
7.5 迁移检查清单
synchronized → ReentrantLock 迁移步骤:
□ 1. 识别所有 synchronized 块/方法中包含阻塞操作的位置
□ 2. 将 synchronized 替换为 ReentrantLock
□ 3. 确保 unlock() 在 finally 块中
□ 4. 读多写少场景考虑 ReadWriteLock
□ 5. 运行 -Djdk.tracePinnedThreads=full 验证无 Pin
□ 6. 压测对比迁移前后吞吐量
八、ThreadLocal 与 ScopedValue
8.1 ThreadLocal 在虚拟线程中的问题
public class ThreadLocalProblem {
// 问题:每个虚拟线程都会持有 ThreadLocal 副本
private static final ThreadLocal<byte[]> BUFFER =
ThreadLocal.withInitial(() -> new byte[1024 * 1024]); // 1MB buffer
public static void main(String[] args) throws Exception {
try (var executor = Executors.newVirtualThreadPerTaskExecutor()) {
for (int i = 0; i < 100_000; i++) {
executor.submit(() -> {
byte[] buf = BUFFER.get(); // 每个虚拟线程分配 1MB
// 100,000 × 1MB = 100GB → OOM!
process(buf);
});
}
}
}
static void process(byte[] buf) { /* ... */ }
}
8.2 ThreadLocal 正确使用方式
public class ThreadLocalBestPractice {
// 小型、不可变的上下文信息仍可使用 ThreadLocal
private static final ThreadLocal<String> REQUEST_ID =
ThreadLocal.withInitial(() -> UUID.randomUUID().toString());
// 但要注意清理(虚拟线程虽然短命,但养成习惯)
public void handleRequest() {
try {
String reqId = REQUEST_ID.get();
System.out.println("Processing request: " + reqId);
// 业务逻辑
} finally {
REQUEST_ID.remove(); // 显式清理
}
}
// 更好的方式:使用参数传递代替 ThreadLocal
public void handleRequestExplicit(String requestId) {
System.out.println("Processing request: " + requestId);
processOrder(requestId); // 显式传递
}
private void processOrder(String requestId) {
// 无需 ThreadLocal,参数已在调用链中
}
}
8.3 ScopedValue 基本用法
// ScopedValue(JDK 21 预览,需 --enable-preview)
public class ScopedValueDemo {
// 声明:static final,不可变
private static final ScopedValue<String> USER = ScopedValue.newInstance();
private static final ScopedValue<String> REQUEST_ID = ScopedValue.newInstance();
public static void main(String[] args) {
// 绑定值并执行
ScopedValue.where(USER, "alice")
.where(REQUEST_ID, "req-001")
.run(() -> {
// 在作用域内读取
System.out.println("User: " + USER.get()); // alice
System.out.println("ReqId: " + REQUEST_ID.get()); // req-001
// 调用下游方法,无需传参
processOrder();
});
// 作用域外不可访问
// USER.get(); // 抛 NoSuchElementException
}
private static void processOrder() {
// 深层调用链中直接读取,无需层层传参
System.out.println("Processing order for: " + USER.get());
auditLog();
}
private static void auditLog() {
System.out.println("Audit: user=" + USER.get() + ", req=" + REQUEST_ID.get());
}
}
8.4 ScopedValue 嵌套与覆盖
public class ScopedValueNesting {
private static final ScopedValue<String> TENANT = ScopedValue.newInstance();
public static void main(String[] args) {
ScopedValue.where(TENANT, "tenant-A").run(() -> {
System.out.println("Outer: " + TENANT.get()); // tenant-A
// 嵌套作用域:覆盖值
ScopedValue.where(TENANT, "tenant-B").run(() -> {
System.out.println("Inner: " + TENANT.get()); // tenant-B
});
// 回到外层作用域:值恢复
System.out.println("After inner: " + TENANT.get()); // tenant-A
});
}
}
8.5 ScopedValue 与结构化并发集成
public class ScopedValueWithConcurrency {
private static final ScopedValue<String> TRACE_ID = ScopedValue.newInstance();
public static void main(String[] args) throws Exception {
ScopedValue.where(TRACE_ID, "trace-xyz").run(() -> {
try {
// 结构化并发中,子任务自动继承 ScopedValue
var scope = new StructuredTaskScope.ShutdownOnFailure();
try (scope) {
scope.fork(() -> {
// 子虚拟线程中可读取父作用域的 ScopedValue
System.out.println("Child sees: " + TRACE_ID.get());
// 输出: Child sees: trace-xyz
return queryDB();
});
scope.fork(() -> {
System.out.println("Child sees: " + TRACE_ID.get());
return callAPI();
});
scope.join();
}
} catch (Exception e) {
e.printStackTrace();
}
});
}
static String queryDB() throws Exception { Thread.sleep(100); return "db-result"; }
static String callAPI() throws Exception { Thread.sleep(100); return "api-result"; }
}
8.6 ThreadLocal vs ScopedValue 对比
| 维度 | ThreadLocal | ScopedValue |
|---|---|---|
| 可变性 | 可变(set/get) | 不可变(绑定后只读) |
| 生命周期 | 手动管理(需 remove) | 自动(作用域结束即释放) |
| 内存开销 | 每线程一份拷贝 | 无拷贝,共享引用 |
| 虚拟线程友好度 | 差(百万线程×大对象) | 优(无额外内存) |
| 继承 | InheritableThreadLocal | 结构化并发自动继承 |
| 线程安全 | 需额外保证 | 天然不可变,线程安全 |
| 状态 | 正式(JDK 1.2+) | 预览(JDK 21+) |
| 适用场景 | 兼容旧代码、小型上下文 | 新代码、请求上下文传递 |
九、虚拟线程池模式(每任务一线程)
9.1 为什么不需要池化虚拟线程
public class NoPoolingNeeded {
public static void main(String[] args) throws Exception {
// 错误做法:池化虚拟线程(限制并发,违背设计初衷)
// ExecutorService wrong = Executors.newFixedThreadPool(10,
// Thread.ofVirtual().factory()); // 不要这样做!
// 正确做法:每任务一线程
try (var executor = Executors.newVirtualThreadPerTaskExecutor()) {
// 提交 10 万个任务,每个任务立即获得虚拟线程
for (int i = 0; i < 100_000; i++) {
final int id = i;
executor.submit(() -> {
// 每个任务独占一个虚拟线程
// 阻塞时自动 yield,不影响其他任务
Thread.sleep(Duration.ofMillis(50));
return "result-" + id;
});
}
}
// 虚拟线程创建成本 ~1μs,无需复用
// 阻塞时自动释放 Carrier,无需担心资源浪费
}
}
9.2 使用 Semaphore 限流
public class SemaphoreRateLimit {
// 限制对下游数据库的并发访问数
private static final Semaphore DB_LIMITER = new Semaphore(20);
// 限制对外部 API 的并发调用数
private static final Semaphore API_LIMITER = new Semaphore(100);
public static void main(String[] args) throws Exception {
try (var executor = Executors.newVirtualThreadPerTaskExecutor()) {
for (int i = 0; i < 10_000; i++) {
executor.submit(() -> processRequest());
}
}
}
static void processRequest() throws Exception {
// 10000 个虚拟线程并发,但数据库访问限制为 20
DB_LIMITER.acquire();
try {
queryDatabase(); // 最多 20 个并发 DB 查询
} finally {
DB_LIMITER.release();
}
// API 调用限制为 100 并发
API_LIMITER.acquire();
try {
callExternalApi(); // 最多 100 个并发 API 调用
} finally {
API_LIMITER.release();
}
}
static void queryDatabase() throws Exception {
Thread.sleep(50); // 模拟 DB 查询
}
static void callExternalApi() throws Exception {
Thread.sleep(100); // 模拟 API 调用
}
}
9.3 虚拟线程 + 生产者消费者模式
public class VirtualThreadProducerConsumer {
public static void main(String[] args) throws Exception {
BlockingQueue<String> queue = new LinkedBlockingQueue<>(1000);
try (var executor = Executors.newVirtualThreadPerTaskExecutor()) {
// 生产者:1 个虚拟线程
executor.submit(() -> {
for (int i = 0; i < 10000; i++) {
queue.put("item-" + i); // 队列满时阻塞 → yield
}
queue.put("POISON"); // 结束信号
});
// 消费者:多个虚拟线程并发消费
for (int i = 0; i < 50; i++) {
executor.submit(() -> {
while (true) {
String item = queue.take(); // 队列空时阻塞 → yield
if ("POISON".equals(item)) {
queue.put("POISON"); // 传递结束信号
break;
}
process(item);
}
});
}
}
}
static void process(String item) {
// 处理逻辑
}
}
十、虚拟线程与 IO 密集型任务
10.1 HTTP 服务器(每请求一线程)
public class HttpServerWithVirtualThreads {
public static void main(String[] args) throws Exception {
// JDK 内置 HTTP 服务器 + 虚拟线程
var server = com.sun.net.httpserver.HttpServer.create(
new java.net.InetSocketAddress(8080), 0);
// 使用虚拟线程处理每个请求
server.setExecutor(Executors.newVirtualThreadPerTaskExecutor());
server.createContext("/api/users", exchange -> {
try {
// 每个请求在独立虚拟线程中处理
// 可以随意阻塞(DB查询、RPC调用)而不影响吞吐量
String userId = extractUserId(exchange);
String userData = queryDatabase(userId); // 阻塞 DB 查询
String orderData = callOrderService(userId); // 阻塞 RPC 调用
byte[] response = (userData + orderData).getBytes();
exchange.sendResponseHeaders(200, response.length);
exchange.getResponseBody().write(response);
exchange.close();
} catch (Exception e) {
e.printStackTrace();
}
});
server.start();
System.out.println("Server started on port 8080 with virtual threads");
}
static String extractUserId(com.sun.net.httpserver.HttpExchange ex) {
return ex.getRequestURI().getQuery().replace("id=", "");
}
static String queryDatabase(String id) throws Exception {
Thread.sleep(50); return "{\"user\":\"" + id + "\"}";
}
static String callOrderService(String id) throws Exception {
Thread.sleep(30); return "{\"orders\":[]}";
}
}
10.2 并发 HTTP 客户端请求
public class ConcurrentHttpClient {
private static final HttpClient CLIENT = HttpClient.newBuilder()
.connectTimeout(Duration.ofSeconds(5))
.build();
public static void main(String[] args) throws Exception {
List<String> urls = List.of(
"https://httpbin.org/delay/1",
"https://httpbin.org/delay/2",
"https://httpbin.org/delay/1",
"https://httpbin.org/delay/3",
"https://httpbin.org/delay/1"
);
long start = System.nanoTime();
// 使用虚拟线程并发请求所有 URL
try (var executor = Executors.newVirtualThreadPerTaskExecutor()) {
List<Future<Integer>> futures = urls.stream()
.map(url -> executor.submit(() -> {
var request = HttpRequest.newBuilder()
.uri(URI.create(url))
.timeout(Duration.ofSeconds(10))
.build();
var response = CLIENT.send(request,
HttpResponse.BodyHandlers.ofString());
return response.statusCode();
}))
.toList();
// 等待所有结果
for (Future<Integer> f : futures) {
System.out.println("Status: " + f.get());
}
}
long elapsed = (System.nanoTime() - start) / 1_000_000;
System.out.println("Total time: " + elapsed + "ms");
// 总耗时 ≈ 最慢请求耗时(~3s),而非所有请求之和(~8s)
}
}
10.3 数据库批量操作
public class DatabaseBatchWithVirtualThreads {
// 数据库连接池仍然需要(连接是昂贵资源)
private static final int POOL_SIZE = 20;
private static final Semaphore CONN_LIMITER = new Semaphore(POOL_SIZE);
public static void main(String[] args) throws Exception {
List<Integer> userIds = IntStream.rangeClosed(1, 10000).boxed().toList();
long start = System.nanoTime();
try (var executor = Executors.newVirtualThreadPerTaskExecutor()) {
List<Future<String>> futures = userIds.stream()
.map(id -> executor.submit(() -> queryUser(id)))
.toList();
long successCount = futures.stream()
.map(f -> {
try { return f.get(); }
catch (Exception e) { return null; }
})
.filter(Objects::nonNull)
.count();
System.out.println("Success: " + successCount + "/10000");
}
long elapsed = (System.nanoTime() - start) / 1_000_000;
System.out.println("Elapsed: " + elapsed + "ms");
}
static String queryUser(int id) throws Exception {
// 用信号量限制并发数据库连接数
CONN_LIMITER.acquire();
try {
// 模拟数据库查询
Thread.sleep(10);
return "user-" + id;
} finally {
CONN_LIMITER.release();
}
}
}
10.4 文件批量处理
public class FileBatchProcessing {
public static void main(String[] args) throws Exception {
Path dir = Path.of("/data/logs");
List<Path> files;
try (var stream = Files.list(dir)) {
files = stream.filter(p -> p.toString().endsWith(".log")).toList();
}
System.out.println("Processing " + files.size() + " files");
try (var executor = Executors.newVirtualThreadPerTaskExecutor()) {
List<Future<Long>> futures = files.stream()
.map(file -> executor.submit(() -> processFile(file)))
.toList();
long totalLines = 0;
for (Future<Long> f : futures) {
totalLines += f.get();
}
System.out.println("Total lines processed: " + totalLines);
}
}
static long processFile(Path file) throws Exception {
// 文件 IO 阻塞时虚拟线程自动 yield
try (var reader = Files.newBufferedReader(file)) {
long count = 0;
String line;
while ((line = reader.readLine()) != null) {
if (line.contains("ERROR")) {
count++;
}
}
return count;
}
}
}
十一、虚拟线程监控
11.1 JFR(Java Flight Recorder)事件
// 虚拟线程相关 JFR 事件:
// jdk.VirtualThreadStart - 虚拟线程启动
// jdk.VirtualThreadEnd - 虚拟线程结束
// jdk.VirtualThreadPinned - Pin 事件(关键!)
// jdk.VirtualThreadSubmitFailed - 提交到调度器失败
public class JfrMonitoring {
public static void main(String[] args) throws Exception {
// 编程方式启动 JFR 记录
try (var rs = new RecordingStream()) {
// 监控 Pin 事件(阈值 20ms)
rs.enable("jdk.VirtualThreadPinned")
.withThreshold(Duration.ofMillis(20));
// 监控虚拟线程创建
rs.enable("jdk.VirtualThreadStart");
rs.onEvent("jdk.VirtualThreadPinned", event -> {
System.out.printf("[PIN] duration=%s, thread=%s%n",
event.getDuration(),
event.getThread().getJavaName());
});
rs.onEvent("jdk.VirtualThreadStart", event -> {
// 可统计虚拟线程创建速率
});
rs.startAsync();
// 运行业务逻辑...
runWorkload();
Thread.sleep(5000); // 等待事件处理
}
}
static void runWorkload() throws Exception {
try (var executor = Executors.newVirtualThreadPerTaskExecutor()) {
for (int i = 0; i < 1000; i++) {
executor.submit(() -> {
Thread.sleep(10);
});
}
}
}
}
11.2 jcmd 线程转储
# 获取包含虚拟线程的线程转储(JSON 格式)
jcmd <pid> Thread.dump_to_file -format=json threads.json
# 传统格式(也包含虚拟线程)
jcmd <pid> Thread.dump_to_file threads.txt
# 输出示例(JSON):
# {
# "threadContainers": [
# {
# "container": "ForkJoinPool-1",
# "parent": null,
# "owner": null,
# "threads": ["ForkJoinPool-1-worker-1", ...],
# "threadCount": 8
# },
# {
# "container": "VirtualThread[#25]/runnable@ForkJoinPool-1-worker-1",
# "threads": ["task-0"],
# "stackTrace": [...]
# }
# ]
# }
11.3 自定义监控指标
public class VirtualThreadMetrics {
private static final AtomicLong VT_CREATED = new AtomicLong();
private static final AtomicLong VT_COMPLETED = new AtomicLong();
private static final AtomicLong VT_FAILED = new AtomicLong();
// 包装 Executor 添加监控
public static ExecutorService monitoredExecutor() {
ExecutorService delegate = Executors.newVirtualThreadPerTaskExecutor();
return new MonitoredExecutorService(delegate);
}
static class MonitoredExecutorService extends AbstractExecutorService {
private final ExecutorService delegate;
MonitoredExecutorService(ExecutorService delegate) {
this.delegate = delegate;
}
@Override
public void execute(Runnable command) {
VT_CREATED.incrementAndGet();
delegate.execute(() -> {
try {
command.run();
VT_COMPLETED.incrementAndGet();
} catch (Exception e) {
VT_FAILED.incrementAndGet();
throw e;
}
});
}
// 获取指标
public static String metrics() {
return String.format("created=%d, completed=%d, failed=%d, active=%d",
VT_CREATED.get(), VT_COMPLETED.get(), VT_FAILED.get(),
VT_CREATED.get() - VT_COMPLETED.get() - VT_FAILED.get());
}
@Override public void shutdown() { delegate.shutdown(); }
@Override public List<Runnable> shutdownNow() { return delegate.shutdownNow(); }
@Override public boolean isShutdown() { return delegate.isShutdown(); }
@Override public boolean isTerminated() { return delegate.isTerminated(); }
@Override public boolean awaitTermination(long t, TimeUnit u)
throws InterruptedException { return delegate.awaitTermination(t, u); }
}
}
11.4 ThreadMXBean 注意事项
public class ThreadMXBeanNote {
public static void main(String[] args) {
ThreadMXBean mxBean = ManagementFactory.getThreadMXBean();
// 注意:ThreadMXBean.getThreadCount() 不包含虚拟线程
// 它只统计平台线程(包括 Carrier 线程)
System.out.println("Thread count (platform only): " + mxBean.getThreadCount());
// 虚拟线程对 ThreadMXBean 不可见(设计决策)
// 原因:虚拟线程可能有数百万个,枚举成本太高
// 替代方案:
// 1. JFR 事件统计
// 2. 自定义计数器(如上)
// 3. jcmd Thread.dump_to_file
}
}
十二、从线程池迁移到虚拟线程
12.1 迁移前后对比
// ===== 迁移前:传统线程池 =====
public class LegacyService {
private final ExecutorService pool = new ThreadPoolExecutor(
50, // corePoolSize
200, // maximumPoolSize
60L, TimeUnit.SECONDS, // keepAliveTime
new LinkedBlockingQueue<>(1000), // workQueue
new ThreadFactory() { // threadFactory
private final AtomicInteger counter = new AtomicInteger();
@Override
public Thread newThread(Runnable r) {
Thread t = new Thread(r, "worker-" + counter.incrementAndGet());
t.setDaemon(true);
return t;
}
},
new ThreadPoolExecutor.CallerRunsPolicy() // rejectedHandler
);
public void handleRequests(List<Request> requests) throws Exception {
List<Future<Response>> futures = new ArrayList<>();
for (Request req : requests) {
futures.add(pool.submit(() -> process(req)));
}
for (Future<Response> f : futures) {
f.get(); // 等待结果
}
}
private Response process(Request req) throws Exception {
// 阻塞 IO 操作
Thread.sleep(100);
return new Response("ok");
}
}
// ===== 迁移后:虚拟线程 =====
public class ModernService {
// 无需配置线程池参数!
public void handleRequests(List<Request> requests) throws Exception {
try (var executor = Executors.newVirtualThreadPerTaskExecutor()) {
List<Future<Response>> futures = new ArrayList<>();
for (Request req : requests) {
futures.add(executor.submit(() -> process(req)));
}
for (Future<Response> f : futures) {
f.get();
}
}
// close() 自动等待所有任务完成
}
private Response process(Request req) throws Exception {
Thread.sleep(100); // 阻塞时自动 yield,无需担心线程浪费
return new Response("ok");
}
}
12.2 迁移步骤详解
// Step 1: 替换 Executor 创建
// Before:
ExecutorService executor = Executors.newFixedThreadPool(100);
// After:
ExecutorService executor = Executors.newVirtualThreadPerTaskExecutor();
// Step 2: 移除线程池参数调优
// 不再需要:corePoolSize, maxPoolSize, keepAlive, workQueue, rejectedHandler
// Step 3: 替换 synchronized + 阻塞IO
// Before:
public synchronized String getData() throws Exception {
return httpClient.send(request, BodyHandlers.ofString()).body();
}
// After:
private final ReentrantLock lock = new ReentrantLock();
public String getData() throws Exception {
lock.lock();
try {
return httpClient.send(request, BodyHandlers.ofString()).body();
} finally {
lock.unlock();
}
}
// Step 4: 线程数限流 → Semaphore
// Before: 通过 fixedThreadPool(50) 限制并发
// After:
private static final Semaphore LIMITER = new Semaphore(50);
public void limitedTask() throws Exception {
LIMITER.acquire();
try {
// 业务逻辑
} finally {
LIMITER.release();
}
}
// Step 5: 大型 ThreadLocal → 参数传递或 ScopedValue
// Before:
private static final ThreadLocal<Connection> CONN = new ThreadLocal<>();
// After: 使用连接池 + 参数传递,或 ScopedValue
12.3 Spring Boot 集成
// Spring Boot 3.2+ 支持虚拟线程
// application.properties:
// spring.threads.virtual.enabled=true
// 或手动配置:
@Configuration
public class VirtualThreadConfig {
// Tomcat 使用虚拟线程处理请求
@Bean
public TomcatProtocolBuilderCustomizer<?> protocolCustomizer() {
return protocol -> {
protocol.setExecutor(Executors.newVirtualThreadPerTaskExecutor());
};
}
// 异步任务使用虚拟线程
@Bean
public AsyncTaskExecutor asyncExecutor() {
return new TaskExecutorAdapter(Executors.newVirtualThreadPerTaskExecutor());
}
// @Async 方法使用虚拟线程
@Bean(name = "applicationTaskExecutor")
public Executor applicationTaskExecutor() {
return Executors.newVirtualThreadPerTaskExecutor();
}
}
// 使用
@Service
public class OrderService {
@Async // 在虚拟线程中执行
public CompletableFuture<Order> processOrderAsync(String orderId) {
// 阻塞式调用多个微服务(虚拟线程中安全)
User user = userService.getUser(orderId);
Payment payment = paymentService.charge(orderId);
Inventory inv = inventoryService.reserve(orderId);
return CompletableFuture.completedFuture(new Order(user, payment, inv));
}
}
12.4 迁移验证清单
□ 1. JDK 版本 ≥ 21
□ 2. 开启 -Djdk.tracePinnedThreads=short 运行测试
□ 3. 检查所有 synchronized + 阻塞IO → 改为 ReentrantLock
□ 4. 检查第三方库(JDBC驱动、HTTP客户端)是否触发 Pin
□ 5. 大型 ThreadLocal → ScopedValue 或参数传递
□ 6. 线程池限流逻辑 → Semaphore
□ 7. 移除线程池参数配置
□ 8. 压测:对比迁移前后 QPS、P99 延迟、CPU 利用率
□ 9. 监控:JFR 观察 Pin 事件频率
□ 10. 灰度上线,观察 1-2 周
十三、虚拟线程适用场景与反模式
13.1 适用场景
// 场景 1:高并发 HTTP 服务(每请求一线程)
// 优势:简化编程模型,同步代码获得异步吞吐量
server.setExecutor(Executors.newVirtualThreadPerTaskExecutor());
// 场景 2:微服务编排(并发调用多个下游)
try (var executor = Executors.newVirtualThreadPerTaskExecutor()) {
Future<User> user = executor.submit(() -> userService.get(id));
Future<List<Order>> orders = executor.submit(() -> orderService.list(id));
Future<Profile> profile = executor.submit(() -> profileService.get(id));
// 三个调用并发执行,总耗时 = max(三者耗时)
}
// 场景 3:批量数据处理(IO 密集)
// 10万条记录并发写入数据库(Semaphore 限制连接数)
// 场景 4:消息队列并发消费
// 每条消息一个虚拟线程,轻松扩展消费并发度
// 场景 5:网关/代理(IO 转发为主)
// 大量并发连接,每个连接一个虚拟线程
13.2 不适用场景
// 反模式 1:CPU 密集型计算
// 虚拟线程无优势(Carrier 数 = CPU 核数),反而增加调度开销
// 正确:使用 ForkJoinPool 或平台线程池
try (var executor = Executors.newVirtualThreadPerTaskExecutor()) {
// 错误!CPU 密集任务不应使用虚拟线程
executor.submit(() -> computePrimeNumbers(1_000_000));
}
// 正确做法:
ForkJoinPool.commonPool().submit(() -> computePrimeNumbers(1_000_000));
// 反模式 2:池化虚拟线程
ExecutorService wrong = Executors.newFixedThreadPool(10,
Thread.ofVirtual().factory()); // 限制了并发,违背初衷
// 反模式 3:synchronized 包裹阻塞 IO
public synchronized void badMethod() throws Exception {
database.query("SELECT ..."); // Pin! Carrier 被占用
}
// 反模式 4:虚拟线程中存大型 ThreadLocal
private static final ThreadLocal<byte[]> BAD =
ThreadLocal.withInitial(() -> new byte[10 * 1024 * 1024]); // 10MB × 百万线程
// 反模式 5:用虚拟线程数量做限流
// 虚拟线程创建极廉价,"限制线程数"无意义
// 正确:用 Semaphore 限制对下游资源的并发访问
13.3 决策矩阵
| 任务类型 | 推荐方案 | 原因 |
|---|---|---|
| HTTP 请求处理 | 虚拟线程 | IO 密集,高并发 |
| 数据库 CRUD | 虚拟线程 + Semaphore | IO 密集,需限制连接数 |
| 微服务 RPC 调用 | 虚拟线程 | IO 密集,并发调用 |
| 文件读写批处理 | 虚拟线程 | IO 密集 |
| 图片/视频编码 | 平台线程池/ForkJoinPool | CPU 密集 |
| 机器学习推理 | 平台线程池 | CPU/GPU 密集 |
| 科学计算 | ForkJoinPool | CPU 密集,可分治 |
| 定时任务 | ScheduledThreadPool | 需要精确调度 |
十四、性能对比测试
14.1 吞吐量对比(IO 密集)
public class ThroughputBenchmark {
private static final int TASK_COUNT = 10_000;
private static final int IO_DELAY_MS = 100;
public static void main(String[] args) throws Exception {
// 预热
benchmarkPlatformThreads(200);
benchmarkVirtualThreads();
// 正式测试
System.out.println("=== IO 密集型任务吞吐量对比 ===");
System.out.println("任务数: " + TASK_COUNT + ", 每任务 IO 延迟: " + IO_DELAY_MS + "ms");
System.out.println();
// 平台线程池(200 线程)
long t1 = benchmarkPlatformThreads(200);
System.out.printf("平台线程池(200): %d ms, QPS=%.0f%n", t1, TASK_COUNT * 1000.0 / t1);
// 平台线程池(500 线程)
long t2 = benchmarkPlatformThreads(500);
System.out.printf("平台线程池(500): %d ms, QPS=%.0f%n", t2, TASK_COUNT * 1000.0 / t2);
// 虚拟线程
long t3 = benchmarkVirtualThreads();
System.out.printf("虚拟线程: %d ms, QPS=%.0f%n", t3, TASK_COUNT * 1000.0 / t3);
// 典型输出:
// 平台线程池(200): 5123 ms, QPS=1952
// 平台线程池(500): 2310 ms, QPS=4329
// 虚拟线程: 1089 ms, QPS=9183
}
static long benchmarkPlatformThreads(int poolSize) throws Exception {
ExecutorService pool = Executors.newFixedThreadPool(poolSize);
long start = System.nanoTime();
List<Future<?>> futures = new ArrayList<>();
for (int i = 0; i < TASK_COUNT; i++) {
futures.add(pool.submit(() -> {
try { Thread.sleep(IO_DELAY_MS); }
catch (InterruptedException e) { Thread.currentThread().interrupt(); }
}));
}
for (Future<?> f : futures) f.get();
long elapsed = (System.nanoTime() - start) / 1_000_000;
pool.shutdown();
return elapsed;
}
static long benchmarkVirtualThreads() throws Exception {
long start = System.nanoTime();
try (var executor = Executors.newVirtualThreadPerTaskExecutor()) {
List<Future<?>> futures = new ArrayList<>();
for (int i = 0; i < TASK_COUNT; i++) {
futures.add(executor.submit(() -> {
Thread.sleep(IO_DELAY_MS);
}));
}
for (Future<?> f : futures) f.get();
}
return (System.nanoTime() - start) / 1_000_000;
}
}
14.2 创建开销对比
public class CreationOverhead {
private static final int COUNT = 100_000;
public static void main(String[] args) throws Exception {
// 平台线程创建开销
long start1 = System.nanoTime();
for (int i = 0; i < 1000; i++) { // 只创建 1000 个(太多会 OOM)
Thread t = new Thread(() -> {});
t.start();
t.join();
}
long platformTime = (System.nanoTime() - start1) / 1_000_000;
System.out.println("1000 平台线程创建+执行: " + platformTime + " ms");
System.out.println("平均: " + (platformTime * 1000.0 / 1000) + " μs/线程");
// 虚拟线程创建开销
long start2 = System.nanoTime();
CountDownLatch latch = new CountDownLatch(COUNT);
for (int i = 0; i < COUNT; i++) {
Thread.startVirtualThread(latch::countDown);
}
latch.await();
long virtualTime = (System.nanoTime() - start2) / 1_000_000;
System.out.println(COUNT + " 虚拟线程创建+执行: " + virtualTime + " ms");
System.out.println("平均: " + (virtualTime * 1000.0 / COUNT) + " μs/线程");
// 典型输出:
// 1000 平台线程创建+执行: 850 ms → 平均 850 μs/线程
// 100000 虚拟线程创建+执行: 320 ms → 平均 3.2 μs/线程
}
}
14.3 内存占用对比
public class MemoryBenchmark {
public static void main(String[] args) throws Exception {
Runtime rt = Runtime.getRuntime();
// 基线内存
System.gc();
long baseline = rt.totalMemory() - rt.freeMemory();
// 创建 10 万个虚拟线程(保持存活)
int count = 100_000;
CountDownLatch ready = new CountDownLatch(count);
CountDownLatch hold = new CountDownLatch(1);
try (var executor = Executors.newVirtualThreadPerTaskExecutor()) {
for (int i = 0; i < count; i++) {
executor.submit(() -> {
ready.countDown();
hold.await(); // 保持线程存活
});
}
ready.await();
System.gc();
long withVT = rt.totalMemory() - rt.freeMemory();
long vtMemory = withVT - baseline;
System.out.println(count + " 虚拟线程内存占用: " + vtMemory / 1024 / 1024 + " MB");
System.out.println("平均每线程: " + vtMemory / count / 1024 + " KB");
// 典型输出: 约 150 MB 总计, 每线程约 1.5 KB
hold.countDown(); // 释放
}
// 对比:10万平台线程需要 ~100GB 栈内存(不可能)
System.out.println("10万平台线程理论栈内存: " + (count * 1024 * 1024L) / 1024 / 1024 / 1024 + " GB");
}
}
14.4 CPU 密集型对比
public class CpuBoundBenchmark {
public static void main(String[] args) throws Exception {
int taskCount = Runtime.getRuntime().availableProcessors() * 2;
System.out.println("=== CPU 密集型任务对比(" + taskCount + " 个任务)===");
// 平台线程池
long t1 = benchmarkCpu(Executors.newFixedThreadPool(taskCount));
System.out.println("平台线程池: " + t1 + " ms");
// 虚拟线程
long t2 = benchmarkCpu(Executors.newVirtualThreadPerTaskExecutor());
System.out.println("虚拟线程: " + t2 + " ms");
// ForkJoinPool
long t3 = benchmarkCpu(new ForkJoinPool());
System.out.println("ForkJoinPool: " + t3 + " ms");
// 典型输出(CPU 密集任务三者接近):
// 平台线程池: 2100 ms
// 虚拟线程: 2150 ms(略慢,有调度开销)
// ForkJoinPool: 2050 ms
}
static long benchmarkCpu(ExecutorService executor) throws Exception {
long start = System.nanoTime();
List<Future<Long>> futures = new ArrayList<>();
for (int i = 0; i < Runtime.getRuntime().availableProcessors() * 2; i++) {
futures.add(executor.submit(() -> {
long sum = 0;
for (long j = 0; j < 100_000_000L; j++) {
sum += j * j;
}
return sum;
}));
}
for (Future<Long> f : futures) f.get();
long elapsed = (System.nanoTime() - start) / 1_000_000;
executor.shutdown();
return elapsed;
}
}
十五、最佳实践
15.1 核心原则
/**
* 虚拟线程最佳实践总结
*/
public class VirtualThreadBestPractices {
// 原则 1:不要池化虚拟线程
// ✅ 正确
ExecutorService good = Executors.newVirtualThreadPerTaskExecutor();
// ❌ 错误
// ExecutorService bad = Executors.newFixedThreadPool(10, Thread.ofVirtual().factory());
// 原则 2:用 Semaphore 限制下游资源并发
private static final Semaphore DB_PERMITS = new Semaphore(20);
public void accessDatabase() throws Exception {
DB_PERMITS.acquire();
try {
// 数据库操作
} finally {
DB_PERMITS.release();
}
}
// 原则 3:避免 synchronized + 阻塞 IO
private final ReentrantLock lock = new ReentrantLock();
public void safeBlockingOperation() throws Exception {
lock.lock();
try {
// 阻塞 IO 安全
} finally {
lock.unlock();
}
}
// 原则 4:CPU 密集任务用平台线程
public void cpuIntensiveTask() {
ForkJoinPool.commonPool().submit(() -> {
// 计算密集型
});
}
// 原则 5:避免大型 ThreadLocal
// ✅ 使用 ScopedValue 或参数传递
private static final ScopedValue<String> CONTEXT = ScopedValue.newInstance();
}
15.2 生产环境配置建议
# JVM 启动参数(生产环境)
java \
# 载体线程数(默认=CPU核数,一般无需调整)
-Djdk.virtualThreadScheduler.parallelism=8 \
# Pin 时最大扩展载体数
-Djdk.virtualThreadScheduler.maxPoolSize=128 \
# JFR 持续记录(低开销)
-XX:StartFlightRecording=disk=true,maxsize=500m,maxage=24h,settings=profile \
-jar application.jar
# 开发/测试环境:开启 Pin 检测
java \
-Djdk.tracePinnedThreads=short \
-jar application.jar
15.3 异常处理模式
public class ExceptionHandling {
public static void main(String[] args) throws Exception {
try (var executor = Executors.newVirtualThreadPerTaskExecutor()) {
List<Future<Result>> futures = new ArrayList<>();
for (int i = 0; i < 100; i++) {
final int id = i;
futures.add(executor.submit(() -> {
try {
return doWork(id);
} catch (RetryableException e) {
// 重试逻辑
return retry(id, 3);
} catch (Exception e) {
// 记录日志,返回降级结果
log.error("Task {} failed", id, e);
return Result.fallback(id);
}
}));
}
// 收集结果,容忍部分失败
List<Result> results = futures.stream()
.map(f -> {
try { return f.get(5, TimeUnit.SECONDS); }
catch (TimeoutException e) { return Result.timeout(); }
catch (Exception e) { return Result.error(e); }
})
.toList();
}
}
static Result doWork(int id) throws Exception {
Thread.sleep(50);
if (id % 10 == 0) throw new RetryableException("transient");
return new Result("ok-" + id);
}
static Result retry(int id, int maxRetries) throws Exception {
for (int i = 0; i < maxRetries; i++) {
try {
Thread.sleep(100 * (i + 1)); // 退避
return doWork(id);
} catch (RetryableException e) {
if (i == maxRetries - 1) throw e;
}
}
throw new RuntimeException("unreachable");
}
record Result(String data) {
static Result fallback(int id) { return new Result("fallback-" + id); }
static Result timeout() { return new Result("timeout"); }
static Result error(Exception e) { return new Result("error: " + e.getMessage()); }
}
static class RetryableException extends Exception {
RetryableException(String msg) { super(msg); }
}
static class log {
static void error(String fmt, Object... args) {
System.err.println(fmt.formatted(args));
}
}
}
15.4 测试策略
public class VirtualThreadTesting {
// 单元测试:验证虚拟线程行为
@Test
void testVirtualThreadCreation() throws Exception {
Thread vt = Thread.startVirtualThread(() -> {
assertTrue(Thread.currentThread().isVirtual());
});
vt.join();
assertEquals(Thread.State.TERMINATED, vt.getState());
}
// 并发正确性测试
@Test
void testConcurrentAccess() throws Exception {
AtomicInteger counter = new AtomicInteger();
int taskCount = 10_000;
try (var executor = Executors.newVirtualThreadPerTaskExecutor()) {
List<Future<?>> futures = new ArrayList<>();
for (int i = 0; i < taskCount; i++) {
futures.add(executor.submit(counter::incrementAndGet));
}
for (Future<?> f : futures) f.get();
}
assertEquals(taskCount, counter.get());
}
// Pin 检测测试
@Test
void testNoPinning() throws Exception {
// 在 CI 中开启 -Djdk.tracePinnedThreads=short
// 如果有 Pin 事件,测试输出中会出现警告
var lock = new ReentrantLock();
try (var executor = Executors.newVirtualThreadPerTaskExecutor()) {
for (int i = 0; i < 100; i++) {
executor.submit(() -> {
lock.lock();
try {
Thread.sleep(10); // 不应 Pin
} catch (InterruptedException e) {
Thread.currentThread().interrupt();
} finally {
lock.unlock();
}
});
}
}
}
// 压力测试
@Test
void testHighConcurrency() throws Exception {
int concurrency = 100_000;
CountDownLatch latch = new CountDownLatch(concurrency);
try (var executor = Executors.newVirtualThreadPerTaskExecutor()) {
for (int i = 0; i < concurrency; i++) {
executor.submit(() -> {
try {
Thread.sleep(10);
} catch (InterruptedException e) {
Thread.currentThread().interrupt();
} finally {
latch.countDown();
}
});
}
assertTrue(latch.await(30, TimeUnit.SECONDS));
}
}
}
15.5 常见陷阱与解决方案
| 陷阱 | 症状 | 解决方案 |
|---|---|---|
| synchronized + 阻塞IO | 吞吐量骤降,Carrier 饥饿 | 改用 ReentrantLock |
| 大型 ThreadLocal | OOM 或 GC 压力大 | ScopedValue / 参数传递 |
| 池化虚拟线程 | 并发度被人为限制 | newVirtualThreadPerTaskExecutor |
| 无限流创建下游连接 | 下游服务过载 | Semaphore 限流 |
| CPU 密集任务用虚拟线程 | 无性能提升 | ForkJoinPool / 平台线程 |
| 依赖 Thread.interrupt 停止 | 虚拟线程中断语义相同但需正确处理 | 检查 isInterrupted |
| 第三方库内部 synchronized | 隐性 Pin | 升级库 / 联系维护者 |
| 虚拟线程中用 ThreadLocal 传连接 | 连接泄漏 | 连接池 + try-with-resources |
15.6 总结
虚拟线程核心价值:
┌─────────────────────────────────────────────────────────┐
│ 同步的编程模型 + 异步的吞吐量 + 极低的资源开销 │
└─────────────────────────────────────────────────────────┘
适用判断口诀:
IO 密集 → 虚拟线程
CPU 密集 → 平台线程 / ForkJoinPool
阻塞多 → 虚拟线程
计算多 → 平台线程
迁移优先级:
1. HTTP 服务层(收益最大)
2. 微服务调用层
3. 数据库访问层(配合 Semaphore)
4. 消息消费层
5. 批处理任务

openEuler 是由开放原子开源基金会孵化的全场景开源操作系统项目,面向数字基础设施四大核心场景(服务器、云计算、边缘计算、嵌入式),全面支持 ARM、x86、RISC-V、loongArch、PowerPC、SW-64 等多样性计算架构
更多推荐
所有评论(0)