Java 21 虚拟线程的并发控制

用一个并行等待两个下游的最小服务,对照平台线程与虚拟线程,观察入口排队、许可等待、总截止时间和取消,再用 JDK 21 的真实 pinning 栈定位持锁阻塞。

把应用部署到土耳其|BRNCHOST · 土耳其 VDS
云服务器,积分可续期|雨云 · 国内外节点 · 积分兑换权益
低价年付,大流量 VPS|RackNerd · SSD 存储 · 1Gbps 端口
香港轻量,搭个小站|晚安云 · 香港云服务器
香港 VPS,大带宽可选|野草云 · BGP 直连
大陆访问,精品线路|搬瓦工 · CN2 GIA / CTGNet 套餐
资料归档,交给 AI 整理|WorkBuddy · 本地文件处理
建站起步,先看应用镜像|腾讯云 · 轻量应用服务器
CN2 GIA,中国方向优化|DMIT · Premium 网络
双 ISP 原生住宅 IP|丽萨主机 · 美国 9929 精品线路
高频 CPU,多地部署|Evoxt · 云服务器 · 每周异地备份
京东云轻量云主机:129元/年,新人专享,限购1台

一个聚合接口要同时读账户资料和订单摘要。两个下游各花八十毫秒,接口拿齐结果后返回。把入口线程池从八个线程换成虚拟线程,二十四个请求可以同时开始;如果两个下游各只给四个连接,多出来的请求仍要等待。升级评估要回答的具体问题是:等待发生在哪里,谁承担等待期间的内存,预算耗尽后还有哪些任务在运行。

本文用本地实验跟踪这条路径。运行环境为 Homebrew OpenJDK 21.0.12.1,构建号同为 21.0.12.1,实验日期为二〇二六年九月八日。下游用可中断的八十毫秒等待模拟,许可用信号量模拟。这里没有真实数据库连接,没有网络传输,也没有线上吞吐成绩。这样的安排便于单独观察调度与排队,同时把结论限制在代码实际覆盖的范围内。

下文内嵌完整源码与编译命令,读者可以在独立目录运行。本文表格与诊断片段来自本轮实际输出;原始证据另由编辑交付留存,不依赖线上附件。

先把线程池承担的工作拆开

平台线程版本有两个执行器。入口执行器有八个线程,负责接收请求、提交两个子任务和等待结果;下游执行器有十六个线程,执行模拟调用。这里刻意分开两个池。若父任务和子任务都放进同一个八线程池,八个父任务可能占满工作线程,随后一起等待还在队列中的子任务。那样测到的是线程饥饿,与虚拟线程的等待成本混在一起,无法解释比较结果。

虚拟线程版本把这两个执行器都替换为每任务创建一个虚拟线程的执行器。父任务仍提交两个子任务,读取相同的结果,沿用相同的许可控制和截止时间。代码没有建立虚拟线程池,也没有把所有业务状态放进线程本地缓存。你可以把它理解为同一请求编排换了执行资源,业务依赖关系保持一致。

ExecutorService requests = virtual
    ? Executors.newVirtualThreadPerTaskExecutor()
    : Executors.newFixedThreadPool(8);
ExecutorService io = virtual
    ? Executors.newVirtualThreadPerTaskExecutor()
    : Executors.newFixedThreadPool(16);

一个虚拟线程在适合卸载的阻塞点等待时,运行时可以让承载线程执行其他任务。计算密集的循环仍要消耗处理器时间;加密、压缩和大对象转换不会因为线程变轻就变成零成本。Oracle 的 Java 21 虚拟线程指南说明了执行器与阻塞行为,证据目录保存了此次读取的官方页面,便于区分版本文档与本地测量。

固定八线程池并非平台线程能力上限。你可以把它扩到二十四个入口线程,再配足子任务线程,缩小这里的时间差,但要承担更多平台线程及其管理成本。本实验选择八线程作为可解释的受限入口,不能拿这个配置证明虚拟线程在任意负载下更快。严肃的容量评估还需要给平台方案一个合理调优后的基线,并记录两种方案的资源用量。

五组运行把排队位置分开记录

每组瞬时提交二十四个请求。前两组给两个下游各二十四个许可,总预算五秒;中间两组把许可缩到四个,其余保持一致。最后一组只允许十二个请求进入,总预算缩到一百八十毫秒。预算从提交时计算,包含入口队列停留时间。

