Java 基础体系 · 第 68/100 篇。示例统一以 Java 25 LTS 为语言和 JVM 基线;框架示例使用与其兼容的现代稳定版本。

Java 25 结构化并发:任务作用域、失败传播、取消与预览边界

Java 的并发代码长期存在一个容易被忽略的问题:启动任务很容易,证明任务何时结束却很难

例如,父任务启动两个线程完成一次聚合查询:

Future<String> user = executor.submit(this::loadUser);
Future<List<Order>> orders = executor.submit(this::loadOrders);

return combine(user.get(), orders.get());

表面上,get() 等待了两个结果;但异常路径并不完整:

  • loadUser() 失败后,loadOrders() 可能仍在后台运行;
  • 调用方取消聚合任务时,两个子任务是否都能收到取消信号取决于额外代码;
  • executor 的生命周期通常比这次请求更长,任务可能脱离请求继续运行;
  • 子任务可能访问请求已经失效的上下文或资源;
  • 代码结构没有直接表达“这两个任务属于这次调用,并且调用返回前必须处理它们”。

结构化并发试图把并发任务组织成类似方法调用的词法结构:在一个任务作用域中 fork 子任务,等待或决定如何处理它们,然后离开作用域。Java 25 提供的 StructuredTaskScope 正是这一模型的预览 API。

但必须先明确版本边界:

Java 25 是 LTS 版本,但结构化并发在 Java 25 仍是预览特性,而不是标准稳定 API。编译和运行需要 --enable-preview,未来版本仍可能调整 API。

官方 Java SE 25 文档中可以查看 StructuredTaskScope 和相关语言规范。语言规范本身不定义该并发类的具体行为;具体 API 语义以 Java SE API 文档和对应的 JEP 说明为准。


一、先建立模型:任务必须有父子关系和生命周期边界

1. 结构化并发解决什么问题

设父任务为 P,它创建子任务集合:

C(P)={c1,c2,,cn}C(P) = \{c_1, c_2, \ldots, c_n\}

结构化并发要求这些子任务满足几个重要关系:

  1. 词法包含:子任务在父任务的作用域中创建;
  2. 生命周期包含:父任务离开作用域时,子任务不能无边界地继续存在;
  3. 完成可见:父任务在返回结果前,必须经过规定的 join 或关闭过程;
  4. 失败可传播:子任务的失败能沿父子关系向上传递;
  5. 取消可传播:父任务取消时,取消信号能向未完成的子任务传播。

可以把一次结构化调用抽象成:

父任务进入 scope
    ├── fork 子任务 A
    ├── fork 子任务 B
    ├── join,等待或根据策略提前结束
    ├── 取得聚合结果,或抛出失败
父任务离开 scope

对应的非结构化模型通常是:

父任务启动 A、B
父任务可能返回、抛异常或被取消
A、B 仍可能存活于共享 executor 中

后者并不一定错误。后台刷新、消息消费、定时任务本来就可能需要独立生命周期。但对于“请求内部并行执行若干子步骤”这种任务,非结构化生命周期会制造大量额外管理工作。

2. 结构化并发不等于“所有线程都自动结束”

结构化并发约束的是任务之间的生命周期关系,不是强行杀死线程。

如果子任务收到中断后仍然执行不可中断操作,取消可能只能停留在“已发出请求”的状态。Java 没有安全的通用机制强制终止任意线程,因此结构化作用域依赖任务合作:

  • 阻塞方法应正确响应 InterruptedException
  • 循环应检查中断状态;
  • 外部客户端和数据库驱动应提供可取消或超时能力;
  • finally 中应释放资源。

因此,“作用域关闭”不代表可以无条件瞬间终止任意代码,而是代表作用域会执行规定的关闭和取消动作,并等待任务完成关闭过程。


二、Java 25 的 StructuredTaskScope

1. 作用域的基本生命周期

Java 25 的结构化并发 API 以 StructuredTaskScope 为核心。典型生命周期是:

  1. 打开作用域;
  2. 在作用域中调用 fork 创建子任务;
  3. 调用 join 或带截止时间的等待方法;
  4. 从作用域的 Joiner 获取结果;
  5. 通过 try-with-resources 关闭作用域。

