使用虚拟线程运行并行流:Java中ForkJoinWorkerThread的VirtualThread等效方案
截至JDK 20,不存在与ForkJoinWorkerThread直接等效的虚拟线程,ForkJoinPool仅能管理ForkJoinWorkerThread实例,因此无法直接将虚拟线程接入ForkJoinPool来运行并行流。针对你的场景(阻塞式分页数据源+CPU密集型后续操作),可以通过以下两种思路实现需求:
一、拆分任务:虚拟线程处理阻塞数据源,并行流处理CPU密集操作
你的场景中,数据源等待是阻塞型操作(适合虚拟线程节省OS线程资源),后续操作是CPU密集型(适合用ForkJoinPool的并行流提升效率),可以将两者拆分:
用虚拟线程并发获取分页数据
利用Executors.newVirtualThreadPerTaskExecutor()创建虚拟线程池,并发请求所有分页的数据源,避免阻塞OS线程:// 假设getPageData(int page)是阻塞的分页数据获取方法 List<Data> allData = Executors.newVirtualThreadPerTaskExecutor().invokeAll( IntStream.range(1, totalPages + 1) .mapToObj(page -> (Callable<List<Data>>) () -> getPageData(page)) .toList() ).stream() .map(future -> { try { return future.get(); } catch (Exception e) { throw new RuntimeException(e); } }) .flatMap(List::stream) .toList();用并行流处理CPU密集操作
将汇总后的全量数据传入并行流,利用ForkJoinPool处理CPU密集型的中间/终端操作:allData.parallelStream() .map(this::cpuIntensiveTransform) .collect(Collectors.toList());
二、自定义工作窃取式任务调度(基于虚拟线程)
如果希望完全用虚拟线程实现工作窃取模式,可以借助ForkJoinPool作为虚拟线程的调度器,手动拆分任务并提交:
ForkJoinPool本身是工作窃取的调度器,虚拟线程可以被调度到ForkJoinPool的工作线程上执行。你可以将CPU密集型任务拆分为多个子任务,提交到ForkJoinPool,同时用虚拟线程处理阻塞的数据源获取:
// 用ForkJoinPool作为虚拟线程的调度器 ExecutorService scheduler = ForkJoinPool.commonPool(); try (var executor = Executors.newThreadPerTaskExecutor(scheduler)) { // 虚拟线程处理阻塞数据源 Future<List<Data>> dataFuture = executor.submit(() -> getPageData(1)); List<Data> data = dataFuture.get(); // 拆分CPU密集任务到ForkJoinPool(工作窃取) ForkJoinTask<List<Result>> task = ForkJoinTask.adapt(() -> data.stream() .map(this::cpuIntensiveTransform) .toList() ); task.fork(); List<Result> results = task.join(); }
需要注意的是,CPU密集型任务用虚拟线程并没有性能优势——虚拟线程在执行CPU密集操作时会一直绑定到OS线程,无法被调度让出,反而可能增加调度开销。因此更推荐第一种拆分方案,让虚拟线程专注处理阻塞场景,并行流处理CPU密集场景。
内容的提问来源于stack exchange,提问作者Warren Nocos

