上下文、CompletableFuture、ForkJoin 与虚拟线程
一次异步调用同时涉及三份信息:任务在哪里执行,执行时能读取什么上下文,以及结果何时交还给谁。Executor 安排执行,ThreadLocal 或显式参数提供上下文,Future 与 CompletableFuture 表示结果。它们可以组合,但不会自动替彼此完成传播和清理。
提交者 ── 捕获只读数据 ──→ 工作任务
│ │
└── 等待或注册回调 ←── 完成结果
│
异常 / 超时 / 取消平台线程池复用线程,虚拟线程通常每任务新建。两种执行方式都需要定义下游调用预算与结果处理,区别在于线程占用和调度成本。
上下文的保存、安装与恢复
ThreadLocal 值存在线程内部
ThreadLocal 为当前线程保存一个值。线程池 worker 结束一次请求后仍会存活,未清除的租户或大对象会随该线程保留。请求开始时安装上下文,离开时清理;存在嵌套调用时,则需要恢复外层值。
下面的 TenantContext 拒绝 null 租户,以 null 表示未绑定,并提供嵌套恢复与提交时捕获:
import java.util.Objects;
import java.util.function.Supplier;
public final class TenantContext {
private static final ThreadLocal<String> CURRENT = new ThreadLocal<>();
private TenantContext() { }
public static String current() { return CURRENT.get(); }
public static <T> T with(String tenant, Supplier<T> action) {
Objects.requireNonNull(tenant);
String old = CURRENT.get();
try { CURRENT.set(tenant); return action.get(); }
finally { if (old == null) CURRENT.remove(); else CURRENT.set(old); }
}
public static <T> Supplier<T> capture(Supplier<T> action) {
String captured = Objects.requireNonNull(current(), "tenant not bound");
return () -> with(captured, action);
}
public static void main(String[] args) {
String result = with("tenant-42", () -> {
String inner = with("tenant-7", TenantContext::current);
if (!"tenant-7".equals(inner) || !"tenant-42".equals(current()))
throw new AssertionError("restore");
return current();
});
if (current() != null) throw new AssertionError("leak");
System.out.println("tenant=" + result + " nested-restored=true cleared=true");
}
}Linux 普通用户保存为 TenantContext.java,使用完整 JDK。安装入口见Java 版本基线。设置实际 JAVA_HOME,公共示例编译为 Java 17:
export JAVA_HOME=/opt/jdk-25
unset JAVA_TOOL_OPTIONS JDK_JAVA_OPTIONS _JAVA_OPTIONS
LAB_OUT=$(mktemp -d /tmp/async-learning.XXXXXX)
"$JAVA_HOME/bin/javac" --release 17 -Xlint:all -Werror -d "$LAB_OUT" TenantContext.java
"$JAVA_HOME/bin/java" -cp "$LAB_OUT" TenantContext输出 tenant=tenant-42 nested-restored=true cleared=true。内层绑定 tenant-7,退出后恢复 tenant-42;最外层退出清空。finally 也覆盖业务异常。该组件使用默认 null 初始值,因此可以用 null 决定 remove;若项目允许显式 null 或自定义 initialValue,需要额外保存“是否存在绑定”。
capture 在提交线程读取值,返回的 Supplier 在执行线程安装并恢复。捕获的是字符串引用;对可变上下文对象应构造只读快照,不共享尚在修改的请求对象。租户来源还需要在入口完成身份校验,ThreadLocal 本身不提供授权验证。
弱 key 与强 value 的存活关系
固定 OpenJDK jdk-25+36 ThreadLocal 实现中,ThreadLocalMap 的 Entry 用弱引用保存 key,强引用保存 value。key 被回收只会留下陈旧项,value 仍可能沿线程引用链存活,直到访问触发清理或线程结束。
长寿命 worker → ThreadLocalMap → Entry → value
└─ 弱 key:可能已经被回收内部表使用开放寻址。清理陈旧槽位还要调整探测链,后续 get/set/remove 并不保证立刻遍历清空所有陈旧值。显式 remove 才是组件可以控制的释放动作。ThreadLocal 的使用契约见 API。
未清理的有效 key 会导致请求串值;key 已回收但 value 留存则可能造成内存滞留。两者引用形态不同,处理都应落到正确生命周期,而不是等待 GC 偶然清扫。
真实跨线程与异常清理
下载异步实验包,解压进入 context-future-virtual,沿用 LAB_OUT:
"$JAVA_HOME/bin/javac" --release 17 -Xlint:all -Werror -d "$LAB_OUT" *.java
"$JAVA_HOME/bin/java" -cp "$LAB_OUT" AsyncLab上下文部分输出:
unwrapped=old-tenant cleanup=true
captured=tenant-42 exceptional-exit-cleared=true第一个负例连续在同一个 worker 上运行两个任务,第二个读到第一个遗留值,然后明确清理。正例在提交时捕获 tenant-42,由 worker 安装;异常任务退出后,再提交查询断言该线程上下文为空。
InheritableThreadLocal 在创建子线程时继承值,而非每次提交时继承。早已存在的池线程不会因此自动刷新;默认继承也可能继续共享可变对象。定义见 InheritableThreadLocal API。显式参数适合少量清晰依赖,ThreadLocal 适合兼容需要隐式访问的框架,二者都不应携带跨线程共用的数据库事务连接。
CompletableFuture 连接完成阶段
先选择依赖关系
CompletableFuture 保存结果,并连接“前置完成后执行什么”。它不是线程池,也不保证一条链始终在同一线程运行。
| 操作 | 表达的依赖 |
|---|---|
| thenApply | 将一个值转换为另一个值 |
| thenAccept / thenRun | 消费结果或仅在完成后执行动作 |
| thenCompose | 将返回的异步阶段展开,避免嵌套 Future |
| thenCombine | 两个独立阶段正常完成后组合结果 |
| allOf | 等所有输入完成,返回 Void 阶段,结果仍从各输入读取 |
| anyOf | 任一输入完成便完成,可能是异常,不是“首个成功” |
详细规则见 CompletableFuture API。allOf 不负责聚合业务结果,也不保证某个分支失败后立即取消其他任务;anyOf 得到结果后,其他分支仍可能运行。空集合时 allOf 已完成,anyOf 则保持未完成,动态扇出需要处理零分支。
下面的组合片段依赖现有 client 和 executor:
CompletableFuture<User> user =
CompletableFuture.supplyAsync(() -> client.loadUser(id), userExecutor);
CompletableFuture<List<Order>> orders =
CompletableFuture.supplyAsync(() -> client.loadOrders(id), orderExecutor);
CompletableFuture<View> view = user.thenCombine(orders, View::new);两个加载独立时可以并行;若订单查询需要用户结果,应以 thenCompose 建立真实依赖。仅把阻塞 get 写进回调会隐藏等待资源,可能重现线程池同池饥饿。
回调在哪个线程执行
非 Async 阶段可能由完成前置阶段的线程执行;注册时前置已经完成,也可能由注册者执行。Async 阶段交给相应 Executor,未显式指定时通常使用 commonPool;commonPool 并行度不足 2 时,默认实现有创建新线程的后备策略。
Async 表示通过执行设施安排动作,若自定义 Executor 本身直接调用 Runnable.run,仍会在当前线程执行。因此,执行器是否真正转线程也需要看具体契约。
AsyncLab 注册两个阶段,再由 main 调用 complete,输出:
inline=value@main async=value@owned-executor这是受控条件下的执行线程,不是 thenApply 永远在 main。非 Async 回调应避免长阻塞,否则完成网络结果的 I/O 线程可能被后续业务占住。上下文也必须在真正执行的回调处安装,不能只包装 supplyAsync 的第一段。
固定 OpenJDK jdk-25+36 CompletableFuture保存结果与依赖完成栈。注册动作与发布结果可能竞争,完成方或注册方会帮助推进依赖。节点组织属于实现,应用应依赖完成关系而非猜测栈遍历顺序。
正常、异常与观察
exceptionally 将异常转换为正常替代值,handle 同时处理成功或失败,whenComplete 主要用于观察并继续传播。观察回调自己抛异常也会影响下游;记录日志时不要意外覆盖成功结果或改变主要故障语义。
get 可响应中断,并通过 ExecutionException 交还计算失败;join 通过 CompletionException 报告失败,不以 InterruptedException 退出等待。取消还有 CancellationException,统一解包时应保留原始 cause,不能把全部异常转为 null 后当作“无数据”。
多个分支失败时,组合结果未必保存全部错误。关键分支分别记录结果与耗时,汇总层再按业务决定缺一不可还是允许部分降级。共享给只需订阅结果的调用方时,可以暴露 CompletionStage 或受限视图,减少外部随意 complete 的权限。
结果超时与底层任务的不同终点
Future.get(timeout) 只限制等待者。orTimeout 让同一个 CompletableFuture 异常完成,completeOnTimeout 则用默认值正常完成;二者都会影响共享该对象的其他订阅者,并不会创建独立的等待视图。
CompletableFuture.cancel(true) 的 mayInterruptIfRunning 参数在这个实现中不控制线程中断。它改变完成状态并影响尚未完成的依赖阶段,不能借此保证 supplyAsync 内的 HTTP 或 SQL 已被终止。
AsyncLab 启动受闩锁阻挡的供给任务,确认任务已进入后分别 cancel 和 orTimeout。结果已完成时,任务仍停在闩锁;最后由测试释放任务并等待退出:
cancel result-done-before-work=true work-released=true
orTimeout result-done-before-work=true work-released=true这组实验同时观察 Future 状态和工作任务终点,避免只看 isDone。生产中需要给底层客户端设置连接、读取和整体请求预算,保留可用的中止句柄;共享加载的某个等待者超时,也不能随意取消其他等待者仍需要的工作。
扇出会放大下游并发。一次入口请求启动十个分支时,应分配同一总截止预算,而非每层重新设置完整超时。执行器排队、获取连接和重试都消耗该预算;已经产生副作用的超时结果要按业务幂等和状态查询处理。
ForkJoin 与共享执行资源
ForkJoinPool 主要通过 work-stealing 调度任务:worker 从自己的工作队列处理任务,空闲时尝试从其他队列取得工作,适合能够拆成足够粗粒度子任务的计算。分解过细会让调度与对象分配超过计算收益。
RecursiveTask 返回值,RecursiveAction 不返回值。fork 安排子任务,join 获取结果;在适合的 ForkJoin 计算中,等待方可以协助执行任务。这个机制不保证所有外部阻塞都能被自动补偿。大量 JDBC、网络等待或任意 monitor 阻塞仍可能耗尽有效执行能力。
ManagedBlocker 给池提供可管理阻塞的配合入口,需要正确实现“是否已可继续”和“实际阻塞”两部分。它有助于调整 worker,却不创造数据库连接或无限 CPU。具体契约见 ForkJoinPool API。
默认 supplyAsync、并行流等可能共享 commonPool。一个组件阻塞大量公共 worker 会影响其他组件,核心阻塞调用应使用有容量设计的独立执行器或虚拟线程模型。不要为了一个模块直接改全进程 common parallelism;先观察实际调用者与等待位置。
JDK 25 的 ForkJoinPool 还实现 ScheduledExecutorService,增加延迟与周期任务方法,不能把旧版“ForkJoinPool 不支持调度”作为跨版本结论。新 API 不改变共享池的资源隔离问题,Java 17 源码也不能直接调用这些方法。
虚拟线程与 ScopedValue
每任务一个虚拟线程,限制的是下游资源
虚拟线程在 Java 21 成为正式 API。它仍是 Thread,在支持卸载的阻塞操作中暂停后,载体线程可以执行其他虚拟线程;计算期间仍消耗载体 CPU。适合大量等待型任务,不会直接加速一段 CPU 计算。
Executors.newVirtualThreadPerTaskExecutor 为每个任务创建虚拟线程。通常不再池化虚拟线程以限制数量,而是在数据库、远端调用等入口使用许可或连接池限制并发,并限制等待者积压。指南见 Java 25 虚拟线程。
JDK 24 的 monitor 改进使 synchronized 中可卸载的阻塞不再仅因持锁固定载体,具体同步语义见Monitor 与 synchronized。native、外部函数及实际库路径仍需测量。锁本身的互斥没有消失,持锁慢调用仍会阻挡同锁任务。
虚拟线程结束可以缩短每任务 ThreadLocal 的存活期,但为大量线程各建一个重型连接或缓存仍昂贵。连接池应按共享下游容量管理,不用 ThreadLocal 为每个虚拟线程缓存一条数据库连接。
只读绑定的作用域
Java 25 的 ScopedValue 是正式 API,绑定在限定的调用范围内可见,正常或异常退出后恢复外层绑定。它适合租户、跟踪标识等只读上下文;绑定不可变不意味着对象自动深度不可变,跨线程共享的值仍应不可变或有同步。
private static final ScopedValue<String> TENANT = ScopedValue.newInstance();
ScopedValue.where(TENANT, "tenant-42").run(() -> {
service.handle(TENANT.get());
});片段需要 Java 25。普通线程池 submit 不会自动继承这个绑定,实验在每个虚拟任务内显式 where;嵌套绑定结束后恢复外层。规则见 ScopedValue API。
StructuredTaskScope 将相关子任务限制在父任务作用域内,提供统一完成、失败与取消策略,并可继承受控作用域中的 ScopedValue。它在 Java 25 仍为第五次 Preview,使用必须开启预览,并按该版本 API 编译;不能照搬早期 ShutdownOnFailure 示例。状态见 StructuredTaskScope API。稳定生产基线未采用预览时,仍可用显式执行器和任务拥有者实现关闭与取消协议。
运行真实 Java 25 任务
实验包的 java25/VirtualContextLab.java 创建 12 个虚拟任务,在每个任务内绑定 tenant-42,以 Semaphore 限制最多 3 个同时进入模拟下游。任务结束归还许可,另测嵌套绑定与异常退出恢复。使用 Java 25 编译并运行:
"$JAVA_HOME/bin/javac" --release 25 -Xlint:all -Werror -d "$LAB_OUT" java25/VirtualContextLab.java
"$JAVA_HOME/bin/java" -cp "$LAB_OUT" VirtualContextLab输出:
virtual-tasks=12 peak-within-3=true permits-restored=3 scoped-unbound=true每个任务都断言 isVirtual,许可峰值按实际上界判断,不要求调度恰好达到 3。3 是示例配额,实际取值来自下游容量。该实验无需 --enable-preview;使用 Java 17 编译会失败,应保持公共源码与 java25 目录分开。
示例用 try-with-resources 关闭执行器,close 会等待任务终止。即使某次 Future.get 有超时,离开 try 时仍可能在 close 等待,不能据此保证方法按 get 的时限返回。任务必须有自己的有界操作与取消协议。
异步调用的排障与迁移
| 现象 | 先区分的关系 | 修复后验证 |
|---|---|---|
| 同线程后续请求读到旧租户 | 入口清理、嵌套恢复、异常路径 | 同 worker 连续任务不串值 |
| 第一阶段有上下文,后续回调没有 | 每个回调的实际执行器与捕获时机 | 正常和异常完成线程都能读取正确值 |
| Future 超时但连接仍被占用 | 结果状态与底层调用是否分开结束 | 截止后在途调用、连接占用实际回落 |
| commonPool 停滞 | 阻塞分支与共享使用者 | 分离阻塞负载后无关计算恢复 |
| 虚拟线程很多、吞吐不增 | CPU、连接池或许可等待 | 相同负载下延迟和下游错误共同改善 |
| 内存随任务增长 | 上下文载荷、队列、每线程缓存 | 任务结束后引用释放,等待总量受限 |
迁移时保持相同入口速率、任务类型、超时与下游配额,比较平台池与虚拟线程的吞吐、尾延迟、CPU、堆和连接等待。只有线程创建更轻,不能支撑无限接收任务的设计。活跃目标的 JSON 线程转储与 JFR 采集见并发诊断。
完整实验与清理
公共实验在解压目录执行 bash run.sh;Java 25 全部实验使用 bash run.sh java25。两个运行版本的源码分别严格编译,脚本中各 Java 进程有 20 秒外部上限。
Linux 主机有 Docker 权限时:
LAB_SOURCE="$(pwd -P)"
JDK_IMAGE=eclipse-temurin:25.0.4_7-jdk
docker pull "$JDK_IMAGE"
docker run --pull never --rm --network none --read-only \
--user 10001:10001 --cap-drop ALL --security-opt no-new-privileges \
--mount "type=bind,src=$LAB_SOURCE,dst=/lab,readonly" \
--tmpfs /tmp:rw,nosuid,nodev,size=128m,mode=1777 \
"$JDK_IMAGE" bash /lab/run.sh java25容器 UID 10001 读取源码,编译写入 /tmp,无需网络。离线镜像提前从可信来源校验并导入。Java 17 对照镜像为 eclipse-temurin:17.0.20_8-jdk,运行时去掉末尾 java25,只执行公共实验。最终输出 PASS async-lab;任何断言或超时都要按失败处理。
确认目标退出后回收手工目录:
case "$LAB_OUT" in
/tmp/async-learning.*) rm -r -- "$LAB_OUT" ;;
*) printf '%s\n' '保留未知目录' ;;
esac
unset LAB_OUT权威资料与规范地址
按方法语义查 API,按实现字段查固定版本源码。
| 资料 | 完整地址 |
|---|---|
| OpenJDK jdk-25+36 ThreadLocal 实现 | https://github.com/openjdk/jdk/blob/jdk-25%2B36/src/java.base/share/classes/java/lang/ThreadLocal.java |
| API | https://docs.oracle.com/en/java/javase/25/docs/api/java.base/java/lang/ThreadLocal.html |
| InheritableThreadLocal API | https://docs.oracle.com/en/java/javase/25/docs/api/java.base/java/lang/InheritableThreadLocal.html |
| CompletableFuture API | https://docs.oracle.com/en/java/javase/25/docs/api/java.base/java/util/concurrent/CompletableFuture.html |
| OpenJDK jdk-25+36 CompletableFuture | https://github.com/openjdk/jdk/blob/jdk-25%2B36/src/java.base/share/classes/java/util/concurrent/CompletableFuture.java |
| ForkJoinPool API | https://docs.oracle.com/en/java/javase/25/docs/api/java.base/java/util/concurrent/ForkJoinPool.html |
| Java 25 虚拟线程 | https://docs.oracle.com/en/java/javase/25/core/virtual-threads.html |
| ScopedValue API | https://docs.oracle.com/en/java/javase/25/docs/api/java.base/java/lang/ScopedValue.html |
| StructuredTaskScope API | https://docs.oracle.com/en/java/javase/25/docs/api/java.base/java/util/concurrent/StructuredTaskScope.html |