一个完整示例:

import java.util.concurrent.StructuredTaskScope;

public class StructuredDemo {
    static String loadUser() throws InterruptedException {
        Thread.sleep(200);
        return "user-42";
    }

    static String loadOrders() throws InterruptedException {
        Thread.sleep(300);
        return "orders-of-user-42";
    }

    static String aggregate() throws Exception {
        try (var scope = StructuredTaskScope.open(
                StructuredTaskScope.Joiner.<String>allSuccessfulOrThrow())) {

            var user = scope.fork(StructuredDemo::loadUser);
            var orders = scope.fork(StructuredDemo::loadOrders);

            scope.join();

            // allSuccessfulOrThrow() 成功时,Joiner 才能生成结果
            return scope.joiner().result();
        }
    }

    public static void main(String[] args) throws Exception {
        System.out.println(aggregate());
    }
}

使用 JDK 25 编译和运行:

javac --enable-preview --release 25 StructuredDemo.java
java --enable-preview StructuredDemo

预期输出:

[ user-42, orders-of-user-42 ]

具体结果表示形式由 Java 25 API 中该 Joiner 的实现定义;这里重要的是生命周期和错误路径,而不是集合的字符串格式。

代码中的每一步有明确含义:

  • StructuredTaskScope.open(...) 创建一个新的任务作用域;
  • fork(...) 创建并启动一个子任务;
  • 子任务通常使用虚拟线程执行,具体线程工厂和配置由作用域配置决定;
  • scope.join() 等待作用域按照 Joiner 的规则达到可返回状态;
  • scope.joiner().result() 将已完成的子任务结果聚合起来;
  • try-with-resources 确保正常返回、异常退出和取消路径都会执行关闭操作。

var uservar ordersStructuredTaskScope.Subtask<String>。它们不是传统意义上已经完成的 Future;fork 返回后,任务可能仍在运行。只有作用域完成相应的等待和策略处理后,才能安全地取得成功结果。

2. 作用域的词法结构

try 代码块本身构成清晰的生命周期边界:

try (var scope = StructuredTaskScope.open(joiner)) {
    scope.fork(taskA);
    scope.fork(taskB);
    scope.join();
    return scope.joiner().result();
}

可以将状态变化简化为:

stateDiagram-v2
    [*] --> OPEN: open
    OPEN --> RUNNING: fork
    RUNNING --> JOINING: join
    JOINING --> SUCCEEDED: Joiner 达成成功条件
    JOINING --> FAILED: Joiner 达成失败条件
    RUNNING --> SHUTTING_DOWN: shutdown
    JOINING --> SHUTTING_DOWN: 失败/成功策略触发取消
    SHUTTING_DOWN --> CLOSED: 子任务结束后 close
    SUCCEEDED --> CLOSED: try-with-resources
    FAILED --> CLOSED: try-with-resources
    CLOSED --> [*]

实际 API 还包含更细的非法状态检查。例如作用域关闭后不能继续 fork,父线程未按要求完成 join 时也不能把作用域当作普通集合直接读取。这里的核心不是记忆每个异常类型,而是理解:

StructuredTaskScope 不是一个可无限期复用的线程池;它是一次任务树分支的生命周期对象。


三、Joiner:失败传播和成功条件由策略决定

结构化并发并没有规定所有任务必须采用同一种完成策略。两个常见需求完全不同:

1. “全部成功,否则整体失败”

聚合用户资料时,用户信息和权限信息都必须存在:

try (var scope = StructuredTaskScope.open(
        StructuredTaskScope.Joiner.<String>allSuccessfulOrThrow())) {

    scope.fork(this::loadUser);
    scope.fork(this::loadPermissions);

    scope.join();
    return scope.joiner().result();
}

allSuccessfulOrThrow() 的逻辑可以形式化为:

成功    ciC(P),  ci 成功\text{成功} \iff \forall c_i \in C(P),\; c_i \text{ 成功}

只要存在一个失败子任务:

ciC(P),  ci failed\exists c_i \in C(P),\; c_i \text{ failed}

