Java21虚拟线程正式案例

wen java案例 1

本文目录导读:

Java21虚拟线程正式案例

  1. 案例1:高并发Web请求处理
  2. 案例2:并发HTTP服务器
  3. 案例3:并发数据爬虫
  4. 案例4:生产级应用示例
  5. 案例5:性能监控工具

我来为您提供几个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("所有任务完成");
        }
    }
}
  1. 创建方式:使用Executors.newVirtualThreadPerTaskExecutor()Thread.ofVirtual()

  2. 性能优势:能够轻松处理数万个并发任务,内存占用远低于平台线程

  3. 注意事项

    • 阻塞操作(如Thread.sleep()、IO)在虚拟线程中不会阻塞底层OS线程
    • 避免在虚拟线程中使用synchronized,使用ReentrantLock替代
    • 虚拟线程适合IO密集型任务,不适合CPU密集型计算
  4. 生产环境建议

    • 使用StructuredTaskScope进行结构化并发
    • 设置适当的超时和重试机制
    • 监控虚拟线程池的使用情况

这些案例展示了Java 21虚拟线程的实际应用场景,从简单的并发任务到复杂的生产级应用,虚拟线程特别适合高并发的IO密集型应用,如Web服务、微服务、数据库连接等场景。

抱歉,评论功能暂时关闭!