CodeGym /课程 /JAVA 25 SELF /Structured Concurrency

Structured Concurrency

JAVA 25 SELF
第 58 级 , 课程 0
可用

1. 引言

当线程各自为政

传统的多线程常常像一场没有指挥的排练。每个线程——就像各自演奏的乐手,互不倾听。有人早早结束去抽烟,有人卡在一个和弦上,有人干脆弹错了音而抛出错误。结果不是交响乐,而是噪音:想弄清谁在哪儿出错几乎不可能,而想让大家一起停下来——更是难上加难。

Structured Concurrency 解决了这个问题。它把零散的线程变成真正的合奏:所有任务都由同一位“指挥”统一管理。如果他下达停止的指令——乐队就会静音。如果某位乐手失误——其余人会平稳地停止,不破坏整体和谐。所有结果与错误集中收集,而不是散落在代码的各个角落。

想象一下: 你不再让乐手各弹各的,而是把他们聚到同一间大厅。有指挥、有谱子,即便小号走音——乐队也不会崩盘,而是体面地结束演出。

Structured Concurrency 带来什么

  • 统一的任务“作用域”:所有子任务都在同一个代码块内生存,它们的生命周期受该块约束。
  • 可预期的结束方式:父线程不会在所有子任务结束之前结束。
  • 集中化取消:如果某个任务失败或父任务决定结束——所有子任务都会被正确地取消。
  • 一致的错误处理:子任务的错误会被聚合,你可以得到“原因树”(tree of causes)。
  • 整洁且可读的代码:没有“悬挂”的线程,没有被遗忘的任务,也没有争抢取消的竞态。

Structured Concurrency 不只是一个新的 API,而是一种新的思维方式:任务应当像普通代码块一样被结构化(例如 try-with-resources)。

2. Java 中的 Structured Concurrency 状态

在撰写本课程时,Structured Concurrency 处于 Preview 状态(Java 21–23),但预计会在 Java 24/25 推进到 GA(General Availability)。API 位于 jdk.incubator.concurrent 包。在用于生产之前,请务必查看你所用 JDK 版本的最新发行说明!

核心类:

  • StructuredTaskScope —— 用于管理任务组的基类。
  • 变体:StructuredTaskScope.ShutdownOnFailureStructuredTaskScope.ShutdownOnSuccess —— 任务的终止策略。

StructuredTaskScope 的核心概念

模型:fork、join 与结果的“点名”

当指挥(也就是父任务)发出信号——子任务就会奔赴各自的声部。这个时刻称为 fork——好比让乐手们在不同的房间演奏各自的片段。

随后进入 join 时刻——指挥举起指挥棒,所有人回到一起,奏出最终的和弦。

接着你可以向每个参与者询问进展:

  • 通过 resultNow() 立即获得结果(如果一切顺利无误);
  • 通过 throwIfFailed()——确认没人“走音”。如果有人确实乱了套,将抛出统一的异常——就像指挥说:“我们出故障了,重来”。

终止策略

每位指挥都有自己的规则来决定何时让音乐停止。在 Structured Concurrency 中,这由终止策略来设定:

  • ShutdownOnFailure —— 只要有乐手失去节奏,指挥挥手:“停!重来。”其他人立即停止演奏。
  • ShutdownOnSuccess —— 相反,一旦有人完美演奏,指挥就满意地说:“够了,没必要再继续,我们已经有了赢家。”其余人静音——这是“第一个成功者”的策略。

与虚拟线程配合

StructuredTaskScope 将每个子任务运行在虚拟线程中。这就像拥有一支通晓配合、反应迅速且不挑舞台的乐队。可以放心地创建上百、上千个这样的执行者——它们并非重量级线程,而是几乎无重量的音符,只在需要时发声。

3. 示例:HTTP 请求聚合器

来看一个实际问题:我们有三个数据源(例如三个不同的服务器),我们希望要么同时获得所有的成功结果(并进行聚合),要么取第一个成功返回的结果。