整体就不应产生成功结果。该策略还会在失败条件满足后关闭或关闭剩余任务,避免无意义地继续执行。

注意两个不同动作:

  1. 子任务失败;
  2. 父任务把该失败作为自己的失败抛出。

子任务失败首先是一个子任务状态;只有 Joiner 的结果处理或显式检查把它转化为父任务异常,调用方才会看到失败。不能把“某个子任务抛异常”理解成“父线程已经自动抛异常”。

2. “任意一个成功即可”

访问多个同类服务时,可以使用竞速策略:

try (var scope = StructuredTaskScope.open(
        StructuredTaskScope.Joiner.<String>anySuccessfulResultOrThrow())) {

    scope.fork(() -> queryReplica("replica-a"));
    scope.fork(() -> queryReplica("replica-b"));
    scope.fork(() -> queryReplica("replica-c"));

    scope.join();
    return scope.joiner().result();
}

该策略的成功条件是:

成功    ciC(P),  ci 成功\text{成功} \iff \exists c_i \in C(P),\; c_i \text{ 成功}

第一个成功结果达到策略条件后,其余仍未完成的任务应被取消或关闭,避免三个副本在已经得到答案后继续消耗资源。

如果所有副本都失败,则整体失败:

失败    ciC(P),  ci failed\text{失败} \iff \forall c_i \in C(P),\; c_i \text{ failed}

这与“忽略异常并返回第一个完成结果”不同。失败任务不能被错误地当成成功结果,也不能因为某个任务先完成但结果无效,就提前结束作用域。

3. Joiner 不是简单的结果收集器

Joiner 同时描述两件事:

  • 什么时候可以结束等待;
  • 如何从子任务集合产生最终结果。

因此,下面两个策略的并发行为不同:

allSuccessfulOrThrow:
    等待全部成功,任一失败触发失败路径

anySuccessfulResultOrThrow:
    任一成功即可,成功后取消剩余任务;全失败才失败

如果业务策略不是这两种,例如:

  • 至少两个副本成功;
  • 收集所有成功结果,忽略部分失败;
  • 某类异常可降级,另一类异常必须失败;
  • 只在达到截止时间后选择当前最优结果;

则应根据 Java 25 API 提供的 Joiner 扩展机制定义明确策略,不能靠调用方在多个线程间共享可变变量来模拟。自定义策略尤其要明确:

  • 哪些子任务状态算成功;
  • 哪些异常可以忽略;
  • 何时触发 shutdown
  • 结果是否允许部分成功;
  • 截止时间到达时如何构造最终结果。

四、失败传播:从子任务到父任务的完整路径

考虑如下代码:

try (var scope = StructuredTaskScope.open(
        StructuredTaskScope.Joiner.<String>allSuccessfulOrThrow())) {

    scope.fork(() -> {
        throw new IllegalStateException("primary failed");
    });

    scope.fork(() -> {
        Thread.sleep(10_000);
        return "slow result";
    });

    scope.join();
    return scope.joiner().result();
}

失败路径不是“一个线程抛异常,其他线程立即消失”,而是大致经过以下步骤:

  1. 第一个子任务进入失败状态;
  2. Joiner 观察到失败;
  3. allSuccessfulOrThrow() 不能再生成成功结果;
  4. 作用域对尚未完成的子任务执行关闭/取消动作;
  5. 慢任务收到中断;
  6. scope.join() 完成失败路径;
  7. scope.joiner().result() 向父任务报告聚合失败;
  8. try-with-resources 继续执行作用域关闭;
  9. 父任务最终向自己的调用方抛出异常。

如果慢任务写成:

scope.fork(() -> {
    while (true) {
        // 忽略中断
    }
});

那么它不会因为 Java 想要取消它就自动停止。该代码可能导致作用域关闭迟迟不能完成,并且持续占用 CPU。正确的合作式写法至少应保留中断响应:

scope.fork(() -> {
    while (!Thread.currentThread().isInterrupted()) {
        doOnePieceOfWork();
    }
    return "stopped";
});

对于可中断阻塞:

scope.fork(() -> {
    try {
        return blockingCall();
    } catch (InterruptedException e) {
        Thread.currentThread().interrupt();
        throw e;
    }
});