场景 成功 / 到期 / 拒绝 整批耗时毫秒 入口等待累计毫秒 许可等待累计毫秒 下游 A/B 峰值
平台线程,宽许可 24 / 0 / 0 261.7 2037.8 0.1 8 / 8
虚拟线程,宽许可 24 / 0 / 0 93.0 26.7 0.1 24 / 24
平台线程,各四许可 24 / 0 / 0 503.5 3484.6 2965.3 4 / 4
虚拟线程,各四许可 24 / 0 / 0 502.1 3.4 9979.5 4 / 4
虚拟线程,入口十二,总预算一百八十毫秒 6 / 6 / 12 182.8 0.4 2044.8 4 / 4

这是一次执行的原始观测,不是多轮基准的统计结论。时间值会随调度、机器负载和预热变化。最后一行的成功数也不是协议规定的常数,两个下游各自竞争许可,获得 A 许可的请求未必同时获得 B 许可,不能按波次数简单推断必定成功八个。测试断言只检查请求计数守恒、许可峰值没有越界、结束后许可全部归还。

入口等待从提交任务计到入口工作线程开始运行。许可等待从尝试获取许可计到成功获得许可,只累计成功获取的等待;两个下游的值相加。累计值可以超过整批墙钟时间,因为多条线程在同一时段等待。它既不是单个请求的平均延迟,也不是尾延迟。这里没有记录完整延迟分布,因此不能据此填写百分之九十九分位。

宽许可场景下,平台版本分批处理入口工作,虚拟版本让更多请求同时进入等待。四许可场景下,两种方案都需要大致六轮下游服务,整批耗时接近;虚拟版本入口等待减少,但许可等待累计上升。你若只看入口线程队列长度,会误以为拥堵已经消除。真实系统的相应证据是连接池等待、下游活动请求、超时来源和请求持有的上下文大小。

信号量只模拟连接额度,不包含建连、连接健康检查、驱动锁竞争和连接回收。把这里的许可等待直接写成某个数据库驱动的实测连接池等待,会夸大证据。接入真实驱动后,应同时记录申请连接与执行查询的起止时间,区分额度不足、慢查询和失效连接重建。三种问题需要不同的调整,增加入口线程对后两种问题未必有帮助。

入口和下游分别决定接受多少工作

下游许可必须由应用级对象持有。示例对 A、B 各创建一个 Gate,二十四个请求共享它们。如果在请求方法里创建新的四许可信号量,每个请求都能拿到自己的许可,全局下游并发仍可能达到二十四。代码的外形还像有限流,限制的却只是单个请求的局部动作。

入口另有一个信号量,在提交任务之前使用不等待的 tryAcquire()。本次突发提交观测到十二个请求进入、十二个请求被拒绝;入口同时占用不超过十二。提交循环与任务执行并发进行,若提交线程暂停到已有任务释放许可,后面的请求仍可能获得位置,因此拒绝数量不是任意调度下的固定值。入口位置在请求收尾时释放。下游信号量限制某一种资源的同时使用量,入口信号量限制服务愿意接纳的在途请求,两者不能互相代替。

这也解释了虚拟线程迁移中的内存风险。一个等待连接的请求还可能持有请求正文、解析对象、追踪上下文和结果缓冲。线程本身变轻后,这些对象仍然存在。若允许无限制创建任务,长尾下游会让在途对象积累。应按入口的内存预算确定上限,并在网关、队列消费者或 HTTP 接入层定义拒绝与重试规则。

示例按并发数量限流,没有实现每秒请求速率限制。四个许可在下游变慢时会降低完成速率,在下游变快时会提高完成速率。若供应商合同限定每秒调用次数,还要增加令牌桶一类的速率控制。把并发上限和速率上限写成同一个配置项,会在响应时间变化时破坏原来的容量假设。

生产环境也不应让每个副本各自占满供应商的全部额度。假设某下游可分给这项服务四十个连接,十个副本每个四个只是一个静态起点;滚动发布同时保留新旧副本时,实际副本数可能增加。扩缩容和重试都要纳入额度分配。本实验只有一个进程,所以表中的峰值只证明进程内约束。

用一个截止时间覆盖排队与执行

常见错误是在获取许可时等四十毫秒,拿到后给网络调用一百八十毫秒,再给等待另一个子任务一百八十毫秒。每一段都声明“有超时”,总请求却可能超过用户预算。示例在请求提交时用单调时钟计算截止时间,此后只传递同一个值,每个阻塞点重新计算剩余时间。