方案 1:“全部必须成功”(ShutdownOnFailure

import jdk.incubator.concurrent.StructuredTaskScope;
import java.util.concurrent.Future;

public class AggregatorAllSuccess {
    public static void main(String[] args) throws Exception {
        try (var scope = new StructuredTaskScope.ShutdownOnFailure()) {
            Future<String> f1 = scope.fork(() -> fetchFromSource1());
            Future<String> f2 = scope.fork(() -> fetchFromSource2());
            Future<String> f3 = scope.fork(() -> fetchFromSource3());

            scope.join(); // 等待所有任务完成
            scope.throwIfFailed(); // 只要有一个失败——抛出异常

            // 全部成功——可以聚合结果
            String result = f1.resultNow() + f2.resultNow() + f3.resultNow();
            System.out.println("聚合结果:" + result);
        }
    }

    static String fetchFromSource1() { /* ... */ return "A"; }
    static String fetchFromSource2() { /* ... */ return "B"; }
    static String fetchFromSource3() { /* ... */ return "C"; }
}

发生了什么:

  • 三项任务并行启动(在虚拟线程中)。
  • 只要其中一个失败——其余任务会被取消,并抛出异常。
  • 如果全部成功——就可以安全地聚合结果。

方案 2:“第一个有效即为成功”(ShutdownOnSuccess

import jdk.incubator.concurrent.StructuredTaskScope;
import java.util.concurrent.Future;

public class AggregatorFirstSuccess {
    public static void main(String[] args) throws Exception {
        try (var scope = new StructuredTaskScope.ShutdownOnSuccess<String>()) {
            Future<String> f1 = scope.fork(() -> fetchFromSource1());
            Future<String> f2 = scope.fork(() -> fetchFromSource2());
            Future<String> f3 = scope.fork(() -> fetchFromSource3());

            scope.join(); // 等到第一个成功者
            scope.throwIfFailed(); // 如果全部失败——抛出异常

            String result = scope.result(); // 第一个成功任务的结果
            System.out.println("第一个成功的结果:" + result);
        }
    }

    static String fetchFromSource1() { /* ... */ return "A"; }
    static String fetchFromSource2() { /* ... */ return "B"; }
    static String fetchFromSource3() { /* ... */ return "C"; }
}

发生了什么:

  • 一旦某个任务成功完成——其余任务将被取消。
  • 如果全部失败——会抛出异常。

4. 自动取消与降级

StructuredTaskScope 会在策略需要时自动取消其余任务。例如,当一个任务失败(ShutdownOnFailure)或一个任务成功完成(ShutdownOnSuccess),其余任务会收到取消信号(interrupt)。

示例:带超时的优雅收尾

import jdk.incubator.concurrent.StructuredTaskScope;
import java.time.Instant;
import java.util.concurrent.Future;

try (var scope = new StructuredTaskScope.ShutdownOnFailure()) {
    Future<String> f1 = scope.fork(() -> fetchWithTimeout());
    Future<String> f2 = scope.fork(() -> fetchWithTimeout());

    scope.joinUntil(Instant.now().plusSeconds(2)); // 最多等待 2 秒
    scope.throwIfFailed();

    String result = f1.resultNow() + f2.resultNow();
    System.out.println(result);
}

如果任务未在 2 秒内完成——将抛出异常,所有任务都会被取消。

5. 错误与异常处理

子任务的异常如何路由到 scope

有时在演出中总会有人打错音——StructuredTaskScope 不会装作什么也没发生。它会记录是谁“跑调”,随后把完整报告交给指挥。当你调用 throwIfFailed() 时,它会抛出一个聚合异常——类似一份汇总报告:“以下是今天走音的名单”。如有需要,你可以展开这棵“原因树”,看看究竟是谁拖了后腿。而如果你想只了解某位演奏者的情况——Future.exceptionNow() 会告诉你他的结局。

何时“取消”不等于“失败”

务必记住:取消任务并不总是错误。如果指挥说“演出结束”,乐手们只是收起乐器——这是 cancelled,而不是 failed。只有真正演错了才算错误,而这类异常会进入汇总。

示例:原因树

import jdk.incubator.concurrent.StructuredTaskScope;
import java.util.concurrent.Future;

try (var scope = new StructuredTaskScope.ShutdownOnFailure()) {
    Future<String> f1 = scope.fork(() -> { throw new RuntimeException("错误 1"); });
    Future<String> f2 = scope.fork(() -> { throw new RuntimeException("错误 2"); });

    scope.join();
    scope.throwIfFailed(); // 将抛出包含两个原因的异常
} catch (Exception e) {
    e.printStackTrace();
    // 可通过 e.getSuppressed() 获取 suppressed 异常
}

6. 与 CompletableFuture 的对比

StructuredTaskScopeCompletableFuture 都可以启动并行任务,但:

  • StructuredTaskScope 适用于任务在逻辑上相关、应当一起结束/取消的场景(任务层级)。
  • CompletableFuture 更适合没有层级的任务组合(例如转换链、响应式场景)。

何时 StructuredTaskScope 能简化代码:

  • 需要保证在离开代码块之前,所有子任务都已结束。
  • 需要集中化的取消与错误处理。
  • 需要确保不会留下“悬挂”的任务。

何时 CompletableFuture 更方便:

  • 任务彼此独立,可以各自运行。
  • 需要复杂组合(如 thenCombinethenCompose 等)。

7. 实战:HTTP 请求聚合器

任务:向 3 个来源发起请求,获取第一个成功的响应

import jdk.incubator.concurrent.StructuredTaskScope;
import java.util.concurrent.Future;

public class HttpAggregator {
    public static void main(String[] args) throws Exception {
        try (var scope = new StructuredTaskScope.ShutdownOnSuccess<String>()) {
            Future<String> f1 = scope.fork(() -> httpRequest("https://api1.example.com"));
            Future<String> f2 = scope.fork(() -> httpRequest("https://api2.example.com"));
            Future<String> f3 = scope.fork(() -> httpRequest("https://api3.example.com"));

            scope.join();
            scope.throwIfFailed();

            String result = scope.result();
            System.out.println("第一个成功的响应:" + result);
        }
    }

    static String httpRequest(String url) throws Exception {
        // 请求模拟(可使用 HttpClient)
        Thread.sleep((long) (Math.random() * 1000));
        if (Math.random() < 0.3) throw new RuntimeException("请求错误:" + url);
        return "来自 " + url + " 的响应";
    }
}

任务:若某个子任务失败——平滑地熄灭其余任务

import jdk.incubator.concurrent.StructuredTaskScope;
import java.util.concurrent.Future;

try (var scope = new StructuredTaskScope.ShutdownOnFailure()) {
    Future<String> f1 = scope.fork(() -> httpRequest("https://api1.example.com"));
    Future<String> f2 = scope.fork(() -> httpRequest("https://api2.example.com"));

    scope.join();
    scope.throwIfFailed();

    String result = f1.resultNow() + f2.resultNow();
    System.out.println("两个响应:" + result);
} catch (Exception e) {
    System.err.println("某个任务出错:" + e.getMessage());
}

8. 使用 StructuredTaskScope 的常见错误

错误 1:忘记调用 join()throwIfFailed()
如果不调用 join(),任务可能在离开代码块前仍未结束。如果不调用 throwIfFailed(),子任务的错误将被悄然忽略。

错误 2:在任务完成之前尝试获取结果。
在任务结束前调用 resultNow() 会抛出 IllegalStateException。请先通过 join() 等待完成。

错误 3:忽略取消信号。
如果任务已被取消(例如由于 scope 的策略),请不要尝试获取它的结果——否则会抛出异常。

错误 4:混用不同的终止策略。
不要在 scope 内手动取消任务——请使用 ShutdownOnFailureShutdownOnSuccess 策略。

错误 5:在虚拟线程中运行耗时的 CPU 密集型任务。
StructuredTaskScope 默认使用虚拟线程——它们非常适合 I/O 密集型任务,但不会加速沉重的计算。

错误 6:忘记关闭 scope(未使用 try-with-resources)。
StructuredTaskScope 实现了 AutoCloseable——务必使用 try-with-resources 来保证所有任务都能被正确结束。

评论
TO VIEW ALL COMMENTS OR TO MAKE A COMMENT,
GO TO FULL VERSION