重新设置中断标志很重要,因为捕获 InterruptedException 后,中断状态通常已经被清除。若决定继续向上抛出,恢复中断状态可以让更上层逻辑仍能观察到取消。


五、取消:shutdown 是请求,不是强制终止

1. 取消的传播方向

结构化任务通常形成树:

请求任务 P
├── 子任务 A
│   └── 子作用域 A1
└── 子任务 B
    └── 子作用域 B1

P 的作用域关闭或执行关闭操作时,取消主要向下传播:

P 的 scope shutdown
    ├── interrupt A
    │   └── A 的子作用域进入关闭路径
    └── interrupt B
        └── B 的子作用域进入关闭路径

当某个子任务失败时,是否触发兄弟任务取消,取决于 Joiner 或显式的 shutdown 调用。结构化关系提供传播边界,但不会替业务自动决定“一个失败是否足以取消所有兄弟任务”。

2. 父线程被取消时

如果父线程在 join()joinUntil(...) 中被中断,父任务应进入取消或异常处理路径,并确保作用域关闭:

try (var scope = StructuredTaskScope.open(
        StructuredTaskScope.Joiner.<String>allSuccessfulOrThrow())) {
    scope.fork(this::readFromA);
    scope.fork(this::readFromB);

    scope.join();
    return scope.joiner().result();
} catch (InterruptedException e) {
    Thread.currentThread().interrupt();
    throw new CancellationException("parent task cancelled");
}

这里的 CancellationException 是否适合作为对外异常,取决于应用协议;关键要求是不能吞掉中断。

3. close() 的真实边界

使用 try-with-resources 不是格式偏好,而是生命周期保证的一部分:

try (var scope = StructuredTaskScope.open(joiner)) {
    // fork、join、result
}

当代码因异常提前离开 try 块时,close() 仍会执行。关闭过程会处理尚未结束的子任务,而不是让它们自动成为脱离作用域的后台任务。

但有三个现实边界:

  1. 子任务必须合作响应中断;
  2. 第三方库可能在一段时间内不响应中断;
  3. 作用域关闭需要等待任务完成清理,因此不一定是常数时间操作。

所以结构化并发改善的是取消的可达性和责任边界,不是制造不可失败的取消机制。


六、超时:截止时间必须进入任务树

如果调用方只允许等待两秒,应把截止时间用于作用域的 join,而不是只在返回后检查耗时:

import java.time.Duration;
import java.time.Instant;
import java.util.concurrent.StructuredTaskScope;
import java.util.concurrent.TimeoutException;

static String queryWithTimeout() throws Exception {
    try (var scope = StructuredTaskScope.open(
            StructuredTaskScope.Joiner.<String>allSuccessfulOrThrow())) {

        scope.fork(() -> {
            Thread.sleep(5_000);
            return "slow";
        });

        try {
            scope.joinUntil(Instant.now().plus(Duration.ofSeconds(2)));
            return scope.joiner().result();
        } catch (TimeoutException e) {
            scope.shutdown();
            throw new TimeoutException("structured operation timed out");
        }
    }
}

这里的因果关系是:

  1. joinUntil 等待到截止时间;
  2. 时间耗尽时,父任务得到超时信号;
  3. 显式 shutdown 请求取消未完成子任务;
  4. try-with-resources 再次保证关闭路径;
  5. 子任务只有在实际响应中断后才会结束。

不能只写:

long start = System.nanoTime();
scope.join();
if (elapsed(start) > timeout) {
    throw new TimeoutException();
}

因为这种代码仍然可能无限等待,直到所有任务完成;“事后发现超时”不等于“在截止时间取消等待”。

同时,作用域超时和底层操作超时是两层机制:

作用域超时:父任务不再等待,并请求取消子任务
HTTP/数据库超时:底层调用本身停止等待外部资源

生产代码通常需要两者都配置。若 HTTP 客户端完全不响应线程中断,作用域只能发出取消请求,无法保证底层连接立即释放。


七、与 ExecutorService 的关键区别

下面是传统写法的典型问题:

ExecutorService executor = Executors.newVirtualThreadPerTaskExecutor();

