结构化并发与 Scoped Values

深入结构化并发与 Scoped Values:StructuredTaskScope 的 Joiner 模型与 JDK 版本演进、ShutdownOnFailure/ShutdownOnSuccess 语义、自定义 Joiner、ScopedValue 替代 ThreadLocal 的作用域传递、取消传播与虚拟线程固定问题,以及生产迁移与可观测性实践。

ExecutorService 提交任务后返回一个游离的 Future,任务的生命周期与提交它的代码块完全脱钩:方法返回了,子任务可能还在跑;主线程抛异常了,子任务照旧执行。这种「失控的并发」在虚拟线程时代会被放大——你可以轻松创建上百万个线程,却很难保证它们被正确回收。结构化并发(Structured Concurrency)用词法作用域把子任务的生命周期绑定到代码块,Scoped Values 则用不可变作用域替代可变的 ThreadLocal。本文把这两个 JEP 讲透。

一、非结构化并发的三个顽疾

// 典型写法:Future 游离,生命周期不可控
public Response handle(Request req) {
    Future<User> user = pool.submit(() -> loadUser(req));
    Future<Orders> orders = pool.submit(() -> loadOrders(req));
    try {
        return new Response(user.get(), orders.get());
    } catch (Exception e) {
        // 一个失败,另一个还在跑;取消逻辑要手写
        throw new RuntimeException(e);
    }
}
顽疾表现后果
取消不传播一个子任务失败,其他照跑资源浪费、延迟尾部长
异常丢失子任务异常无人 get静默失败,难排查
生命周期失控方法返回后任务仍在跑线程泄漏、观测困难
结构化并发的核心承诺:
  1. 子任务的生命周期 ≤ 父作用域
  2. 作用域退出前,所有子任务必须结束(正常/取消/失败)
  3. 取消沿作用域树自动向下传播
  4. 失败沿作用域树向上冒泡,不丢失

一句话总结: 结构化并发把「并发任务的树形结构」显式化——作用域不退出,子任务不逃逸,这解决了 Future 模型的取消与异常两大痛点。

二、StructuredTaskScope 的版本演进

JDK 版本状态关键变化
JDK 19孵化初版 StructuredTaskScope,fork/join
JDK 20孵化ShutdownOnSuccess 引入
JDK 21预览API 大改:Joiner 抽象、工厂方法
JDK 22/23预览join() 语义与异常类型调整
JDK 24预览ScopedValue 配套稳定
JDK 25正式StructuredTaskScope 与 ScopedValue 转正
# 预览期需要显式开启
javac --release 21 --enable-preview App.java
java  --enable-preview App
JDK 21 后的 API 形状:
  try (var scope = StructuredTaskScope.open(joiner)) {
      var a = scope.fork(taskA);
      var b = scope.fork(taskB);
      scope.join();          // 等待并聚合
      return combine(a.get(), b.get());
  }
  // 退出 try:作用域关闭,所有子任务已结束

一句话总结: 结构化并发经历了 6 个版本的孵化与预览才转正,JDK 21 的 Joiner 抽象是分水岭——早期示例代码在 JDK 25 上基本不能直接编译,升级时务必对照官方迁移说明。

三、Joiner:失败与成功的语义

3.1 ShutdownOnFailure

任一子任务失败即取消其余,适合「所有子任务都必须成功」的场景:

Response handle(Request req) throws Exception {
    try (var scope = StructuredTaskScope.open(
            Joiner.<Response>awaitAllSuccessfulOrThrow())) {
        var user   = scope.fork(() -> loadUser(req));
        var orders = scope.fork(() -> loadOrders(req));
        scope.join();                    // 任一失败 → 抛异常并取消其余
        return new Response(user.get(), orders.get());
    }
}
语义要点:
  awaitAllSuccessfulOrThrow()
    - 全部成功 → join() 正常返回
    - 任一失败 → 取消其他子任务,join() 抛出失败原因
    - 抛出的异常类型是 StructuredTaskScope.FailedException,
      其 cause 是真正的业务异常

3.2 ShutdownOnSuccess

任一子任务成功即取消其余,适合「多源竞速、取最快结果」:

String fetchFastest(String key) throws Exception {
    try (var scope = StructuredTaskScope.open(
            Joiner.<String>anySuccessfulResultOrThrow())) {
        scope.fork(() -> queryCache(key));
        scope.fork(() -> queryReplicaA(key));
        scope.fork(() -> queryReplicaB(key));
        return scope.join();             // 返回第一个成功结果
    }
}
竞速模式的注意点:
  1. 被取消的子任务要能响应中断(阻塞调用需可中断)
  2. 结果不确定 —— 三个副本数据不一致时以最快为准
  3. 取消不是「立刻停」,是设置中断标志,配合可中断 API 生效