long deadline = submitted + TimeUnit.MILLISECONDS.toNanos(budget);
long remaining = deadline - System.nanoTime();
if (remaining <= 0) throw new TimeoutException("request");
return future.get(remaining, TimeUnit.NANOSECONDS);

Gate 先用剩余时间等待许可,拿到许可后再次计算。模拟操作睡眠取八十毫秒和剩余时间的较小值,醒来后检查是否到期。父任务等待第一个结果后,再次以相同截止时间等待第二个结果,不给第二次等待重新发放完整预算。这样入口等待、许可等待和实际操作共享同一份请求时间。

System.nanoTime() 适合计算本进程内经过的时间,不能把它的绝对数值传给远端作为跨机器时间戳。实际 HTTP 调用可以把剩余时长转换为客户端超时设置,并按照协议向下游传播预算;远端再用自己的单调时钟计时。还要为序列化、响应写出和清理预留时间,否则业务任务在截止点返回,用户仍可能拿不到及时的响应。

表中预算组整批耗时一百八十二点八毫秒,没有精确停在一百八十。线程唤醒和调度有开销,收尾也占时间。示例实现的是协作式预算检查,不提供实时系统的硬期限保证。测试还把捕获到的异常统一计入 expired,便于最小实验展示;接入真实系统时必须拆分许可超时、执行异常、中断与调用方取消,避免把业务错误计成预算耗尽。

父任务在 finally 中取消两个子任务。Future.cancel(true) 发出中断请求,模拟的睡眠和许可等待会响应中断,所以许可最终归还。驱动或本地调用若忽略中断,子任务可能继续运行。执行器的关闭还会等待任务结束;把它放在请求级 try 块里并不能凭空取得硬超时。此处执行器属于整轮实验,真实服务通常在生命周期结束时统一关闭。

一次远端写入在调用方取消前可能已经提交。中断 Java 线程不会撤回另一台机器上的事务。聚合只读数据比较适合先验证虚拟线程;付款、发券之类的写入还需要明确的业务键和结果查询。这里没有用取消实验推导分布式写入安全性。

从真实栈定位 Java 21 的持锁阻塞

pinning 探针另起一个 JVM:启动虚拟线程,在 synchronized 块中睡眠一百二十毫秒,然后等待线程退出。与对照实验分开运行可以防止诊断输出扰动耗时。本次使用的诊断命令如下,随后列出实际捕获的关键栈。

java -Djdk.tracePinnedThreads=full -cp . Lab pin

本次输出包含以下片段,行号对应证据中的源码:

