结构化并发(Java 21+)
提出问题
你写过这样的代码吗?提交一组任务到 ExecutorService,拿到一堆 Future,然后逐条 get() 阻塞等待。如果某个任务异常退出,其他任务还在跑,你得手动取消——但常常忘了,导致线程泄漏、资源空转。更糟的是,子任务崩溃后主线程可能一直卡在 future.get() 上,错误信息被吞没。
这些问题根源于传统线程模型里任务的生命周期没有跟作用域绑定。线程是"发射后不管"的——你启动了一个线程,就失去了对它的控制。Java 21 引入的结构化并发(JEP 453,StructuredTaskScope)正是为了解决这个问题:把并发任务当作一个代码块,块内所有子任务的生命周期不会超过这个块,要么全部完成,要么全部取消。
结构化并发解决什么
结构化并发不是 Java 独创的。早在 2016 年,Martin Sústrik 在 C 语言中提出了结构化并发的概念,随后被 Go 的 goroutine 生命周期管理、Python 的 Trio 库、Swift 的 async/await 逐步采纳。Java 21 的 JEP 453 是这一思想在 JVM 生态的正式落地。
核心承诺只有一条:子任务的生命周期 ≤ 父作用域的生命周期。翻译成人话就是:你开了一个 try 块,块里 fork 出去的所有子线程,在 try 块结束之前,要么执行完,要么被强制取消,不会有任何一个漏在外面。
为什么传统模式做不到
传统 ExecutorService.submit() 返回的 Future 是一个"句柄",你拿着它,但 JVM 不保证你一定会调用 cancel()。下面的代码在线上真实发生过:
ExecutorService es = Executors.newFixedThreadPool(10);
for (int i = 0; i < 100; i++) {
int finalI = i;
Future<String> f = es.submit(() -> fetchData(finalI));
// 忘记把 f 存到 list 里
// 这些任务跑到地老天荒,但没人能取消它们
}更隐蔽的情况是异常路径。子任务抛了 RuntimeException,主线程调 f.get() 拿到 ExecutionException 包裹的异常,但其他正在跑的 Future 你根本不知道,也没人取消它们。这些任务会一直运行到完成,浪费 CPU 和 IO 资源,甚至导致线程池里的活跃线程一直占着坑,让新任务排队。
StructuredTaskScope 的核心模型
StructuredTaskScope 是 java.util.concurrent 包下的一个密封类(sealed class),它的设计思路很直接:任务的生命周期绑定到 try-with-resources 块。
try (var scope = new StructuredTaskScope.ShutdownOnFailure()) {
Future<String> user = scope.fork(() -> fetchUser(id));
Future<Order> order = scope.fork(() -> fetchOrder(id));
scope.join(); // 等待所有子任务完成
scope.throwIfFailed(); // 任何一个失败则抛出
return new Response(user.resultNow(), order.resultNow());
}这段代码的精髓在于:try 块结束时,scope.close() 被自动调用,它会确保所有 fork 出的子任务要么已经结束,要么被强制取消。你不需要手动追踪 Future 列表,不会有"漏取消"的情况。结构化并发把并发任务的生命周期变成了作用域可见的、可组合的。
底层原理:close() 怎么保证不漏
StructuredTaskScope.close() 的内部实现大致是:
- 检查 scope 是否已经 shutdown(通过
shutdown()方法标记)。如果还没有,说明join()之后子任务已经全部正常结束,直接返回。 - 如果 scope 已经 shutdown 但还有未完成的子任务(比如调用了
ShutdownOnFailure策略后还有正在运行的子任务),会调用Future.cancel(true)强制中断。 - 关键限制:
close()方法会阻塞等待所有子任务结束,但最多等一小段时间(默认同join()的超时时间)。如果子任务无视中断信号(比如卡在不可中断的 IO 上),close()会抛出IllegalStateException。
这意味着结构化并发不是"零延迟取消",而是"尽力取消 + 超时报错"。如果子任务里调了 InputStream.read() 这种不可中断的阻塞调用,中断信号发过去也没用,子任务会继续阻塞,直到 close() 超时抛出异常。
时序图:一次完整的结构化并发调用
主线程 scope.fork() scope.fork()
| | |
|--- scope.fork(f1) ---->| |
|<--- Future<A> ---------| |
| | |
|--- scope.fork(f2) --------------------------->|
|<--- Future<B> ---------------------------------|
| | |
|--- scope.join() ------>| |
| |--- f1 执行完毕 --------|
| | |--- f2 执行完毕
|<--- join() 返回 ------| |
| | |
|--- scope.throwIfFailed() ---> 无异常,继续
| | |
|--- scope.close() ----->| |
| | (子任务已结束,直接返回)
| | |
|--- return Response --- 组合结果如果 f1 失败了:
主线程 scope.fork() scope.fork()
| | |
|--- scope.fork(f1) ---->| |
|--- scope.fork(f2) --------------------------->|
| | |
|--- scope.join() ------>| |
| |--- f1 抛出异常 --------|
| | ShutdownOnFailure 触发
| | scope.shutdown() ---> 取消 f2
|<--- join() 返回 ------| |
| |--- f2 被 cancel(true) |
|--- scope.throwIfFailed() ---> 抛出 ExecutionException
| | |
|--- scope.close() ----->| (scope 已 shutdown,直接返回)ShutdownOnFailure 与 ShutdownOnSuccess
StructuredTaskScope 提供了两个内置策略,覆盖了最常用的两种并发场景。
ShutdownOnFailure:全部成功模式
适用场景:所有子任务都必须成功,有一个失败整个操作就无意义。
// 真实场景:查询订单详情,同时查用户信息和商品信息
try (var scope = new StructuredTaskScope.ShutdownOnFailure()) {
Future<UserInfo> user = scope.fork(() -> userService.getUser(order.getUserId()));
Future<ProductInfo> product = scope.fork(() -> productService.getProduct(order.getProductId()));
Future<List<Logistics>> logistics = scope.fork(() -> logisticsService.query(order.getOrderId()));
scope.join(); // 三个同时跑,谁快谁慢不关心,一起等
scope.throwIfFailed(); // 任何一个抛异常,这里就抛出来
return new OrderDetail(
user.resultNow(),
product.resultNow(),
logistics.resultNow()
);
}如果 userService.getUser() 抛了 UserNotFoundException,shutdown 会被触发,product 和 logistics 两个子任务会被自动取消,主线程不会等它们浪费的时间。
ShutdownOnSuccess:谁先到用谁
适用场景:冗余调用,用最快的结果。
// 真实场景:多数据中心冗余查询,减少 tail latency
try (var scope = new StructuredTaskScope.ShutdownOnSuccess<String>()) {
scope.fork(() -> queryFromDC("beijing", query));
scope.fork(() -> queryFromDC("shanghai", query));
scope.fork(() -> queryFromDC("guangzhou", query));
String result = scope.join().result(); // 取最先完成的
return result; // 其他两个自动取消
}这个模式对降低 p99 延迟特别有效。假设每个数据中心有 50ms 的波动,单个数据中心 p99 是 200ms,三个数据中心冗余并发的 p99 可以降到 120ms 左右(因为三个独立波动取最小)。
自定义策略
内置策略不够用的时候,可以继承 StructuredTaskScope:
// 自定义策略:收集所有结果,但如果有超过 2 个子任务失败就提前退出
public class QuorumScope<T> extends StructuredTaskScope<T> {
private final int maxFailures;
private final AtomicInteger failureCount = new AtomicInteger(0);
private final CopyOnWriteArrayList<T> results = new CopyOnWriteArrayList<>();
public QuorumScope(int maxFailures) {
this.maxFailures = maxFailures;
}
@Override
protected void handleComplete(Future<T> future, T result, Throwable error) {
if (error != null) {
if (failureCount.incrementAndGet() > maxFailures) {
shutdown(); // 失败太多,提前结束
}
} else {
results.add(result);
}
}
public List<T> results() { return List.copyOf(results); }
}配合虚拟线程
结构化并发和虚拟线程(Virtual Threads)是 Java 21 Loom 项目的两大支柱,它们配合使用效果最佳。虚拟线程解决了"轻量级线程"的问题,结构化并发解决了"任务生命周期管理"的问题。
┌─────────────────────────────────────────────────────┐
│ Loom 项目(JDK 19-21) │
│ │
│ ┌──────────────────┐ ┌──────────────────────────┐ │
│ │ Virtual Threads │ │ Structured Concurrency │ │
│ │ (JEP 444) │ │ (JEP 453) │ │
│ │ │ │ │ │
│ │ 轻量级线程 │ │ 任务生命周期管理 │ │
│ │ 百万级不崩 │ │ 自动取消 / 错误传播 │ │
│ │ 挂起不阻塞 OS │ │ 调试友好 / 可组合 │ │
│ └────────┬─────────┘ └────────────┬─────────────┘ │
│ │ │ │
│ └─────────── 配合 ─────────┘ │
│ 效果最佳 │
└─────────────────────────────────────────────────────┘// 虚拟线程 + 结构化并发
try (var scope = new StructuredTaskScope.ShutdownOnFailure()) {
// scope.fork 提交的任务默认在虚拟线程中执行
Future<String> r1 = scope.fork(() -> expensiveIoCall1());
Future<String> r2 = scope.fork(() -> expensiveIoCall2());
scope.join();
scope.throwIfFailed();
return combine(r1.resultNow(), r2.resultNow());
}虚拟线程的轻量特性让 fork 大量并发子任务变得低廉,而结构化并发确保这些子任务不会失控。
虚拟线程的线程模型对比
| 维度 | 平台线程 | 虚拟线程 |
|---|---|---|
| 创建成本 | ~1MB 栈空间 + OS 线程 | ~几百字节,随用随建 |
| 最大数量 | 几千(受 OS 限制) | 百万级 |
| 阻塞代价 | 阻塞 OS 线程,线程池占坑 | 挂起到 heap,不占 OS 线程 |
| 适用场景 | CPU 密集型 | IO 密集型 + 大量并发子任务 |
| 与结构化并发 | 配合良好,但 fork 太多会 OOM | 完美配合,随便 fork |
踩坑:虚拟线程 + 结构化并发的常见陷阱
陷阱 1:synchronized 块导致虚拟线程 pinning
虚拟线程在 synchronized 块内执行阻塞操作时,会被 pinned 到平台线程,失去轻量优势。如果 fork 的子任务里用了 synchronized,大量的虚拟线程会被 pinned,导致平台线程耗尽。
解法:用 ReentrantLock 替换 synchronized。
// 错误写法
scope.fork(() -> {
synchronized (lock) { // 虚拟线程被 pinned!
return doIo();
}
});
// 正确写法
scope.fork(() -> {
lock.lock(); // ReentrantLock 不会 pinning
try {
return doIo();
} finally {
lock.unlock();
}
});陷阱 2:fork 太多子任务
虽然虚拟线程轻量,但 StructuredTaskScope 的 join() 要等所有子任务完成。如果 fork 了 10 万个子任务,join() 等待的开销和上下文切换的代价不可忽视。
解法:用 fork 分片,控制并发度。
// 不要 fork 10 万个
for (int i = 0; i < 100_000; i++) {
scope.fork(() -> processItem(i)); // join() 会等到崩溃
}
// 应该分片
int batchSize = 1000;
for (int i = 0; i < 100_000; i += batchSize) {
int start = i;
int end = Math.min(i + batchSize, 100_000);
scope.fork(() -> processBatch(start, end));
}错误传播与取消
结构化并发最关键的改进在错误处理。传统模式下:
ExecutorService es = Executors.newFixedThreadPool(10);
Future<String> f1 = es.submit(() -> { throw new RuntimeException("fail"); });
Future<String> f2 = es.submit(() -> { /* 还在跑,但没人管 */ });
try { f1.get(); } catch (Exception e) { /* 捕捉到了,但 f2 依然在跑 */ }而结构化并发下:
try (var scope = new StructuredTaskScope.ShutdownOnFailure()) {
Future<String> f1 = scope.fork(() -> { throw new RuntimeException("fail"); });
Future<String> f2 = scope.fork(() -> { /* 会被自动取消 */ });
scope.join(); // 因为 f1 失败,ShutdownOnFailure 会关闭 scope
scope.throwIfFailed(); // 抛出 ExecutionException,携带 f1 的原始异常
}这种机制保证了要么全部成功,要么全部取消,避免了"僵尸任务"和"吞没异常"。
生产环境踩坑:超时设置
join() 默认没有超时。如果子任务里有一个卡死(比如网络不通的 HTTP 调用),join() 会一直等下去。必须设超时:
try (var scope = new StructuredTaskScope.ShutdownOnFailure()) {
Future<String> f1 = scope.fork(() -> slowCall());
scope.join(Duration.ofSeconds(5)); // 5 秒超时
scope.throwIfFailed();
return f1.resultNow();
}但注意:join(timeout) 超时返回后,子任务还在跑!close() 时才会尝试取消。所以 join() 超时后应该手动处理:
try (var scope = new StructuredTaskScope.ShutdownOnFailure()) {
Future<String> f1 = scope.fork(() -> slowCall());
try {
scope.join(Duration.ofSeconds(5));
} catch (TimeoutException e) {
// 超时了,手动关闭 scope
scope.shutdown(); // 主动关闭,让 close() 取消子任务
throw new TimeoutException("查询超时");
}
scope.throwIfFailed();
return f1.resultNow();
}调试:jstack 能看到什么
结构化并发最大的隐形收益是可调试性。传统线程池模式下,你看到 jstack 的输出:
"pool-1-thread-3" #13 prio=5 os_prio=0 tid=0x00007f...
at com.example.Worker.run(Worker.java:45)
at java.util.concurrent.ThreadPoolExecutor.runWorker(ThreadPoolExecutor.java:1136)你不知道这个线程在干哪个业务任务,也不知道是谁提交的。结构化并发下:
"ForkJoinPool-1-worker-1" #13
com.example.OrderService.lambda$getOrderDetail$1(OrderService.java:42)
// 能看出是 OrderService.getOrderDetail 里的子任务子任务的名字和调用栈能追溯到父任务,这对线上排查问题极为有用。
不是银弹
结构化并发不是万能的。它不适合:
- 异步事件驱动模型:像 Netty 的 ChannelHandler、Vert.x 的 EventLoop 这种模式,任务的生命周期不由一个作用域控制,而是由事件驱动,结构化并发强行套用会写出反直觉的代码。
- 长时间运行的守护任务:比如后台健康检查、定时轮询。这类任务不需要作用域管理,反而是"发射后不管"的典型场景。
- 需要跨作用域共享结果:比如一个子任务的结果要传给另一个独立的作用域,结构化并发强制父子关系,跨作用域传递需要额外机制。
总结
- 结构化并发把并发任务的生命周期绑定到代码作用域,靠 try-with-resources 确保不会泄漏子任务
- 底层通过
close()配合shutdown()机制实现自动取消,但不可中断的 IO 阻塞会导致取消失败 - ShutdownOnFailure 和 ShutdownOnSuccess 覆盖了"全部成功"和"取最先成功"两种典型场景
- 配合虚拟线程使用,可创建大量子任务而不用担心线程开销,但要注意 synchronized pinning 和 fork 数量控制
- 错误传播机制保证子任务异常不会吞没,取消逻辑自动执行
- 调试时 jstack 能显示父子任务关联,告别"线程池里的线程来自哪里"的困惑
- 必须设
join()超时,否则子任务卡死会导致整个 scope 无法退出 - 不适用于事件驱动模型、守护任务、跨作用域共享场景
参考
JDK 21 JEP 453: Structured Concurrency — https://openjdk.org/jeps/453 JDK 21 JEP 444: Virtual Threads — https://openjdk.org/jeps/444 《Java Concurrency in Practice》第 6 章 — 传统线程模型的问题根源 Martin Sústrik, "Structured Concurrency" — http://250bpm.com/blog:71 Ron Pressler, "Structured Concurrency" (Inside Java Podcast) — https://inside.java/2022/11/30/podcast-33/