一句话总结: Joiner 把「何时关闭作用域、返回什么」策略化:awaitAllSuccessfulOrThrow 用于必全成功,anySuccessfulResultOrThrow 用于竞速取快;两者的取消都是协作式中断。

四、自定义 Joiner

当内置语义不够用时(比如「允许部分失败、返回部分结果」),实现自己的 Joiner:

// 收集成功结果,容忍部分失败
class CollectJoiner<T> implements StructuredTaskScope.Joiner<T, List<T>> {
    private final List<T> results = new ArrayList<>();

    @Override
    public boolean onComplete(Subtask<T> subtask) {
        if (subtask.state() == Subtask.State.SUCCESS) {
            results.add(subtask.get());
        }
        return false;              // false = 继续等待其余子任务
    }

    @Override
    public List<T> result() {
        return List.copyOf(results);
    }
}
List<Price> bestEffort(List<Supplier> suppliers) throws Exception {
    try (var scope = StructuredTaskScope.open(new CollectJoiner<Price>())) {
        suppliers.forEach(s -> scope.fork(() -> s.quote()));
        return scope.join();       // 返回成功结果集合,忽略失败
    }
}
onComplete 返回值的含义:
  true  → 立即关闭作用域,取消剩余子任务
  false → 继续等待其他子任务完成
  该方法可能被多个线程并发调用,实现必须线程安全

一句话总结: 自定义 Joiner 是结构化并发的扩展点,onComplete 的布尔返回值就是「是否提前收工」的开关,实现时注意并发安全。

五、Scoped Values:ThreadLocal 的替代品

5.1 ThreadLocal 的问题

问题说明
可变任何代码都能 set,数据流难追踪
生命周期失控忘记 remove 造成线程池串数据/内存泄漏
继承昂贵InheritableThreadLocal 复制开销大,虚拟线程下几乎不可用
与虚拟线程冲突百万线程各持一份副本,内存爆炸

5.2 ScopedValue 的模型

ScopedValue 是不可变、有作用域、可继承到子任务的隐式参数:

static final ScopedValue<User> CURRENT_USER = ScopedValue.newInstance();

void handle(Request req) {
    User user = authenticate(req);
    ScopedValue.where(CURRENT_USER, user).run(() -> {
        // 作用域内任意深度都能读到,无需逐层传参
        doBusiness();
    });
    // 出了作用域,CURRENT_USER.isBound() == false
}

void doBusiness() {
    User u = CURRENT_USER.get();   // 隐式读取
    audit(u);
}
ScopedValue vs ThreadLocal:
  不可变     —— 作用域内只读,无 set,数据流单向
  自动清理   —— 作用域退出即解绑,无泄漏风险
  零拷贝继承 —— 子线程/子任务直接读同一份,不复制
  虚拟线程友好 —— 百万虚拟线程共享一个绑定,内存恒定

5.3 与结构化并发配合

Response handle(Request req) {
    User user = authenticate(req);
    return ScopedValue.where(CURRENT_USER, user).call(() -> {
        try (var scope = StructuredTaskScope.open(Joiner.<Data>awaitAllSuccessfulOrThrow())) {
            // fork 出的子任务自动继承 CURRENT_USER 绑定
            var a = scope.fork(() -> loadWithAudit(req));
            scope.join();
            return new Response(a.get());
        }
    });
}

一句话总结: ScopedValue 用「作用域 + 不可变」解决了 ThreadLocal 的可变与泄漏问题,且天然向下继承给虚拟线程子任务,是虚拟线程时代上下文传递的标准答案。

六、虚拟线程固定(Pinning)与取消传播

结构化并发跑在虚拟线程上时,最需要警惕的是固定(Pinning):虚拟线程被阻塞在 synchronized 块或本地方法中时,无法从载体线程卸载,会拖垮调度器。

Pinning 的两大来源:
  1. synchronized 块内的阻塞(ReentrantLock 不会 pin)
  2. 本地方法(JNI)阻塞
// 反例:synchronized 内做阻塞 IO,会 pin 住载体线程
synchronized (lock) {
    var data = httpClient.send(request, BodyHandlers.ofString());  // 危险
}

// 正例:用 ReentrantLock 替代 synchronized
lock.lock();
try {
    var data = httpClient.send(request, BodyHandlers.ofString());
} finally {
    lock.unlock();
}
# JDK 21+ 诊断 pinning
-Djdk.tracePinnedThreads=full        # JDK 21 打印固定堆栈
-Djdk.tracePinnedThreads=short