VirtualThread[#20]/runnable@ForkJoinPool-1-worker-1 reason:MONITOR
    java.base/java.lang.VirtualThread.parkNanos(VirtualThread.java:635)
    java.base/java.lang.VirtualThread.sleepNanos(VirtualThread.java:807)
    java.base/java.lang.Thread.sleep(Thread.java:507)
    Lab.lambda$main$3(Lab.java:47) <== monitors:1

先读原因 MONITOR,再沿阻塞栈找到业务帧。这里线程在持有一个监视器时进入睡眠,JDK 21 报告了对应的固定行为。它证明这个构建中的该路径能触发 pinning,没有证明应用里的某段真实 HTTP 客户端代码也发生同样问题。诊断实际服务时,应保留业务帧、阻塞操作类型和发生频率,避免只截一行虚拟线程编号。

处理顺序通常是缩小锁覆盖范围。读取共享状态时持锁,复制出调用所需的不可变值,释放锁后执行网络请求,再在必要时核对版本并提交结果。这样做可能引入“读出后状态已变化”的竞争,必须通过版本检查或状态机约束解决,不能机械地把网络调用移出去后宣告完成。

某些业务需要在等待期间保持互斥,可以评估 ReentrantLock 等机制,但仍要检查临界区吞吐与取消逻辑。换锁只处理承载线程占用的一部分问题;单个全局锁包住八十毫秒调用,即使线程可以卸载,请求仍然串行通过。短小的内存临界区没有同样的阻塞等待成本,不应依据一次探针把所有 synchronized 全局替换。

JFR 也可用于收集相关事件并按持续时间、栈和频次分析。本轮只运行了诊断参数,没有采集 JFR,因此不提供虚构的事件数量或图表。对长期运行服务,团队可以在受控环境增加 JFR 记录,再确认采样与阈值足以覆盖怀疑的路径。

后续 JDK 已改变这一领域。OpenJDK 团队关于 JEP 491 的说明明确指向 JDK 24 的同步改进。不能拿新版本里某个监视器场景的表现解释本轮 Java 21 的栈,也不能用该改进推断所有本地调用导致的固定都已经消失。升级时记录发行商和补丁构建号,比只记“用了 LTS”更有用。

语言特性与迁移决定分开验收

Java 21 的结构化并发仍是预览 API。若采用 StructuredTaskScope,编译和运行都要启用预览,并使用 Java 21 对应签名;本实验没有依赖它。版本语义以 Java SE 21 API为准,未来版本的教程不能直接当成兼容实现。

javac --release 21 --enable-preview Example.java
java --enable-preview Example

这两行是预览代码的参数说明,未在本轮执行。结构化作用域能让父子任务的生命周期更容易表达,但团队仍要为连接额度和远端写入结果负责。若业务不接受预览特性,先用正式执行器接口完成迁移仍是可行路线;预览类型可以留在一个小适配层里,避免扩散到公共业务接口。

模式匹配适合用封闭类型区分找到、缺失和拒绝访问等返回状态。不过,本轮的问题集中在线程资源与请求预算,没有为了展示语法增加另一条主线。实际升级可以把语言重构另开变更,分别审核空值约定与分支覆盖,减少运行时变化和业务语义变化相互遮蔽的机会。

把实验参数换成接口的容量约定

假设账户服务可提供四个并发位置,订单服务可提供八个,且两边处理时间仍按八十毫秒估算。账户侧理想上限约为每秒五十次调用,订单侧约为一百次。每个聚合请求都要访问账户,因此聚合接口不能据此期待每秒一百个完成。这是由假设推导的上限,没有计入网络、排队和失败,也没有在本轮执行,不应写进实测表。

如果产品允许账户失败时仍返回订单摘要,可以改变聚合契约,给失败区域展示独立状态。这样的降级会改变成功的定义,原先“两个结果都齐全”与降级后的“至少一个结果可用”不能放在同一吞吐口径中比较。你需要在响应结构中表达缺失原因,让调用方能区分权限拒绝、预算耗尽和业务上没有数据,再根据用户用途决定是否缓存降级结果。

若两个结果必须同时有效,则提前获知一边失败后取消另一边可以减少浪费。当前示例按固定顺序等待两个 Future,另一边先失败时不一定马上结束父任务。它仍受共同截止时间约束,但没有实现失败即收尾的完成队列。正式实现可以用完成服务或合适的任务作用域传播失败,并测试已完成、等待许可和正在调用三种子任务状态下的取消行为。

入口上限也可以从内存约束出发估算。假设一次在途请求平均持有两百千字节对象,允许一万个请求等待,仅这些对象就约占两千兆字节,尚未计算运行时和缓存。这里的对象大小是示意值,需要通过真实请求堆分析替换;算式的用途是让团队看见无限等待的成本,不能把假设当成虚拟线程自身的内存测量。

升级发布时先保留平台线程路径的可切换配置,再在小流量路径接入虚拟执行器。观察到下游等待增长而完成量没有增加,可以回退线程策略或收紧入口,并保留当前运行时继续排查。若同时升级字节码目标,回退到旧 JDK 还需要旧构建产物,开关线程模式无法解决字节码兼容性。把运行时回退与并发策略回退分别准备,故障处理时才有明确操作。

本轮没有测量处理器密集任务、堆峰值、长时间稳定性或真实驱动取消。下一轮最有价值的补充应由接口瓶颈决定:怀疑等待转移就接连接池指标,怀疑上下文膨胀就记录堆,怀疑取消无效就故障注入下游。不要为了让升级报告看起来全面而加入与当前决策无关的跑分。

根据这次结果,若服务入口有大量阻塞等待、下游仍有可用额度,虚拟线程值得继续接入真实客户端验证;若下游额度已经用满,先明确等待上限和拒绝策略。迁移后的验收至少应能解释一次请求在哪里开始等待、预算如何耗尽、取消后资源何时归还。表中的四许可场景已经给出一个具体提醒:入口队列接近清空时,用户完成这一批请求所需的时间仍可能保持原样。

完整实验与执行方法

下面是本次运行的完整 Lab.java。使用 JDK 21,在独立空目录中编译和执行;不需要数据库、Docker 或第三方依赖。四组对照和预算组共用同一程序,pinning 探针另起 JVM。终端中的时间会随机器负载改变,核对许可上限与计数守恒,不要求复现表格中的小数。

java -version
javac --release 21 Lab.java
java -cp . Lab
java -Djdk.tracePinnedThreads=full -cp . Lab pin
import java.util.*;
import java.util.concurrent.*;
import java.util.concurrent.atomic.*;
public class Lab {
 static class Gate {
  final Semaphore sem; final AtomicInteger active=new AtomicInteger(), peak=new AtomicInteger();
  final LongAdder waits=new LongAdder();
  Gate(int n){sem=new Semaphore(n);}
  String call(long deadline) throws Exception {
   long begin=System.nanoTime(), left=deadline-begin;
   if(left<=0 || !sem.tryAcquire(left,TimeUnit.NANOSECONDS))throw new TimeoutException("permit");
   waits.add(System.nanoTime()-begin);
   peak.accumulateAndGet(active.incrementAndGet(),Math::max);
   try {left=deadline-System.nanoTime();if(left<=0)throw new TimeoutException("before IO");
    long delay=TimeUnit.MILLISECONDS.toNanos(80);
    TimeUnit.NANOSECONDS.sleep(Math.min(left,delay));
    if(System.nanoTime()>=deadline)throw new TimeoutException("IO deadline");return "ok";
   }finally{active.decrementAndGet();sem.release();}
  }
 }
 static String await(Future<String> f,long d)throws Exception{
  long left=d-System.nanoTime();if(left<=0)throw new TimeoutException("request");return f.get(left,TimeUnit.NANOSECONDS);
 }
 static void run(String name,boolean virtual,int permits,int ingress,long budget)throws Exception{
  Gate a=new Gate(permits),b=new Gate(permits);Semaphore admission=new Semaphore(ingress);
  LongAdder queues=new LongAdder();AtomicInteger ok=new AtomicInteger(),expired=new AtomicInteger(),rejected=new AtomicInteger();
  long start=System.nanoTime();List<Future<?>> roots=new ArrayList<>();
  try(ExecutorService io=virtual?Executors.newVirtualThreadPerTaskExecutor():Executors.newFixedThreadPool(16);
      ExecutorService requests=virtual?Executors.newVirtualThreadPerTaskExecutor():Executors.newFixedThreadPool(8)){
   for(int i=0;i<24;i++){
    long submitted=System.nanoTime(),deadline=submitted+TimeUnit.MILLISECONDS.toNanos(budget);
    if(!admission.tryAcquire()){rejected.incrementAndGet();continue;}
    roots.add(requests.submit(()->{Future<String> fa=null,fb=null;
     queues.add(System.nanoTime()-submitted);
     try{if(System.nanoTime()>=deadline)throw new TimeoutException("queue");
      fa=io.submit(()->a.call(deadline));fb=io.submit(()->b.call(deadline));await(fa,deadline);await(fb,deadline);ok.incrementAndGet();
     }catch(Exception e){expired.incrementAndGet();}
     finally{if(fa!=null)fa.cancel(true);if(fb!=null)fb.cancel(true);admission.release();}
    }));
   }
   for(Future<?> f:roots)f.get();
  }
  System.out.printf(Locale.ROOT,"%s ok=%d expired=%d rejected=%d elapsed_ms=%.1f queue_total_ms=%.1f permit_wait_total_ms=%.1f peakA=%d peakB=%d availableA=%d availableB=%d%n",name,ok.get(),expired.get(),rejected.get(),(System.nanoTime()-start)/1e6,queues.sum()/1e6,(a.waits.sum()+b.waits.sum())/1e6,a.peak.get(),b.peak.get(),a.sem.availablePermits(),b.sem.availablePermits());
  if(a.peak.get()>permits||b.peak.get()>permits||a.sem.availablePermits()!=permits||b.sem.availablePermits()!=permits||ok.get()+expired.get()+rejected.get()!=24)throw new AssertionError();
 }
 public static void main(String[] args)throws Exception{
  if(args.length>0){Object lock=new Object();Thread t=Thread.startVirtualThread(()->{synchronized(lock){try{Thread.sleep(120);}catch(InterruptedException e){Thread.currentThread().interrupt();}}});t.join();System.out.println("pinning probe complete");return;}
  run("platform-wide",false,24,24,5000);run("virtual-wide",true,24,24,5000);
  run("platform-gated",false,4,24,5000);run("virtual-gated",true,4,24,5000);
  run("virtual-budget",true,4,12,180);
 }
}