Java——结构化并发实战(StructuredTaskScope)
《虚拟线程实战》解决的是”线程不再贵”,但没解决”谁来管这些线程”。
一个请求扇出到三个下游,其中一个失败了怎么办?超时了剩下的两个还在跑吗?调用方抛出异常时,那些已经没人等待的任务去哪了?
这些问题的答案,在CompletableFuture、ExecutorService里是”你自己记住”,在结构化并发里是”作用域负责”。这篇用同一份需求把三种写法跑一遍,用实测的耗时和线程残留数说话。
全部代码在 JDK 25 上编译运行,StructuredTaskScope目前仍是预览 API,编译需要--enable-preview。
实验环境
1 | $ java -version |
编译与运行都要带预览开关:
1 | javac --enable-preview --release 25 Demo.java |
下面所有示例的公共部分是一个”下游调用”:它会记录自己是否还在运行,并如实报告自己被中断的时刻。
1 | static final AtomicInteger running = new AtomicInteger(); |
running 这个计数器是关键道具:任务结束和任务被取消在日志里长得一样,但”还剩几个在跑”骗不了人。
一、同一个需求,三种写法
需求固定:并发调用三个下游(各 3000ms),其中一个在 100ms 后失败,要求尽快发现失败并停止其它任务,超时时间 200ms。三份代码我都跑了。
1.1 写法一:CompletableFuture.allOf
1 | var pool = Executors.newFixedThreadPool(3); |
1 | === A. CompletableFuture.allOf + 手动超时 === |
超时在 209ms 抛出了,但三个任务一个都没停:allOf 只负责”等”和”抛”,取消要靠调用方自己写 future.cancel(true),而且必须挨个 cancel、并且对方得响应中断。
如果不设超时直接 allOf(...).join(),它的语义是等全部结束。我另外测过一次”一路 100ms 就失败”的场景,join() 是在 3015ms 才抛出异常的——失败发生在第 100ms,但异常要等到所有任务结束才交给你。
1.2 写法二:ExecutorCompletionService 手写快速失败
想快速失败就得自己写循环,完成一个处理一个:
1 | var ecs = new ExecutorCompletionService<String>(pool2); |
1 | === B. ExecutorCompletionService 手写快速失败 === |
这段是”能跑对”的:109ms 就发现了失败,取消也生效了。代价是这份样板代码每个项目都要重写一遍,而且只要有一处漏了 cancel(true),或者某个任务不响应中断,残留就悄悄留下了。
1.3 写法三:结构化并发
1 | try (var scope = StructuredTaskScope.open()) { |
1 | === C. 结构化并发 + withTimeout === |
三个下游同时在 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 | ### StructuredTaskScope (java.util.concurrent.StructuredTaskScope) |
几个要点:
open()是唯一入口,没有公开构造器。不传参数时默认策略是”全部成功,否则抛异常”(等价于老的ShutdownOnFailure)。fork()返回Subtask<T>,不是Future。Subtask只有三个方法:get()、exception()、state()。- 超时不再是”策略”的一部分,而是作用域配置:
open(joiner, cfg -> cfg.withTimeout(...))。 - 虚拟线程默认无名(实测日志里线程名是空的),排查问题时建议给作用域起名:
cfg -> cfg.withName("下游调用")。
2.1 这里有个编译坑
1 | import static java.util.concurrent.StructuredTaskScope.*; |
1 | 错误: 对TimeoutException的引用不明确 |
StructuredTaskScope 自带 TimeoutException 和 FailedException 两个异常类型,静态导入会和 java.util.concurrent 的同名类撞车,而且继承链不一样,抓错了就漏掉:
1 | TimeoutException 继承链: class java.util.concurrent.StructuredTaskScope$TimeoutException -> class java.util.concurrent.TimeoutException |
StructuredTaskScope.TimeoutException 是 java.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 | try (var scope = StructuredTaskScope.open(Joiner.<String>anySuccessfulResultOrThrow())) { |
1 | 最快成功: 镜像站-上海 的结果(总耗时 272ms) |
250ms 的任务赢了,总耗时 272ms(含判优开销),另外两个在 close() 时被取消——不需要写 invokeAny 那种”提交一批、等第一个”的模板代码。
3.2 部分成功也要全部收齐(聚合报表)
如果业务语义是”三路数据源各自独立,能拿到几路算几路”,用 awaitAll():
1 | try (var scope = StructuredTaskScope.open(Joiner.<String>awaitAll(), |
1 | HDFS 成功: HDFS: 7 |
Subtask.State 的三种取值(SUCCESS / FAILED / UNAVAILABLE)把失败和”被取消”分成两种状态:UNAVAILABLE 表示”这个任务的结果不存在,因为它被取消了”。
四、失败与取消是怎么传播的
把 1.1 节里”谁都活下来”的场景换成结构化并发:
1 | try (var scope = StructuredTaskScope.open()) { // 默认 awaitAllSuccessfulOrThrow |
1 | [下游A] 在 108ms 处被中断,提前退出 |
两个 3000ms 的任务,在 108ms 就被中断了,总耗时 137ms。同一场景下 allOf 的耗时是 3015ms——20 倍的差距,全部来自”谁负责取消”这个设计选择。
另外注意 UNAVAILABLE 的语义:它不只表示”被取消”,也包含”因为用了 anySuccessfulResultOrThrow,其余任务的结果无人认领”。想让结果一定存活,就不要用这类丢弃语义的 joiner。
五、与 ScopedValue 的配合
结构化并发和 ScopedValue 配合时有一个额外收益:父线程绑定的 ScopedValue 会被 fork() 出的子任务自动继承,而且不涉及拷贝。
1 | static final ScopedValue<String> USER = ScopedValue.newInstance(); |
1 | 子任务1(虚拟线程) 读到 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() 改抛 ExecutionException;onTimeout() 改名 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 且三个任务全部残留;手写ExecutorCompletionService109ms 停止但需要约 15 行样板;结构化并发 112ms 停止且零取消代码 - JDK 25 的入口只有
open(),策略全部由Joiner表达;超时属于Configuration,不属于 joiner Subtask.state()把FAILED与UNAVAILABLE分开,这是”部分成功”场景的关键表达能力ScopedValue绑定只被fork()的子任务继承,其余七种新建线程方式都读不到- 预览 API 的漂移是真实成本:21→27 之间一共派生了七个预览 JEP(453/462/480/499/505/525/533),其中 25、26、27 三次是实质性的 API 变更;写代码前先确认目标 JDK 版本的签名
参考资料
- JEP 505: Structured Concurrency (Fifth Preview)
- JEP 525: Structured Concurrency (Sixth Preview)
- JEP 533: Structured Concurrency (Seventh Preview)
- JEP 453: Structured Concurrency (Preview)
- JDK 25 API:StructuredTaskScope
系列索引:Java 系列,语言特性与运行时的长文集







