Skip to content

结构化并发(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()。下面的代码在线上真实发生过:

java
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 块

java
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() 的内部实现大致是:

  1. 检查 scope 是否已经 shutdown(通过 shutdown() 方法标记)。如果还没有,说明 join() 之后子任务已经全部正常结束,直接返回。
  2. 如果 scope 已经 shutdown 但还有未完成的子任务(比如调用了 ShutdownOnFailure 策略后还有正在运行的子任务),会调用 Future.cancel(true) 强制中断。
  3. 关键限制: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:全部成功模式

适用场景:所有子任务都必须成功,有一个失败整个操作就无意义。

java
// 真实场景:查询订单详情,同时查用户信息和商品信息
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 会被触发,productlogistics 两个子任务会被自动取消,主线程不会等它们浪费的时间。

ShutdownOnSuccess:谁先到用谁

适用场景:冗余调用,用最快的结果。

java
// 真实场景:多数据中心冗余查询,减少 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

java
// 自定义策略:收集所有结果,但如果有超过 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     │  │  调试友好 / 可组合       │  │
│  └────────┬─────────┘  └────────────┬─────────────┘  │
│           │                          │                │
│           └─────────── 配合 ─────────┘                │
│                     效果最佳                          │
└─────────────────────────────────────────────────────┘
java
// 虚拟线程 + 结构化并发
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

java
// 错误写法
scope.fork(() -> {
    synchronized (lock) {  // 虚拟线程被 pinned!
        return doIo();
    }
});

// 正确写法
scope.fork(() -> {
    lock.lock();  // ReentrantLock 不会 pinning
    try {
        return doIo();
    } finally {
        lock.unlock();
    }
});

陷阱 2:fork 太多子任务

虽然虚拟线程轻量,但 StructuredTaskScopejoin() 要等所有子任务完成。如果 fork 了 10 万个子任务,join() 等待的开销和上下文切换的代价不可忽视。

解法:用 fork 分片,控制并发度。

java
// 不要 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));
}

错误传播与取消

结构化并发最关键的改进在错误处理。传统模式下:

java
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 依然在跑 */ }

而结构化并发下:

java
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() 会一直等下去。必须设超时

java
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() 超时后应该手动处理:

java
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 里的子任务

子任务的名字和调用栈能追溯到父任务,这对线上排查问题极为有用。

不是银弹

结构化并发不是万能的。它不适合:

  1. 异步事件驱动模型:像 Netty 的 ChannelHandler、Vert.x 的 EventLoop 这种模式,任务的生命周期不由一个作用域控制,而是由事件驱动,结构化并发强行套用会写出反直觉的代码。
  2. 长时间运行的守护任务:比如后台健康检查、定时轮询。这类任务不需要作用域管理,反而是"发射后不管"的典型场景。
  3. 需要跨作用域共享结果:比如一个子任务的结果要传给另一个独立的作用域,结构化并发强制父子关系,跨作用域传递需要额外机制。

总结

  • 结构化并发把并发任务的生命周期绑定到代码作用域,靠 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/

手撕 → 框架 → 生产化,一步步把 AI Agent 工程化搞透。