本文目录导读:

我来为您提供几个Java 21虚拟线程的实际应用案例,从简单到复杂。
案例1:高并发Web请求处理
场景:模拟大量并发HTTP请求
import java.time.Duration;
import java.time.Instant;
import java.util.concurrent.*;
import java.util.concurrent.atomic.AtomicInteger;
import java.util.stream.IntStream;
public class VirtualThreadWebDemo {
public static void main(String[] args) throws InterruptedException {
System.out.println("=== 虚拟线程 vs 平台线程 性能对比 ===");
// 测试1: 高并发任务处理
testHighConcurrency();
// 测试2: 混合IO任务
testMixedIOOperations();
}
// 模拟处理大量并发任务
static void testHighConcurrency() throws InterruptedException {
int taskCount = 10_000;
// 使用虚拟线程
Instant vtStart = Instant.now();
try (var executor = Executors.newVirtualThreadPerTaskExecutor()) {
CountDownLatch latch = new CountDownLatch(taskCount);
for (int i = 0; i < taskCount; i++) {
final int taskId = i;
executor.submit(() -> {
try {
simulateIOWork(taskId);
} finally {
latch.countDown();
}
});
}
latch.await();
}
Instant vtEnd = Instant.now();
System.out.printf("虚拟线程处理 %d 个任务耗时: %d ms%n",
taskCount, Duration.between(vtStart, vtEnd).toMillis());
// 使用固定线程池
Instant ptStart = Instant.now();
try (var executor = Executors.newFixedThreadPool(100)) {
CountDownLatch latch = new CountDownLatch(taskCount);
for (int i = 0; i < taskCount; i++) {
final int taskId = i;
executor.submit(() -> {
try {
simulateIOWork(taskId);
} finally {
latch.countDown();
}
});
}
latch.await();
}
Instant ptEnd = Instant.now();
System.out.printf("平台线程(100个)处理 %d 个任务耗时: %d ms%n",
taskCount, Duration.between(ptStart, ptEnd).toMillis());
}
// 模拟IO操作(如数据库查询、REST调用等)
static void simulateIOWork(int taskId) {
try {
Thread.sleep(100); // 模拟100ms的IO延迟
// System.out.printf("任务 %d 完成%n", taskId);
} catch (InterruptedException e) {
Thread.currentThread().interrupt();
}
}
// 测试混合操作
static void testMixedIOOperations() throws InterruptedException {
System.out.println("\n=== 混合IO操作测试 ===");
int taskCount = 5_000;
AtomicInteger completedTasks = new AtomicInteger(0);
// 使用虚拟线程处理混合工作负载
try (var executor = Executors.newVirtualThreadPerTaskExecutor()) {
CountDownLatch latch = new CountDownLatch(taskCount);
for (int i = 0; i < taskCount; i++) {
final int taskId = i;
executor.submit(() -> {
try {
// 模拟不同的操作类型
switch (taskId % 3) {
case 0 -> {
// 数据库查询
simulateDatabaseQuery();
}
case 1 -> {
// HTTP调用
simulateHttpCall();
}
case 2 -> {
// 文件操作
simulateFileOperation();
}
}
completedTasks.incrementAndGet();
} catch (Exception e) {
System.err.println("任务失败: " + e.getMessage());
} finally {
latch.countDown();
}
});
}
// 等待所有任务完成,带超时
if (!latch.await(30, TimeUnit.SECONDS)) {
System.err.println("任务超时!");
}
System.out.printf("成功完成任务: %d/%d%n", completedTasks.get(), taskCount);
}
}
static void simulateDatabaseQuery() throws InterruptedException {
Thread.sleep(50 + (int)(Math.random() * 100));
}
static void simulateHttpCall() throws InterruptedException {
Thread.sleep(100 + (int)(Math.random() * 200));
}
static void simulateFileOperation() throws InterruptedException {
Thread.sleep(30 + (int)(Math.random() * 50));
}
}
案例2:并发HTTP服务器
import com.sun.net.httpserver.HttpServer;
import com.sun.net.httpserver.HttpExchange;
import java.io.IOException;
import java.io.OutputStream;
import java.net.InetSocketAddress;
import java.util.concurrent.Executors;
public class VirtualThreadHttpServer {
public static void main(String[] args) throws IOException {
// 创建HttpServer
HttpServer server = HttpServer.create(new InetSocketAddress(8080), 0);
// 创建上下文
server.createContext("/api/data", exchange -> {
handleRequest(exchange);
});
// 使用虚拟线程作为执行器
server.setExecutor(Executors.newVirtualThreadPerTaskExecutor());
server.start();
System.out.println("服务器启动在端口 8080");
System.out.println("测试: curl http://localhost:8080/api/data");
}
static void handleRequest(HttpExchange exchange) throws IOException {
try {
// 模拟处理延迟
Thread.sleep(100);
String response = """
{
"status": "success",
"message": "使用虚拟线程处理",
"thread": "%s",
"time": "%s"
}
""".formatted(
Thread.currentThread().getName(),
java.time.LocalDateTime.now()
);
exchange.getResponseHeaders().set("Content-Type", "application/json");
exchange.sendResponseHeaders(200, response.getBytes().length);
try (OutputStream os = exchange.getResponseBody()) {
os.write(response.getBytes());
}
} catch (InterruptedException e) {
Thread.currentThread().interrupt();
exchange.sendResponseHeaders(500, -1);
}
}
}
案例3:并发数据爬虫
import java.io.IOException;
import java.net.URI;
import java.net.http.HttpClient;
import java.net.http.HttpRequest;
import java.net.http.HttpResponse;
import java.util.*;
import java.util.concurrent.*;
import java.util.concurrent.atomic.AtomicInteger;
public class VirtualThreadCrawler {
private static final List<String> urls = List.of(
"https://example.com/page1",
"https://example.com/page2",
"https://example.com/page3",
// ... 更多URL
);
public static void main(String[] args) throws InterruptedException {
System.out.println("=== 并发网页爬虫(使用虚拟线程) ===");
// 方案1: 使用虚拟线程
crawlWithVirtualThreads();
// 方案2: 使用平台线程对比
// crawlWithPlatformThreads();
}
static void crawlWithVirtualThreads() throws InterruptedException {
AtomicInteger completedCount = new AtomicInteger(0);
AtomicInteger failedCount = new AtomicInteger(0);
List<String> results = new CopyOnWriteArrayList<>();
try (var executor = Executors.newVirtualThreadPerTaskExecutor()) {
CountDownLatch latch = new CountDownLatch(urls.size());
for (String url : urls) {
executor.submit(() -> {
try {
String content = fetchUrl(url);
results.add(content.substring(0, Math.min(100, content.length())));
completedCount.incrementAndGet();
System.out.println("抓取成功: " + url);
} catch (Exception e) {
failedCount.incrementAndGet();
System.err.println("抓取失败: " + url + " - " + e.getMessage());
} finally {
latch.countDown();
}
});
}
latch.await();
System.out.printf("%n统计: 成功 %d, 失败 %d%n",
completedCount.get(), failedCount.get());
}
}
static String fetchUrl(String url) throws IOException, InterruptedException {
HttpClient client = HttpClient.newBuilder()
.connectTimeout(Duration.ofSeconds(10))
.build();
HttpRequest request = HttpRequest.newBuilder()
.uri(URI.create(url))
.timeout(Duration.ofSeconds(30))
.build();
HttpResponse<String> response = client.send(request,
HttpResponse.BodyHandlers.ofString());
// 简单模拟内容处理
Thread.sleep(100);
return response.body();
}
}
案例4:生产级应用示例
import java.util.concurrent.*;
import java.util.concurrent.atomic.*;
import java.util.function.Supplier;
public class ProductionVirtualThreadDemo {
// 配置
private static final int MAX_CONCURRENT_TASKS = 100_000;
private static final int TASK_TIMEOUT_SECONDS = 10;
public static void main(String[] args) {
System.out.println("=== 生产级虚拟线程应用 ===");
// 创建虚拟线程池
ExecutorService virtualThreadPool = Executors.newVirtualThreadPerTaskExecutor();
// 自定义线程工厂
ExecutorService customPool = Executors.newThreadPerTaskExecutor(
Thread.ofVirtual()
.name("virtual-worker-", 0)
.factory()
);
try {
// 演示批量任务处理
processBatchTasks(virtualThreadPool);
// 演示并发控制
demoConcurrencyControl(virtualThreadPool);
// 演示结构化并发
demoStructuredConcurrency();
} finally {
virtualThreadPool.shutdown();
customPool.shutdown();
}
}
// 批量任务处理
static void processBatchTasks(ExecutorService executor) {
System.out.println("\n--- 批量任务处理 ---");
int batchSize = 100;
CompletableFuture<Void>[] futures = new CompletableFuture[batchSize];
for (int i = 0; i < batchSize; i++) {
final int taskId = i;
futures[i] = CompletableFuture.runAsync(() -> {
try {
processTask(taskId);
} catch (Exception e) {
System.err.println("任务 " + taskId + " 失败: " + e.getMessage());
}
}, executor);
}
// 等待所有任务完成
CompletableFuture.allOf(futures).join();
System.out.println("所有 " + batchSize + " 个任务完成");
}
// 演示并发控制
static void demoConcurrencyControl(ExecutorService executor) {
System.out.println("\n--- 并发控制演示 ---");
Semaphore semaphore = new Semaphore(10); // 最多10个并发
for (int i = 0; i < 20; i++) {
final int taskId = i;
CompletableFuture.runAsync(() -> {
try {
semaphore.acquire();
try {
System.out.printf("任务 %d 开始 (并发: %d)%n",
taskId, 20 - semaphore.availablePermits());
Thread.sleep(100);
} finally {
semaphore.release();
}
} catch (InterruptedException e) {
Thread.currentThread().interrupt();
}
}, executor).join();
}
System.out.println("并发控制测试完成");
}
// 演示结构化并发(Java 21新特性)
static void demoStructuredConcurrency() {
System.out.println("\n--- 结构化并发演示 ---");
try (var scope = new StructuredTaskScope.ShutdownOnFailure()) {
Future<String> user = scope.fork(() -> fetchUserData());
Future<String> order = scope.fork(() -> fetchOrderData());
Future<String> settings = scope.fork(() -> fetchSettings());
scope.join();
scope.throwIfFailed();
String result = String.format("用户数据: %s, 订单: %s, 设置: %s",
user.resultNow(), order.resultNow(), settings.resultNow());
System.out.println("结构化并发结果: " + result);
} catch (Exception e) {
System.err.println("结构化并发失败: " + e.getMessage());
}
}
static String fetchUserData() throws InterruptedException {
Thread.sleep(100);
return "用户A";
}
static String fetchOrderData() throws InterruptedException {
Thread.sleep(150);
return "订单#123";
}
static String fetchSettings() throws InterruptedException {
Thread.sleep(200);
return "默认设置";
}
// 模拟处理任务
static void processTask(int taskId) throws InterruptedException {
// 模拟复杂操作
Thread.sleep(10);
if (taskId % 20 == 0) {
System.out.printf("任务 %d 处理完成(特殊任务)%n", taskId);
}
}
}
案例5:性能监控工具
import java.time.Duration;
import java.time.Instant;
import java.util.concurrent.*;
import java.util.concurrent.atomic.*;
public class VirtualThreadMonitor {
private static final AtomicLong activeTasks = new AtomicLong(0);
private static final AtomicLong completedTasks = new AtomicLong(0);
private static final ConcurrentLinkedQueue<String> taskLog = new ConcurrentLinkedQueue<>();
public static void main(String[] args) throws InterruptedException {
System.out.println("=== 虚拟线程监控工具 ===");
// 启动监控
startMonitor();
// 模拟工作负载
testWorkLoad();
// 等待监控输出
Thread.sleep(2000);
}
static void startMonitor() {
ScheduledExecutorService monitor = Executors.newSingleThreadScheduledExecutor();
monitor.scheduleAtFixedRate(() -> {
System.out.printf("活跃任务: %d, 完成任务: %d%n",
activeTasks.get(), completedTasks.get());
// 展示线程状态
ThreadMXBean threadMXBean = ManagementFactory.getThreadMXBean();
System.out.printf("虚拟线程数: %d, 平台线程数: %d%n%n",
threadMXBean.getThreadCount(),
Thread.getAllStackTraces().keySet().size());
}, 0, 2, TimeUnit.SECONDS);
}
static void testWorkLoad() throws InterruptedException {
int totalTasks = 500;
CountDownLatch latch = new CountDownLatch(totalTasks);
try (var executor = Executors.newVirtualThreadPerTaskExecutor()) {
for (int i = 0; i < totalTasks; i++) {
final int taskId = i;
executor.submit(() -> {
activeTasks.incrementAndGet();
try {
// 模拟工作
Thread.sleep(100);
completedTasks.incrementAndGet();
} catch (InterruptedException e) {
Thread.currentThread().interrupt();
} finally {
activeTasks.decrementAndGet();
latch.countDown();
}
});
}
latch.await();
System.out.println("所有任务完成");
}
}
}
-
创建方式:使用
Executors.newVirtualThreadPerTaskExecutor()或Thread.ofVirtual() -
性能优势:能够轻松处理数万个并发任务,内存占用远低于平台线程
-
注意事项:
- 阻塞操作(如
Thread.sleep()、IO)在虚拟线程中不会阻塞底层OS线程 - 避免在虚拟线程中使用
synchronized,使用ReentrantLock替代 - 虚拟线程适合IO密集型任务,不适合CPU密集型计算
- 阻塞操作(如
-
生产环境建议:
- 使用
StructuredTaskScope进行结构化并发 - 设置适当的超时和重试机制
- 监控虚拟线程池的使用情况
- 使用
这些案例展示了Java 21虚拟线程的实际应用场景,从简单的并发任务到复杂的生产级应用,虚拟线程特别适合高并发的IO密集型应用,如Web服务、微服务、数据库连接等场景。