登录
推荐 文章 Go 技术 课程 下载 专题 AI
首页 >  文章 >  java教程

Java Structured Concurrency 汇总子任务异常

来源:17golang原创

时间:2026-10-10 13:33:06 260浏览 收藏

Java Structured Concurrency 的默认异常策略适合“任何一个子任务失败,整组工作就没有继续价值”的场景:首个失败会让作用域取消其余任务,join() 抛出 ExecutionException,原因是一个失败子任务的异常。但批量校验、并行探测和多数据源健康检查经常需要另一种结果:让彼此独立的子任务全部结束,再一次性汇总所有失败。

在 JDK 27 的预览 API 中,可以用 StructuredTaskScope.Joiner.allUntil(subtask -> false) 构造“不因任务完成而提前取消”的策略。join() 返回全部 Subtask 句柄,随后按状态读取成功值或异常,最后把多个根因附加到一个业务聚合异常中。

官方 API:https://docs.oracle.com/en/java/javase/27/docs/api/java.base/java/util/concurrent/StructuredTaskScope.html

默认策略为什么看不到全部异常

最小的 Structured Concurrency 写法通常直接调用 StructuredTaskScope.open()。Java SE 27 文档说明,这等价于使用“全部成功,否则抛出”的 Joiner:所有任务成功时 join() 返回 null;任一任务失败时,作用域取消尚未完成的任务,并以第一个失败作为 ExecutionException 的 cause。

import java.util.concurrent.ExecutionException;
import java.util.concurrent.StructuredTaskScope;

static void failFast() throws InterruptedException {
    try (var scope = StructuredTaskScope.open()) {
        // 两个子任务属于同一个词法作用域
        scope.fork(() -> loadCustomer());
        scope.fork(() -> loadOrders());

        // 任一任务失败后,默认策略取消其余任务并抛出首个失败
        scope.join();
    } catch (ExecutionException e) {
        // cause 代表触发默认策略的一个失败,而不是所有失败的集合
        reportPrimaryFailure(e.getCause());
    }
}

这里的“看不到全部异常”并不是异常被偷偷吞掉,而是策略主动选择了短路。取消发生后,尚未完成或在取消后才结束的子任务可能处于 UNAVAILABLE,它们没有可读取的结果或异常。即使两个任务几乎同时出错,默认契约也只承诺把一个失败放进 ExecutionException。

StructuredTaskScope 默认 Joiner 与两个子任务、首个失败、ExecutionException 和被取消任务的静态关系
图1:默认失败策略结构图。作用域通过默认 Joiner 以首个失败为 cause 产生 ExecutionException,并可能让其他任务进入取消后的不可用状态;这是静态说明图,不是运行截图。

最小配方:让 allUntil 不触发提前取消

Joiner.allUntil(Predicate) 会在每个子任务完成后检查谓词;谓词返回 true 时取消作用域,否则继续等待。把谓词固定为 false,就能让没有超时的作用域等待全部独立子任务结束,并让 join() 返回所有 Subtask。

import java.util.List;
import java.util.concurrent.Callable;
import java.util.concurrent.StructuredTaskScope;

static List> collectAll(
        List> tasks) throws InterruptedException {

    var joiner = StructuredTaskScope.Joiner
            .allUntil(subtask -> false); // 恒假表示不因完成状态而取消

    try (var scope = StructuredTaskScope.open(joiner)) {
        for (var task : tasks) {
            // 每个 Callable 默认由独立虚拟线程执行
            scope.fork(task);
        }

        // 返回按 fork 顺序保存的全部 Subtask 句柄
        return scope.join();
    }
}

这里不要改用 allSuccessfulOrThrow()。它虽然能在全部成功时返回结果列表,但仍会在任一任务失败后取消作用域,并只传播一个失败。汇总全部异常的关键不是“换一个异常容器”,而是先选择不短路的完成策略。

join 之后才能按状态读取结果

Subtask 提供三种状态:SUCCESS、FAILED 和 UNAVAILABLE。官方 API 对读取方法有严格前置条件:只有成功状态才能调用 get(),只有失败状态才能调用 exception();作用域所有者还必须先完成 join()。