Future<String> a = executor.submit(this::taskA);
Future<String> b = executor.submit(this::taskB);

try {
    return a.get() + b.get();
} finally {
    // 这里若遗漏,任务和 executor 生命周期就不清晰
}

即使改为依次调用 Future.cancel(true),仍然需要自己处理:

  • 哪个任务失败后取消哪些任务;
  • 取消后如何等待任务真正结束;
  • Future.get() 抛出的异常如何聚合;
  • 超时发生在 a.get() 还是 b.get() 时如何取消另一个;
  • 当前请求返回后 executor 是否仍然存活;
  • 多层业务方法如何传播取消。

结构化作用域把其中一部分关系固化到 API 中:

问题 Future 组合 StructuredTaskScope
子任务归属 通常靠变量和约定 由作用域表达
失败策略 调用方手工编排 Joiner 描述
兄弟任务取消 通常手工 cancel 作用域策略可触发 shutdown
作用域退出 可能遗留后台任务 try-with-resources 负责关闭
任务树 不自然 可嵌套建立
API 稳定性 稳定 Java 25 仍为预览

这并不意味着 ExecutorService 被取代。长期运行的线程池、独立后台消费者、定时调度器和跨请求共享的任务,通常不适合放进一次性的结构化作用域。两者解决的是不同生命周期的问题。


八、嵌套作用域和任务树

子任务内部可以继续打开子作用域:

static String parentTask() throws Exception {
    try (var outer = StructuredTaskScope.open(
            StructuredTaskScope.Joiner.<String>allSuccessfulOrThrow())) {

        outer.fork(() -> childAggregate());
        outer.join();

        return outer.joiner().result();
    }
}

static String childAggregate() throws Exception {
    try (var inner = StructuredTaskScope.open(
            StructuredTaskScope.Joiner.<String>allSuccessfulOrThrow())) {

        inner.fork(() -> "left");
        inner.fork(() -> "right");
        inner.join();

        return inner.joiner().result().toString();
    }
}

其生命周期关系是:

inner scopechild taskouter scopeparent task\text{inner scope} \subset \text{child task} \subset \text{outer scope} \subset \text{parent task}

如果外层作用域因超时关闭:

  1. 外层向子任务发送取消;
  2. 子任务中的内层作用域进入关闭路径;
  3. 内层未完成任务收到取消;
  4. 内层清理完成后,子任务结束;
  5. 外层关闭完成。

这种嵌套关系让取消路径与调用栈更接近。相比之下,若每层都向全局 executor 提交任务,父层很难知道孙任务是否仍在运行。


九、常见误解和反例

误解一:fork 后立刻读取子任务结果

错误思路:

var subtask = scope.fork(this::load);
return subtask.get(); // 没有先完成规定的 join

fork 返回的是任务句柄,不代表任务已经完成。结构化并发要求通过作用域的 join 生命周期和 Joiner 规则观察结果,而不是把 Subtask 当作普通同步值。

正确思路是:

scope.fork(this::load);
scope.join();
return scope.joiner().result();

误解二:结构化并发自动让失败立即中止所有任务

失败传播分为三层:

子任务抛异常
    ↓
Joiner 判断是否达到失败条件
    ↓
作用域 shutdown 未完成任务
    ↓
父任务取得失败结果

如果自定义 Joiner 允许部分失败,某个子任务失败可能只是被记录,而不是导致整体失败。失败策略必须显式设计。

误解三:虚拟线程意味着取消一定很快

虚拟线程降低了等待型任务的线程资源成本,但不改变中断语义。以下代码仍可能无法及时取消:

scope.fork(() -> {
    synchronized (lock) {
        while (true) {
            // 长时间占用锁,并且不检查中断
        }
    }
});

虚拟线程不是强制终止机制,也不能修复不合作的任务代码。

误解四:Java 25 LTS 意味着结构化并发已经稳定

LTS 描述的是 Java 25 发行版的长期支持属性,不代表其中每个新 API 都已永久定型。结构化并发在 Java 25 的使用边界仍应写成:

javac --enable-preview --release 25 ...
java --enable-preview ...

