虚拟线程实战》解决的是”线程不再贵”,但没解决”谁来管这些线程”。
一个请求扇出到三个下游,其中一个失败了怎么办?超时了剩下的两个还在跑吗?调用方抛出异常时,那些已经没人等待的任务去哪了?
这些问题的答案,在 CompletableFutureExecutorService 里是”你自己记住”,在结构化并发里是”作用域负责”。这篇用同一份需求把三种写法跑一遍,用实测的耗时和线程残留数说话。
全部代码在 JDK 25 上编译运行,StructuredTaskScope 目前仍是预览 API,编译需要 --enable-preview

实验环境

1
2
3
4
$ java -version
java 25 2025-09-16 LTS
Java(TM) SE Runtime Environment (build 25+37-LTS-3491)
Java HotSpot(TM) 64-Bit Server VM (build 25+37-LTS-3491, mixed mode, sharing)

编译与运行都要带预览开关:

1
2
javac --enable-preview --release 25 Demo.java
java --enable-preview Demo

下面所有示例的公共部分是一个”下游调用”:它会记录自己是否还在运行,并如实报告自己被中断的时刻。

1
2
3
4
5
6
7
8
9
10
11
static final AtomicInteger running = new AtomicInteger();

static String fetch(String name, long ms) throws Exception {
running.incrementAndGet();
long t0 = System.nanoTime();
try { Thread.sleep(ms); return name + " 结果"; }
catch (InterruptedException e) {
System.out.printf(" <- %s 在 %dms 被中断%n", name, (System.nanoTime()-t0)/1_000_000);
throw e;
} finally { running.decrementAndGet(); }
}

running 这个计数器是关键道具:任务结束任务被取消在日志里长得一样,但”还剩几个在跑”骗不了人。

一、同一个需求,三种写法

需求固定:并发调用三个下游(各 3000ms),其中一个在 100ms 后失败,要求尽快发现失败并停止其它任务,超时时间 200ms。三份代码我都跑了。

1.1 写法一:CompletableFuture.allOf

1
2
3
4
5
6
7
8
9
10
11
12
var pool = Executors.newFixedThreadPool(3);
try {
var f = List.of(cf(pool, () -> fetch("下游A", 3000)),
cf(pool, () -> fetch("下游B", 3000)),
cf(pool, () -> fetch("下游C", 3000)));
try { CompletableFuture.allOf(f.toArray(CompletableFuture[]::new)).get(200, TimeUnit.MILLISECONDS); }
catch (TimeoutException e) {
System.out.println(" 200ms 超时抛出 TimeoutException(耗时 " + elapsed() + "ms)");
// 到这里你必须自己想起来取消这些 future,否则它们继续占用线程
}
probe(500); // 打印 500ms 后的 running 计数
} finally { Thread.sleep(2700); pool.shutdown(); }
1
2
3
4
=== A. CompletableFuture.allOf + 手动超时 ===
200ms 超时抛出 TimeoutException(耗时 209ms)
500ms 后仍在运行的下游任务数 = 3
(3 秒后线程池才自然结束)

超时在 209ms 抛出了,但三个任务一个都没停allOf 只负责”等”和”抛”,取消要靠调用方自己写 future.cancel(true),而且必须挨个 cancel、并且对方得响应中断。

如果不设超时直接 allOf(...).join(),它的语义是等全部结束。我另外测过一次”一路 100ms 就失败”的场景,join() 是在 3015ms 才抛出异常的——失败发生在第 100ms,但异常要等到所有任务结束才交给你。

1.2 写法二:ExecutorCompletionService 手写快速失败

想快速失败就得自己写循环,完成一个处理一个:

1
2
3
4
5
6
7
8
9
10
11
var ecs = new ExecutorCompletionService<String>(pool2);
var futures = new ArrayList<Future<String>>();
for (var t : tasks) futures.add(ecs.submit(t));
for (int i = 0; i < tasks.size(); i++) {
try { ecs.take().get(); }
catch (ExecutionException e) {
System.out.println(" 第 " + (i+1) + " 个完成的任务失败: " + e.getCause().getMessage() + ",耗时 " + elapsed() + "ms");
for (var f : futures) f.cancel(true); // 取消要自己写,且是「协作式」的
break;
}
}
1
2
3
4
5
=== B. ExecutorCompletionService 手写快速失败 ===
第 1 个完成的任务失败: 下游B 返回 500,耗时 109ms
<- 下游A 在 109ms 被中断
<- 下游C 在 109ms 被中断
300ms 后仍在运行的下游任务数 = 0