for (var subtask : subtasks) {
    switch (subtask.state()) {
        case SUCCESS -> {
            // 仅 SUCCESS 可以安全读取返回值
            values.add(subtask.get());
        }
        case FAILED -> {
            // 仅 FAILED 可以安全读取异常或 Error
            failures.add(subtask.exception());
        }
        case UNAVAILABLE -> {
            // 超时、取消或未启动时没有可读取的 outcome
            unavailableCount++;
        }
    }
}

如果对 FAILED 调用 get(),或者对 SUCCESS 调用 exception(),会得到 IllegalStateException,反而遮住真正的失败原因。因此状态判断不是可选的防御代码,而是 API 合约的一部分。

把多个失败放进一个可诊断异常

汇总异常时,最容易犯的错误是只拼接异常消息。消息字符串会丢掉异常类型、cause 链和堆栈。Java 自带的 suppressed exceptions 更合适:抛出一个代表整批失败的异常,把每个子任务根因通过 addSuppressed() 附加进去。

final class BatchFailure extends Exception {
    BatchFailure(List extends Throwable> causes) {
        super("并行任务失败数量:" + causes.size());

        // 保留每个根因的类型、消息和原始堆栈
        causes.forEach(this::addSuppressed);
    }
}

如果日志系统已经展开 suppressed exceptions,这个结构能直接保留全部堆栈。若还需要知道“哪个输入对应哪个异常”,可以先用带任务编号或业务键的包装异常包住根因,再作为 suppressed exception 附加。不要依赖完成顺序做关联,因为并发任务的完成顺序不是输入顺序。

完整写法:同时保留成功结果和所有失败

下面把策略、状态分类和聚合异常放到一个可复用方法中。示例假设所有任务返回同一种类型,并且任务彼此独立;只要存在失败,就抛出 BatchFailure,成功结果仍可在需要时改造成报告对象返回。

import java.util.ArrayList;
import java.util.List;
import java.util.concurrent.Callable;
import java.util.concurrent.StructuredTaskScope;

public final class StructuredBatch {

    public static List runAll(List> tasks)
            throws InterruptedException, BatchFailure {

        var joiner = StructuredTaskScope.Joiner
                .allUntil(subtask -> false); // 等待全部独立任务

        try (var scope = StructuredTaskScope.open(joiner)) {
            for (var task : tasks) {
                // fork 返回句柄,但 allUntil 会统一返回全部句柄
                scope.fork(task);
            }

            var subtasks = scope.join();
            var values = new ArrayList();
            var failures = new ArrayList();

            for (int index = 0; index  values.add(subtask.get());
                    case FAILED -> {
                        // 包装任务下标,保留输入与异常的关联
                        var cause = subtask.exception();
                        failures.add(new RuntimeException(
                                "子任务[" + index + "]失败", cause));
                    }
                    case UNAVAILABLE -> failures.add(
                            new IllegalStateException(
                                    "子任务[" + index + "]没有可用结果"));
                }
            }

            if (!failures.isEmpty()) {
                // 一次抛出整批异常,所有根因位于 suppressed 列表
                throw new BatchFailure(failures);
            }

            return List.copyOf(values);
        }
    }

    static final class BatchFailure extends Exception {
        BatchFailure(List extends Throwable> causes) {
            super("并行任务失败数量:" + causes.size());
            causes.forEach(this::addSuppressed); // 不把堆栈压成文本
        }
    }
}
Joiner allUntil 返回 Subtask 列表并按 SUCCESS FAILED UNAVAILABLE 汇总为结果集合和 BatchFailure
图2:全部结果汇总结构图。allUntil 返回 Subtask 列表,调用方按状态把成功值与失败根因分别汇入结果集合和 BatchFailure;这是静态说明图,不是运行截图。

编译和运行 JDK 27 预览 API 时需要显式启用预览特性:

# 编译时声明目标版本并启用预览 API
javac --enable-preview --release 27 StructuredBatch.java

