Java 21引入了虚拟线程(Virtual Threads),这是Project Loom的核心特性。虚拟线程是轻量级线程,可以显著提高应用程序的并发处理能力,特别适合I/O密集型任务。
在Spring Boot 3.2+中,已经内置了对虚拟线程的支持,可以通过简单的配置启用。
在 application.yml 中已经配置:
|
1 2 3 4 |
spring: threads: virtual: enabled: true |
这个配置会自动:
@Async 方法启用虚拟线程执行器可以通过以下方式验证:
|
1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 16 |
@SpringBootTest public class VirtualThreadTest { @Test public void testVirtualThread() { Thread thread = Thread.ofVirtual().start(() -> { System.out.println("虚拟线程名称: " + Thread.currentThread().getName()); System.out.println("是否为虚拟线程: " + Thread.currentThread().isVirtual()); }); try { thread.join(); } catch (InterruptedException e) { Thread.currentThread().interrupt(); } } } |
配置类改造(推荐):
|
1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 16 17 18 19 20 21 22 23 24 25 26 27 28 29 30 31 32 33 34 35 36 37 38 39 40 41 42 43 44 45 46 47 48 49 50 51 52 53 54 55 56 57 58 59 60 61 62 63 64 65 66 67 |
/** * 异步线程池配置(虚拟线程版本) * 当 spring.threads.virtual.enabled=true 时使用虚拟线程 * * @author * @date 2024-10-31 */ @Slf4j @Configuration @EnableAsync public class VirtualThreadAsyncConfig implements AsyncConfigurer { /** * 虚拟线程异步执行器 * Spring Boot 3.2+ 会自动创建虚拟线程执行器,这里提供手动配置示例 */ @Bean(name = "virtualThreadExecutor") @ConditionalOnProperty(name = "spring.threads.virtual.enabled", havingValue = "true", matchIfMissing = false) public Executor virtualThreadExecutor() { return new TaskExecutorAdapter(Executors.newVirtualThreadPerTaskExecutor()); } /** * 默认异步执行器 * 如果启用虚拟线程,Spring Boot会自动使用虚拟线程执行器 */ @Override public Executor getAsyncExecutor() { // Spring Boot 3.2+ 会自动使用虚拟线程执行器(如果启用) // 如果需要手动指定,可以返回 virtualThreadExecutor() return Executors.newVirtualThreadPerTaskExecutor(); } @Override public AsyncUncaughtExceptionHandler getAsyncUncaughtExceptionHandler() { return new SimpleAsyncUncaughtExceptionHandler(); } } 使用示例: @Service @Slf4j public class ExampleService { /** * 使用默认虚拟线程执行器 */ @Async public CompletableFuture<String> asyncMethod1() { log.info("当前线程: {}, 是否为虚拟线程: {}", Thread.currentThread().getName(), Thread.currentThread().isVirtual()); // 执行异步任务 return CompletableFuture.completedFuture("完成"); } /** * 指定使用虚拟线程执行器 */ @Async("virtualThreadExecutor") public CompletableFuture<String> asyncMethod2() { log.info("当前线程: {}, 是否为虚拟线程: {}", Thread.currentThread().getName(), Thread.currentThread().isVirtual()); return CompletableFuture.completedFuture("完成"); } } |
Spring Boot 3.2+ 启用虚拟线程后,所有Web请求会自动使用虚拟线程处理,无需额外配置。
验证方式:
|
1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 16 17 |
@RestController @RequestMapping("/api/test") @Slf4j public class VirtualThreadTestController { @GetMapping("/virtual-thread") public Map<String, Object> testVirtualThread() { Thread currentThread = Thread.currentThread(); Map<String, Object> result = new HashMap<>(); result.put("threadName", currentThread.getName()); result.put("isVirtual", currentThread.isVirtual()); result.put("threadId", currentThread.threadId()); log.info("请求处理线程: {}, 是否为虚拟线程: {}", currentThread.getName(), currentThread.isVirtual()); return result; } } |
配置类改造:
|
1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 16 17 18 19 20 21 22 |
/** * 定时任务配置(虚拟线程版本) * * @author * @date 2024-10-31 */ @Slf4j @Configuration @EnableScheduling @ConditionalOnProperty(name = "spring.threads.virtual.enabled", havingValue = "true") public class VirtualThreadSchedulingConfig implements SchedulingConfigurer { @Override public void configureTasks(ScheduledTaskRegistrar taskRegistrar) { taskRegistrar.setScheduler(taskScheduler()); } @Bean public Executor taskScheduler() { return Executors.newVirtualThreadPerTaskExecutor(); } } |
使用示例:
|
1 2 3 4 5 6 7 8 9 10 11 12 13 14 |
@Component @Slf4j public class ScheduledTaskExample { /** * 定时任务会自动使用虚拟线程执行器 */ @Scheduled(fixedRate = 5000) public void scheduledTask() { log.info("定时任务执行 - 线程: {}, 是否为虚拟线程: {}", Thread.currentThread().getName(), Thread.currentThread().isVirtual()); } } |
|
1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 16 17 18 19 20 21 22 23 24 25 26 27 28 29 30 31 32 33 34 35 36 37 38 39 40 41 42 43 44 45 46 47 48 49 50 51 52 53 54 55 56 57 58 59 60 61 62 63 64 65 66 67 68 69 70 |
/** * 虚拟线程工具类(改进版) * * @author * @date 2024-10-31 */ @Slf4j @Component public class VirtualThreadUtils { private static final ExecutorService EXECUTOR = Executors.newVirtualThreadPerTaskExecutor(); /** * 执行虚拟线程任务 * * @param task 任务 * @return Future */ public static Future<?> exeVirtualThread(Runnable task) { return EXECUTOR.submit(() -> { try { log.debug("虚拟线程执行任务 - 线程: {}, 是否为虚拟线程: {}", Thread.currentThread().getName(), Thread.currentThread().isVirtual()); task.run(); } catch (Exception e) { log.error("虚拟线程执行任务异常", e); throw e; } }); } /** * 执行虚拟线程任务(带返回值) * * @param task 任务 * @param <T> 返回值类型 * @return Future */ public static <T> Future<T> exeVirtualThread(java.util.concurrent.Callable<T> task) { return EXECUTOR.submit(() -> { try { log.debug("虚拟线程执行任务 - 线程: {}, 是否为虚拟线程: {}", Thread.currentThread().getName(), Thread.currentThread().isVirtual()); return task.call(); } catch (Exception e) { log.error("虚拟线程执行任务异常", e); throw e; } }); } @PreDestroy public void destroy() { log.info("关闭虚拟线程执行器"); EXECUTOR.shutdown(); try { if (!EXECUTOR.awaitTermination(60, TimeUnit.SECONDS)) { EXECUTOR.shutdownNow(); if (!EXECUTOR.awaitTermination(60, TimeUnit.SECONDS)) { log.error("虚拟线程执行器未能正常关闭"); } } } catch (InterruptedException e) { EXECUTOR.shutdownNow(); Thread.currentThread().interrupt(); } } } |
|
1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 16 17 18 19 20 21 22 23 24 25 26 27 28 29 30 31 32 33 34 35 36 37 38 39 40 41 42 43 |
@Service @Slf4j public class CompletableFutureExample { /** * 使用虚拟线程执行器执行CompletableFuture */ public CompletableFuture<List<String>> processDataAsync() { Executor virtualExecutor = Executors.newVirtualThreadPerTaskExecutor(); CompletableFuture<List<String>> future = CompletableFuture .supplyAsync(() -> { log.info("异步任务执行 - 线程: {}, 是否为虚拟线程: {}", Thread.currentThread().getName(), Thread.currentThread().isVirtual()); return fetchData(); }, virtualExecutor) .thenApplyAsync(data -> { log.info("处理数据 - 线程: {}, 是否为虚拟线程: {}", Thread.currentThread().getName(), Thread.currentThread().isVirtual()); return processData(data); }, virtualExecutor); return future; } private List<String> fetchData() { // 模拟I/O操作 try { Thread.sleep(100); } catch (InterruptedException e) { Thread.currentThread().interrupt(); } return Arrays.asList("data1", "data2", "data3"); } private List<String> processData(List<String> data) { return data.stream() .map(String::toUpperCase) .collect(Collectors.toList()); } } |
|
1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 16 17 18 19 20 21 |
/** * RabbitMQ虚拟线程配置 * * @author 往事随风去 * @date 2024-10-31 */ @Slf4j @Configuration @ConditionalOnProperty(name = "spring.threads.virtual.enabled", havingValue = "true") public class RabbitMQVirtualThreadConfig { @Bean public SimpleRabbitListenerContainerFactory rabbitListenerContainerFactory( ConnectionFactory connectionFactory) { SimpleRabbitListenerContainerFactory factory = new SimpleRabbitListenerContainerFactory(); factory.setConnectionFactory(connectionFactory); // 使用虚拟线程执行器处理消息 factory.setTaskExecutor(new TaskExecutorAdapter(Executors.newVirtualThreadPerTaskExecutor())); return factory; } } |
对于MyBatis等数据库操作,通常不需要特别配置,因为:
虚拟线程适合:
虚拟线程不适合:
虚拟线程会频繁切换,ThreadLocal的使用需要注意:
|
1 2 3 4 5 |
// ❌ 不推荐:在虚拟线程中使用ThreadLocal存储大量数据 ThreadLocal<List<String>> threadLocal = new ThreadLocal<>(); // ✅ 推荐:使用ScopedValue(Java 21+)或谨慎使用ThreadLocal ScopedValue<String> scopedValue = ScopedValue.newInstance(); |
启用虚拟线程后,不需要配置线程池大小,虚拟线程会自动管理:
|
1 2 3 4 5 6 |
# ❌ 不需要配置这些(虚拟线程会自动管理) server: tomcat: threads: max: 400 # 虚拟线程模式下无效 min-spare: 100 # 虚拟线程模式下无效 |
虚拟线程的监控需要使用新的API:
|
1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 16 17 18 19 20 21 22 |
@Service @Slf4j public class VirtualThreadMonitor { @Scheduled(fixedRate = 60000) public void monitorVirtualThreads() { ThreadMXBean threadBean = ManagementFactory.getThreadMXBean(); long[] threadIds = threadBean.getAllThreadIds(); long virtualThreadCount = Arrays.stream(threadIds) .mapToObj(id -> { ThreadInfo info = threadBean.getThreadInfo(id); return info != null ? Thread.ofPlatform().getThreadGroup() .findThread(id) : null; }) .filter(Objects::nonNull) .filter(Thread::isVirtual) .count(); log.info("虚拟线程数量: {}", virtualThreadCount); } } |
|
1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 16 17 18 19 20 |
@Configuration public class HybridThreadConfig { /** * CPU密集型任务使用平台线程池 */ @Bean("cpuIntensiveExecutor") public Executor cpuIntensiveExecutor() { int processors = Runtime.getRuntime().availableProcessors(); return Executors.newFixedThreadPool(processors); } /** * I/O密集型任务使用虚拟线程 */ @Bean("ioIntensiveExecutor") public Executor ioIntensiveExecutor() { return Executors.newVirtualThreadPerTaskExecutor(); } } |
|
1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 16 17 18 19 20 21 |
@Service public class TaskService { @Autowired @Qualifier("cpuIntensiveExecutor") private Executor cpuExecutor; @Autowired @Qualifier("ioIntensiveExecutor") private Executor ioExecutor; public void processTask() { // I/O操作使用虚拟线程 CompletableFuture<String> data = CompletableFuture .supplyAsync(this::fetchData, ioExecutor); // CPU计算使用平台线程 CompletableFuture<String> result = data .thenApplyAsync(this::heavyComputation, cpuExecutor); } } |
|
1 2 3 4 |
spring: threads: virtual: enabled: true |
ScheduledTaskCofiguration 改造为使用虚拟线程如果出现问题,可以快速回滚:
|
1 2 3 4 |
spring: threads: virtual: enabled: false # 禁用虚拟线程,恢复传统线程池 |
虚拟线程在Spring Boot中的正确使用方式:
spring.threads.virtual.enabled=true@Async 自动使用虚拟线程SchedulingConfigurerExecutors.newVirtualThreadPerTaskExecutor()通过合理使用虚拟线程,可以显著提高I/O密集型应用的并发处理能力,减少线程资源消耗。