# JDK 24+ 用 JFR 事件观察
jfr print --events jdk.VirtualThreadPinned recording.jfr
取消传播的实现:
  1. scope.close() 对每个未完成子任务调用 cancel()
  2. cancel() 触发 Thread.interrupt()
  3. 子任务需在可中断点(阻塞 IO、sleep、await)响应
  4. 不响应中断的忙循环会拖慢作用域关闭

一句话总结: 结构化并发的取消依赖中断协作,写子任务时避免在 synchronized 内做阻塞调用(改用 ReentrantLock),并用 jdk.tracePinnedThreads 或 JFR 事件排查固定。

七、迁移与可观测性实践

7.1 从 ExecutorService 迁移

迁移检查清单:
  1. 把「提交 N 个任务再 join」的代码识别出来 → 改造成作用域
  2. 手写的取消/超时逻辑 → 交给 Joiner 与作用域生命周期
  3. 线程池的并发上限 → 改用 Semaphore 或结构化限流
  4. 逐层传参的上下文 → 改用 ScopedValue
  5. 保留线程池做「资源池化」(如数据库连接)而非「并发控制」
// 用 Semaphore 限制并发度(结构化作用域没有池的概念)
private static final Semaphore PERMITS = new Semaphore(64);

Response handle(Request req) throws Exception {
    PERMITS.acquire();
    try (var scope = StructuredTaskScope.open(Joiner.<Data>awaitAllSuccessfulOrThrow())) {
        var data = scope.fork(() -> load(req));
        scope.join();
        return new Response(data.get());
    } finally {
        PERMITS.release();
    }
}

7.2 可观测性

观测要点:
  1. 线程转储:虚拟线程以「结构化树」呈现,父任务可见子任务
  2. 每个作用域可命名:StructuredTaskScope.open(joiner, cfg -> cfg.withName("load-order"))
  3. JFR 事件:jdk.VirtualThreadStart/End、jdk.VirtualThreadPinned
  4. 与 Micrometer 结合,记录作用域内子任务的成功/失败/取消计数
  5. 超时统一在作用域外层用 orTimeout 或 Joiner 的超时配置
// 给作用域命名,便于线程转储与追踪定位
try (var scope = StructuredTaskScope.open(
        Joiner.<Data>awaitAllSuccessfulOrThrow(),
        cfg -> cfg.withName("order-aggregation"))) {
    // ...
}

一句话总结: 迁移的关键是把「并发控制」与「资源池化」拆开——并发控制交给结构化作用域与 Semaphore,资源池保留线程池;观测上善用作用域命名与 JFR 事件。

八、常见陷阱

陷阱现象对策
作用域内 fork 后不 joinjoin() 前读取结果报错必须先 join 再 get
在作用域外使用 Subtask结果不可靠Subtask 只在作用域内有效
子任务阻塞不可中断关闭作用域卡住使用可中断 API
synchronized 内阻塞虚拟线程固定改 ReentrantLock
ScopedValue 当作可变缓存编译不过 / 逻辑错它只读,缓存用别的机制
忘记异常聚合只看第一个异常FailedException 里遍历 suppressed
作用域嵌套过深线程转储难读命名 + 分层,别超三层
// 正确读取失败原因(含被抑制的异常)
try {
    scope.join();
} catch (StructuredTaskScope.FailedException e) {
    Throwable cause = e.getCause();
    for (Throwable s : cause.getSuppressed()) {
        log.warn("子任务失败", s);
    }
}

一句话总结: 结构化并发的坑大多来自「作用域生命周期」与「虚拟线程固定」;先 join 后 get、用可中断 API、避免 synchronized 阻塞,能规避九成问题。

小结

维度结构化并发Scoped Values
解决的问题任务生命周期与取消传播上下文隐式传递
核心 APIStructuredTaskScope + JoinerScopedValue.where(...).run/call
关键语义作用域不退出,子任务不逃逸不可变、自动清理、零拷贝继承
虚拟线程天然适配,注意 pinning天然适配,内存恒定
落地版本JDK 25 转正JDK 25 转正

结构化并发不是「更花哨的线程池」,而是把并发任务的组织方式从「散点」变成「树」,让取消、异常与生命周期都有明确归属。配合 ScopedValue 的不可变上下文,虚拟线程时代的并发代码第一次拥有了和同步代码一样清晰的控制流。

延伸阅读

继续阅读

探索更多技术文章

浏览归档,发现更多关于系统设计、工具链和工程实践的内容。

全部文章 返回首页

「java」更多文章

  1. Testcontainers 集成测试
  2. Micrometer 可观测性
  3. Spring Batch 批处理