# 运行时也必须启用预览特性
java --enable-preview StructuredBatch

需要短路时,不要照搬全部汇总

异常汇总不是默认策略的升级版,而是另一种业务取舍。以下场景更适合继续使用失败即取消:

  • 多个子任务共同组成一个请求,缺少任一结果就无法生成响应;
  • 后续任务昂贵,首个失败后继续执行只会浪费资源;
  • 子任务会修改外部状态,继续执行可能扩大部分成功带来的补偿成本;
  • 用户只关心请求成功或失败,不需要一次看到全部校验问题。

而批量表单验证、多端点健康检查、多个独立配置文件解析,通常值得等待全部结束。选择标准不是“想看更多日志”,而是剩余任务在一个失败发生后是否仍有业务价值。

兼容与工程坑点

Structured Concurrency 仍是预览 API

Java SE 27 中 StructuredTaskScope、Joiner 和 Subtask 仍标记为 preview。不要把这篇示例直接复制到锁定旧 JDK 的项目;升级时也要重新核对方法签名、异常类型和 Joiner 工厂方法。

超时会产生 UNAVAILABLE

allUntil 本身不会因为子任务失败或配置超时而自动抛出。若给作用域配置了超时,返回列表中可能包含 UNAVAILABLE。它不是“没有异常的失败”,而是没有可读取 outcome 的状态,应单独记录为超时或取消结果。

不要把 Error 一律降级成普通业务错误

Subtask.exception() 返回 Throwable,其中可能包含 Error。示例为了展示完整收集将其统一保留;生产代码应按容错策略决定是否立即重新抛出严重 Error,而不是把所有问题都包装成可重试的业务异常。

子任务必须能响应中断

即便使用全部汇总策略,作用域在所有者中断或关闭时仍会取消未完成任务。阻塞调用若不响应中断,会延迟作用域关闭。把网络和文件操作设置为有界等待,并正确传播中断,比单纯收集异常更重要。

选择策略速查

目标Joiner 选择异常表现
全部成功才继续默认 open() 或 awaitAllSuccessfulOrThrow()首个失败触发取消并抛出
得到全部成功值allSuccessfulOrThrow()任一失败则不返回结果列表
得到任一成功值anySuccessfulOrThrow()全部失败时抛出一个失败
汇总所有独立任务异常allUntil(subtask -> false)调用方检查每个 Subtask 状态
复杂取消和聚合规则自定义 Joiner实现必须线程安全并定义超时行为

常见问题

为什么不直接 catch ExecutionException 后读取其他子任务?

默认策略已经取消作用域,其他任务可能变成 UNAVAILABLE,并不保证每个任务都产生异常。要汇总全部独立失败,必须先采用不会在首个失败时短路的策略。

suppressed exceptions 和 cause 有什么区别?

cause 通常表示一个异常的直接成因;suppressed exceptions 适合附加同一聚合操作中并列发生的其他异常。汇总时保留一个业务异常作为顶层,再附加所有任务根因,日志结构更清楚。

allUntil 的谓词固定为 false 会不会无限等待?

它会等所有已 fork 的任务完成;如果某个任务永不返回,确实会一直等待。因此任务本身要使用超时或可中断调用,必要时再给作用域配置超时,并处理 UNAVAILABLE。

所有任务都成功时结果顺序稳定吗?

allUntil 返回的是该作用域内的 Subtask 列表,工程上仍应显式保存任务键或 fork 下标做关联,不要用完成先后推断输入身份。

Structured Concurrency 把一组并发任务变成一个有明确生命周期的工作单元,但“何时取消、返回什么、传播哪个异常”仍由 Joiner 策略决定。需要完整诊断时,先让独立任务全部形成 outcome,再按状态汇总;需要快速失败时,则保留默认短路策略。把业务目标写进 Joiner 选择,比事后从一个 ExecutionException 里猜测其余任务发生了什么更可靠。

声明:本文转载于:17golang原创 如有侵犯,请联系study_golang@163.com删除
相关阅读
更多>
最新阅读
更多>
课程推荐
更多>