这段是”能跑对”的:109ms 就发现了失败,取消也生效了。代价是这份样板代码每个项目都要重写一遍,而且只要有一处漏了 cancel(true),或者某个任务不响应中断,残留就悄悄留下了。

1.3 写法三:结构化并发

1
2
3
4
5
6
7
try (var scope = StructuredTaskScope.open()) {
scope.fork(() -> fetch("下游A", 3000));
scope.fork(() -> fetch("下游B", 3000));
scope.fork(() -> fetch("下游C", 3000));
try { scope.join(); } catch (TimeoutException e) { ... }
probe(300);
}
1
2
3
4
5
6
7
=== C. 结构化并发 + withTimeout ===
<- 下游C 在 199ms 被中断
<- 下游A 在 200ms 被中断
<- 下游B 在 199ms 被中断
200ms 超时抛出 TimeoutException(耗时 209ms)
300ms 后仍在运行的下游任务数 = 0
出作用域,running = 0

三个下游同时在 200ms 被中断,try-with-resources 块结束时 running 归零。没有一行取消代码——这是 close() 的职责。

1.4 三者对照

CompletableFuture.allOf ExecutorCompletionService 结构化并发
超时后残留任务 3(需自己 cancel) 0(手写 cancel) 0(框架负责)
快速失败耗时 3015ms(等全部) 109ms 112ms
取消代码 调用方写 调用方写
作用域外还有孤儿任务 可能 可能 不可能(close() 保证)

二、JDK 25 里的 API 形态

JDK 25 的 StructuredTaskScope第五次预览(JEP 505)引入的新形态:类变成了 sealed interface,策略类被 Joiner 取代。下面是我用反射打出来的实际签名(javap 也能看到):

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
### StructuredTaskScope  (java.util.concurrent.StructuredTaskScope)
Object join[]
void close[]
StructuredTaskScope open[Joiner, Function]
StructuredTaskScope open[Joiner]
StructuredTaskScope open[]
Subtask fork[Callable]
Subtask fork[Runnable]
boolean isCancelled[]

### Joiner<T,R> 静态工厂
awaitAllSuccessfulOrThrow[]
allSuccessfulOrThrow[]
anySuccessfulResultOrThrow[]
awaitAll[]
allUntil[Predicate]

### Configuration
withThreadFactory[ThreadFactory]
withName[String]
withTimeout[Duration]

几个要点:

  • open() 是唯一入口,没有公开构造器。不传参数时默认策略是”全部成功,否则抛异常”(等价于老的 ShutdownOnFailure)。
  • fork() 返回 Subtask<T>,不是 FutureSubtask 只有三个方法:get()exception()state()
  • 超时不再是”策略”的一部分,而是作用域配置:open(joiner, cfg -> cfg.withTimeout(...))
  • 虚拟线程默认无名(实测日志里线程名是空的),排查问题时建议给作用域起名:cfg -> cfg.withName("下游调用")

2.1 这里有个编译坑

1
2
3
import static java.util.concurrent.StructuredTaskScope.*;
...
catch (TimeoutException e) { } // 编译错误:对 TimeoutException 的引用不明确
1
2
3
错误: 对TimeoutException的引用不明确
StructuredTaskScope 中的类 java.util.concurrent.StructuredTaskScope.TimeoutException
和 java.util.concurrent 中的类 java.util.concurrent.TimeoutException 都匹配

StructuredTaskScope 自带 TimeoutExceptionFailedException 两个异常类型,静态导入会和 java.util.concurrent 的同名类撞车,而且继承链不一样,抓错了就漏掉:

1
2
TimeoutException 继承链: class java.util.concurrent.StructuredTaskScope$TimeoutException -> class java.util.concurrent.TimeoutException
FailedException 继承链: class java.util.concurrent.StructuredTaskScope$FailedException -> class java.util.concurrent.ExecutionException

StructuredTaskScope.TimeoutExceptionjava.util.concurrent.TimeoutException 的子类,所以捕父类也能接到;但写代码时显式一点更省事:catch (StructuredTaskScope.TimeoutException e)

三、预置 Joiner 的四种语义

JDK 25 内置了五个静态工厂,对应四类常见需求:

工厂方法 语义 join() 的返回值 近似的老写法
awaitAllSuccessfulOrThrow() 全部成功才算成功,否则抛 void ShutdownOnFailure
allSuccessfulOrThrow() 同上,非阻塞获取结果 List<T>
anySuccessfulResultOrThrow() 任一成功即可,取最快的 T ShutdownOnSuccess
awaitAll() 等全部结束,逐个子任务判读 void
allUntil(Predicate) 自定义终止条件 void