构建工具、测试工具和生产启动参数都必须一致配置。遗漏运行时参数会导致预览类或预览 API 无法正常运行;未来 JDK 版本也可能改变方法签名、策略类型或生命周期细节。


十、诊断失败、取消和泄漏

1. 先区分三种时间

诊断结构化并发问题时,应分别测量:

任务执行时间:子任务实际运行了多久
父任务等待时间:join 或 joinUntil 阻塞了多久
作用域关闭时间:shutdown 到所有子任务结束用了多久

如果父任务两秒超时,但作用域五秒后才关闭,通常说明:

  • 子任务没有正确响应中断;
  • 底层 I/O 没有可取消能力;
  • 清理逻辑阻塞;
  • 子任务在中断后仍继续工作。

不能只记录父接口返回时间,否则会漏掉关闭路径中的资源占用。

2. 检查任务是否吞掉中断

问题代码:

try {
    blockingCall();
} catch (InterruptedException ignored) {
    // 继续执行
}

这会破坏取消传播。应根据业务决定抛出、返回取消结果,或恢复中断:

try {
    blockingCall();
} catch (InterruptedException e) {
    Thread.currentThread().interrupt();
    throw e;
}

3. 用线程转储观察作用域关闭

当作用域无法结束时,可以使用 JDK 的线程转储工具观察任务线程的状态,例如:

jcmd <pid> Thread.dump_to_file -format=json threads.json

应关注:

  • 相关虚拟线程是否仍处于 RUNNABLE;
  • 是否卡在锁、网络读取或本地方法;
  • 是否已经收到中断但继续运行;
  • 父线程是否卡在 join 或作用域关闭过程。

命令和输出格式以实际 JDK 25 发行版为准。线程转储只能提供观测证据,不能替代对任务取消协议的修复。


十一、什么时候不应使用结构化作用域

结构化作用域适合“父任务等待一组有明确边界的子任务”:

一次请求中的并行查询
一个批处理项中的并行步骤
一次计算中的副本竞速
一个业务方法内部的扇出/汇聚

以下任务通常不应被强行放入该模型:

消息队列长期消费者
应用级定时调度器
跨请求共享的缓存刷新任务
独立于调用方生命周期的后台任务
需要无限期运行的监控循环

原因不是这些任务不能使用线程,而是它们的生命周期本来就不应绑定到某个父调用。若父请求结束后任务仍必须继续运行,那么把它作为请求作用域的子任务反而会产生错误的取消关系。


十二、Java 25 预览边界下的采用方式

采用结构化并发前,需要同时接受三个事实:

  1. 语义价值是稳定方向:任务树、词法作用域、失败传播和取消传播能减少并发管理中的隐含状态;
  2. API 仍可能变化:Java 25 的具体类、方法和 Joiner 设计属于预览范围;
  3. 工程风险不只在代码:编译参数、运行参数、框架集成、监控和异常协议都必须支持预览特性。

比较稳妥的验证顺序是:

先用 JDK 25 编译最小示例
    ↓
验证 all-success 和 any-success 两条失败路径
    ↓
验证父线程取消和 joinUntil 超时
    ↓
验证子任务响应中断
    ↓
验证作用域关闭时间和资源释放
    ↓
再接入 Web、RPC、数据库或消息框架

尤其要测试反例,而不只是测试成功路径:

  • 一个子任务立即失败,另一个任务长时间阻塞;
  • 父线程在 join() 中被中断;
  • 截止时间到达但底层调用不响应中断;
  • 子任务开启嵌套作用域;
  • Joiner 允许部分成功;
  • close() 执行时仍有未完成任务。

结构化并发的核心收益不是“用更少的代码创建线程”,而是让并发任务的归属、完成、失败和取消形成可检查的生命周期。Java 25 的 StructuredTaskScope 已经能表达这套模型,但由于它仍处于预览阶段,应用必须同时尊重它的语义边界和 API 不稳定边界:以作用域组织任务,以 Joiner 明确成功与失败,以中断实现合作式取消,并把预览参数和升级验证纳入完整交付流程。


系列导航与关联阅读

官方资料

本文依据 Java、Spring 与相关项目官方文档重新梳理;正文、示例与生产清单由 WR BLOG 编写。