3.1 取最快成功(fan-out 到多个镜像站)

1
2
3
4
5
6
7
try (var scope = StructuredTaskScope.open(Joiner.<String>anySuccessfulResultOrThrow())) {
scope.fork(() -> fetch("镜像站-北京", 800));
scope.fork(() -> fetch("镜像站-上海", 250));
scope.fork(() -> fetch("镜像站-广州", 600));
String r = scope.join(); // 直接拿到判优结果
System.out.printf("最快成功: %s(总耗时 %dms)%n", r, elapsed());
}
1
最快成功: 镜像站-上海 的结果(总耗时 272ms)

250ms 的任务赢了,总耗时 272ms(含判优开销),另外两个在 close() 时被取消——不需要写 invokeAny 那种”提交一批、等第一个”的模板代码。

3.2 部分成功也要全部收齐(聚合报表)

如果业务语义是”三路数据源各自独立,能拿到几路算几路”,用 awaitAll()

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
try (var scope = StructuredTaskScope.open(Joiner.<String>awaitAll(),
cfg -> cfg.withName("数据源"))) {
var db = scope.fork(() -> "DB: " + 42);
var api = scope.fork(() -> { throw new IllegalStateException("HTTP 503"); });
var hdfs = scope.fork(() -> "HDFS: 7");
scope.join();
for (var e : Map.of("DB", db, "API", api, "HDFS", hdfs).entrySet()) {
var st = e.getValue();
switch (st.state()) {
case SUCCESS -> System.out.printf(" %-5s 成功: %s%n", e.getKey(), st.get());
case FAILED -> System.out.printf(" %-5s 失败: %s%n", e.getKey(), st.exception().getMessage());
case UNAVAILABLE -> System.out.printf(" %-5s 被取消%n", e.getKey());
}
}
}
1
2
3
4
HDFS  成功: HDFS: 7
API 失败: HTTP 503
DB 成功: DB: 42
聚合完成:一路失败不影响其余两路的结果

Subtask.State 的三种取值(SUCCESS / FAILED / UNAVAILABLE)把失败和”被取消”分成两种状态UNAVAILABLE 表示”这个任务的结果不存在,因为它被取消了”。

四、失败与取消是怎么传播的

把 1.1 节里”谁都活下来”的场景换成结构化并发:

1
2
3
4
5
6
7
8
9
10
11
try (var scope = StructuredTaskScope.open()) {          // 默认 awaitAllSuccessfulOrThrow
var ok1 = scope.fork(() -> slow("下游A", 3000));
var bad = scope.fork(() -> { Thread.sleep(100); throw new IllegalStateException("下游B 返回 500"); });
var ok2 = scope.fork(() -> slow("下游C", 3000));
try { scope.join(); } catch (StructuredTaskScope.FailedException e) {
System.out.println("join() 抛出 FailedException,原因: " + e.getCause().getMessage());
}
System.out.println(" 下游A 状态=" + ok1.state());
System.out.println(" 下游B 状态=" + bad.state());
System.out.println(" 下游C 状态=" + ok2.state());
}
1
2
3
4
5
6
7
8
  [下游A] 在 108ms 处被中断,提前退出
[下游C] 在 108ms 处被中断,提前退出
join() 抛出 FailedException,原因: 下游B 返回 500
总耗时 137ms
下游A 状态=UNAVAILABLE
下游B 状态=FAILED
下游C 状态=UNAVAILABLE
作用域外再等 500ms:没有残留线程在跑

两个 3000ms 的任务,在 108ms 就被中断了,总耗时 137ms。同一场景下 allOf 的耗时是 3015ms——20 倍的差距,全部来自”谁负责取消”这个设计选择。

另外注意 UNAVAILABLE 的语义:它不只表示”被取消”,也包含”因为用了 anySuccessfulResultOrThrow,其余任务的结果无人认领”。想让结果一定存活,就不要用这类丢弃语义的 joiner。

五、与 ScopedValue 的配合

结构化并发和 ScopedValue 配合时有一个额外收益:父线程绑定的 ScopedValue 会被 fork() 出的子任务自动继承,而且不涉及拷贝。

1
2
3
4
5
6
7
8
9
10
11
static final ScopedValue<String> USER = ScopedValue.newInstance();

ScopedValue.where(USER, "alice").run(() -> {
try (var scope = StructuredTaskScope.open()) {
var a = scope.fork(() -> "子任务读到 " + USER.get());
var b = scope.fork(() -> "子任务读到 " + USER.get());
scope.join();
System.out.println(a.get());
System.out.println(b.get());
}
});
1
2
子任务1(虚拟线程) 读到 alice
子任务2(虚拟线程) 读到 alice

但这条继承规则只对结构化并发的子任务成立。我用七种方式在同一个作用域里试过读取绑定:

新建线程的方式 能否读到父作用域的绑定
StructuredTaskScope.fork(...)
Thread.ofVirtual().start(...) 不能
Thread.ofPlatform().start(...) 不能
Thread.startVirtualThread(...) 不能
Executors.newVirtualThreadPerTaskExecutor() 不能
Executors.newFixedThreadPool(1) 不能
ForkJoinPool.commonPool() 不能

官方文档把规则写得很直白:绑定由 fork 启动的所有线程继承;而像 ForkJoinPool 这类遗留线程管理类不支持继承,因为它们无法保证子线程会在父线程离开作用域之前结束。这也解释了为什么”请求上下文”这件事,正确做法是 ScopedValue + 结构化并发成对使用——单独的 ScopedValue 只能在同一线程内传递。

六、版本陷阱:JDK 21 到 27 的 API 漂移

结构化并发从 JDK 21 起一直是预览,每次预览都可能改签名,这是它最实际的成本。核对 openjdk.org 的原文,五个版本全为 Status: Closed / Delivered

JDK JEP 标题 API 形态
21 JEP 453 Structured Concurrency (Preview) class StructuredTaskScope<T> + 公开构造器 + ShutdownOnFailure / ShutdownOnSuccess 子类
24 JEP 499 (Fourth Preview) 与 21 相比 without change
25 JEP 505 (Fifth Preview) 改为 sealed interface + open() 工厂 + Joiner 接口,join() 返回值从 scope 变成 R
26 JEP 525 (Sixth Preview) 新增 onTimeout()allSuccessfulOrThrow() 改为返回 List
27 JEP 533 (Seventh Preview) Joiner<T,R,R_X> 增加第三类型参数;join() 改抛 ExecutionExceptiononTimeout() 改名 timeout()删除 awaitAll()

实际影响:

  • 网上搜到的示例大概率编译不过。JEP 505 自己的示例代码里还在用 new StructuredTaskScope.ShutdownOnFailure(),照抄到 JDK 25 上是编译错误——官方 JEP 正文的示例没有跟着最终 API 更新。
  • 跨版本升级要重读一遍签名。26 和 27 又改了三处,本文所有代码以 JDK 25 为准(我逐个编译运行过)。
  • 预览 API 的编译产物带预览标记:用 JDK 25 编译的 class 文件不能在 JDK 26 上直接运行,必须用 --release 25 明确目标版本。

七、什么时候不该用它

  • 只有一个异步任务:没有扇出就没有结构可言,CompletableFuture 或者直接同步调用更直白。
  • 任务之间需要长时间独立演进:作用域要求”子任务必须在作用域结束前结束”,这正是它的价值,但如果你真的需要一个跑几小时的批处理任务,那它不属于任何请求作用域,用普通的线程池更合适。
  • 需要跨请求复用线程上下文:作用域是”一次请求一个”,长生命周期的缓存/连接池不属于这里。
  • 团队还没上 JDK 25:21/24 的 API 形态和 25 不同,为 21 写代码等于为三年后的迁移埋坑。要么等它转正,要么直接以 25 为目标写。

总结

  • 结构化并发管的是**”谁来负责善后”**:try-with-resources 结束时,作用域保证所有子任务都已经结束或被取消
  • 实测差异(同一场景):allOf 等 3015ms 且三个任务全部残留;手写 ExecutorCompletionService 109ms 停止但需要约 15 行样板;结构化并发 112ms 停止且零取消代码
  • JDK 25 的入口只有 open(),策略全部由 Joiner 表达;超时属于 Configuration,不属于 joiner
  • Subtask.state()FAILEDUNAVAILABLE 分开,这是”部分成功”场景的关键表达能力
  • ScopedValue 绑定只被 fork() 的子任务继承,其余七种新建线程方式都读不到
  • 预览 API 的漂移是真实成本:21→27 之间一共派生了七个预览 JEP(453/462/480/499/505/525/533),其中 25、26、27 三次是实质性的 API 变更;写代码前先确认目标 JDK 版本的签名

参考资料

系列索引:Java 系列,语言特性与运行时的长文集