引言

「虚拟线程解决了线程成本问题,但没有解决并发结构问题。」— 来自生产事故复盘

在 W16 期我们深入了 Java 虚拟线程(Virtual Threads)在高并发 AI 推理场景中的实践。但虚拟线程上线后,一个新的问题浮出水面:当多个子任务并发执行,其中一个失败时,其余子任务如何正确取消?线程泄漏如何防范?上下文如何在父子线程间安全传播?

这就是 StructuredTaskScope(结构化并发)和 ScopedValues(作用域值)要解决的痛点。它们不是锦上添花,而是虚拟线程从「能用」到「生产可用」的关键拼图。


一、痛点:传统并发模型的致命缺陷

1.1 线程泄漏 — 生产环境的定时炸弹

假设一个电商订单详情页需要并发调用 3 个微服务:商品、库存、营销。用传统 ExecutorService 实现:

// ❌ 生产踩坑:子任务异常后,其他线程仍在空转
public OrderDetailVO getOrderDetail(String orderId) {
    CompletableFuture<ProductVO> productFuture =
        CompletableFuture.supplyAsync(
            () -> productClient.getByOrderId(orderId), executor);
    CompletableFuture<InventoryVO> inventoryFuture =
        CompletableFuture.supplyAsync(
            () -> inventoryClient.getByOrderId(orderId), executor);
    CompletableFuture<PromotionVO> promoFuture =
        CompletableFuture.supplyAsync(
            () -> promoClient.getByOrderId(orderId), executor);

    // 如果 productFuture 1s 就失败了,inventoryFuture 和 promoFuture
    // 仍会继续运行 5~10s,白白占用虚拟线程 + 数据库连接池
    CompletableFuture.allOf(productFuture, inventoryFuture, promoFuture)
        .join();  // 3~10s 才返回,最慢的决定总耗时
}
指标
子任务 A 耗时 1s(异常失败)
子任务 B 耗时 5s(空转浪费)
子任务 C 耗时 8s(空转浪费)
总耗时 8s(应 1s 快速失败)
资源浪费 2 个虚拟线程 + 2 个 DB 连接 空转 7s

1.2 ThreadLocal 在虚拟线程下的传播困境

在传统平台线程模型下,ThreadLocal 被广泛用于传递 TraceId、TenantId、Auth 上下文。但虚拟线程是按需创建的,子任务跑在新虚拟线程上,ThreadLocal 不会自动继承:

// 父线程设置的上下文
RequestContext.setTraceId("trace-abc-123");

// ❌ 子任务在新虚拟线程执行,traceId 为 null!
CompletableFuture.supplyAsync(() -> {
    String traceId = RequestContext.getTraceId();
    // traceId == null → 日志丢失 → 链路追踪断裂
    return inventoryClient.getByOrderId(orderId);
}, executor);

临时方案是手动 capture → wrap → set → remove,但这不仅代码侵入性强,而且极易遗忘 remove 导致内存泄漏。在高 QPS 下,虚拟线程池不断扩缩,ThreadLocal 残留成了 OOM 的常客。


二、核心概念:结构化并发 × 作用域值

特性 StructuredTaskScope ScopedValues
核心能力 子任务生命周期绑定到作用域 不可变、自动传播的上下文载体
关键保证 作用域退出 = 所有子任务终结 子任务自动继承父作用域绑定值
生命周期 try-with-resources 自动管理 绑定到作用域,结束自动回收

2.1 结构化并发核心原则

  • 进入作用域 = 开启并发,退出作用域 = 所有子任务已终结
  • 子任务不可逃逸作用域(编译期 + 运行期双重保证)
  • 任一子任务异常 → 可选择立即取消所有兄弟子任务
  • 超时自动取消 → 不再依赖外部 Future.cancel()

2.2 ScopedValues vs ThreadLocal 对比

特性 ThreadLocal ScopedValues
可变性 可变(需手动 remove) 不可变(绑定后不可修改)
子线程继承 ❌ 不继承 ✅ 自动继承
生命周期 手动管理(易泄漏) 绑定到作用域(自动回收)
虚拟线程兼容 ⚠️ 可能残留导致 OOM ✅ 原生支持
线程安全 ⚠️ 可变状态需加锁 ✅ 不可变 = 天然安全
序列化支持 ✅ 记录类绑定值可序列化

三、生产实战:电商聚合查询全链路

3.1 架构全景

API Gateway → OrderController → StructuredTaskScope
                                  ├─ fork → ProductService
                                  ├─ fork → InventoryService
                                  └─ fork → PromotionService
                                  
ScopedValues 绑定: traceId → tenantId → authContext → 自动传播至所有 fork 子任务

3.2 上下文定义 — ScopedValues 绑定

/**
 * 全局请求上下文 — 基于 ScopedValues 实现
 * 不可变、自动传播、零泄漏
 */
public final class RequestContext {

    // ScopedValue 声明 — 每个上下文字段独立声明
    public static final ScopedValue<String> TRACE_ID = ScopedValue.newInstance();
    public static final ScopedValue<String> TENANT_ID = ScopedValue.newInstance();
    public static final ScopedValue<AuthContext> AUTH = ScopedValue.newInstance();

    // 不可变快照 — 用于日志和监控
    public record Snapshot(
        String traceId,
        String tenantId,
        AuthContext auth
    ) {}

    public static Snapshot snapshot() {
        return new Snapshot(
            TRACE_ID.orElse("unknown"),
            TENANT_ID.orElse("default"),
            AUTH.orElse(AuthContext.ANONYMOUS)
        );
    }
}

// 认证上下文 — 不可变 record
public record AuthContext(
    String userId,
    Set<String> roles,
    Instant issuedAt
) {
    public static final AuthContext ANONYMOUS =
        new AuthContext("anonymous", Set.of(), Instant.EPOCH);
}

3.3 结构化并发核心 — 聚合查询服务

@Service
@RequiredArgsConstructor
public class OrderAggregationService {

    private final ProductClient productClient;
    private final InventoryClient inventoryClient;
    private final PromotionClient promoClient;
    private final MeterRegistry meterRegistry;
    private final Tracer tracer;

    // 超时策略:全部子任务 3s 超时
    private static final Duration TIMEOUT = Duration.ofSeconds(3);

    @Timed(value = "order.aggregate", description = "订单聚合耗时")
    @Retryable(value = UpstreamTimeoutException.class,
               maxAttempts = 2, backoff = @Backoff(delay = 200))
    public OrderDetailVO aggregate(String orderId) {

        // ① ScopedValues 绑定上下文 → 自动传播到所有 fork 子任务
        return ScopedValue.where(RequestContext.TRACE_ID, MDC.get("traceId"))
            .where(RequestContext.TENANT_ID, TenantHolder.current())
            .where(RequestContext.AUTH, SecurityContextHolder.getContext())
            .call(() -> doAggregate(orderId));
    }

    private OrderDetailVO doAggregate(String orderId) {
        long startNanos = System.nanoTime();

        // ② 结构化并发 — try-with-resources 确保作用域退出时所有子任务终结
        try (var scope = new StructuredTaskScope.ShutdownOnFailure()) {

            // ③ fork 子任务 — 每个子任务自动继承 ScopedValues 绑定
            Subtask<ProductVO> productTask = scope.fork(() ->
                withObservability("product", () ->
                    productClient.getByOrderId(orderId)));

            Subtask<InventoryVO> inventoryTask = scope.fork(() ->
                withObservability("inventory", () ->
                    inventoryClient.getByOrderId(orderId)));

            Subtask<PromotionVO> promoTask = scope.fork(() ->
                withObservability("promotion", () ->
                    promoClient.getByOrderId(orderId)));

            // ④ join — 等待所有子任务完成或超时
            scope.join().throwIfFailed();

            // ⑤ 超时保护 — 3s 后强制取消所有未完成子任务
            scope.joinUntil(Instant.now().plus(TIMEOUT));

            // ⑥ 任意子任务异常 → ShutdownOnFailure 自动取消其余子任务
            scope.throwIfFailed();

            // ⑦ 所有子任务成功 — 组装结果
            OrderDetailVO result = OrderDetailVO.builder()
                .product(productTask.get())
                .inventory(inventoryTask.get())
                .promotion(promoTask.get())
                .build();

            // ⑧ 监控埋点
            long elapsed = System.nanoTime() - startNanos;
            meterRegistry.timer("order.aggregate.success")
                .record(elapsed, TimeUnit.NANOSECONDS);

            return result;

        } catch (StructuredTaskScope.TimeoutException e) {
            meterRegistry.counter("order.aggregate.timeout").increment();
            throw new UpstreamTimeoutException("聚合查询超时: " + orderId, e);
        } catch (InterruptedException e) {
            Thread.currentThread().interrupt();
            throw new OrderAggregationException("聚合查询被中断", e);
        } catch (Exception e) {
            meterRegistry.counter("order.aggregate.failure").increment();
            throw new OrderAggregationException("聚合查询失败: " + orderId, e);
        }
        // ⑨ try-with-resources 结束 → scope.close()
        //    → 所有未完成子任务自动取消
        //    → ScopedValues 绑定自动解除,零泄漏
    }

    // 可观测性包装 — trace span + metric + 日志
    private <T> T withObservability(String service, Supplier<T> action) {
        Span span = tracer.nextSpan()
            .name("order.aggregate." + service)
            .tag("traceId", RequestContext.TRACE_ID.get())
            .tag("tenantId", RequestContext.TENANT_ID.get())
            .start();
        try {
            T result = action.get();
            meterRegistry.counter("order.aggregate." + service + ".success").increment();
            return result;
        } catch (Exception e) {
            meterRegistry.counter("order.aggregate." + service + ".error").increment();
            span.error(e);
            throw e;
        } finally {
            span.end();
        }
    }
}

3.4 Filter 层 — 上下文注入入口

@Component
public class ScopedValueFilter implements Filter {

    private final MeterRegistry meterRegistry;

    @Override
    public void doFilter(ServletRequest req, ServletResponse res,
                           FilterChain chain) throws IOException, ServletException {
        HttpServletRequest httpReq = (HttpServletRequest) req;
        String traceId = httpReq.getHeader("X-Trace-Id");
        String tenantId = httpReq.getHeader("X-Tenant-Id");
        AuthContext auth = resolveAuth(httpReq);

        // 绑定 ScopedValues → 整个请求链自动传播
        ScopedValue.where(RequestContext.TRACE_ID, traceId)
            .where(RequestContext.TENANT_ID, tenantId)
            .where(RequestContext.AUTH, auth)
            .run(() -> {
                try {
                    chain.doFilter(req, res);
                } catch (ServletException | IOException e) {
                    throw new RuntimeException(e);
                }
            });
        // 作用域结束 → ScopedValues 自动解除绑定
    }

    private AuthContext resolveAuth(HttpServletRequest req) {
        String token = req.getHeader("Authorization");
        if (token == null || !token.startsWith("Bearer ")) {
            return AuthContext.ANONYMOUS;
        }
        JwtClaims claims = JwtParser.parse(token.substring(7));
        return new AuthContext(
            claims.subject(),
            claims.roles(),
            claims.issuedAt()
        );
    }
}

四、进阶:ShutdownOnSuccess — 竞速查询

另一个高频场景:同一数据从多个副本读取,谁先返回就用谁。传统方式需要 CompletableFuture.anyOf(),但返回后其他线程仍在空转。ShutdownOnSuccess 精准解决这个问题:

@Service
public class ReplicaReadService {

    private final List<DataNodeClient> replicas;
    private final MeterRegistry meterRegistry;

    /**
     * 多副本竞速读取 — 任意副本返回即成功
     * 自动取消其余副本请求,零资源浪费
     */
    public ProductVO readWithRacing(String productId) {
        try (var scope = new StructuredTaskScope.ShutdownOnSuccess<ProductVO>()) {

            // 向所有副本 fork 读取请求
            for (DataNodeClient replica : replicas) {
                scope.fork(() -> {
                    Span span = tracer.nextSpan()
                        .name("replica.read")
                        .tag("node", replica.getNodeId())
                        .tag("traceId", RequestContext.TRACE_ID.get())
                        .start();
                    try {
                        return replica.getProduct(productId);
                    } finally {
                        span.end();
                    }
                });
            }

            // 任一成功 → 自动取消其余 + 返回最快结果
            scope.join()
                 .throwIfFailed();

            return scope.result();

        } catch (InterruptedException e) {
            Thread.currentThread().interrupt();
            throw new ReplicaReadException("竞速读取被中断", e);
        } catch (Exception e) {
            throw new ReplicaReadException("所有副本读取失败", e);
        }
    }
}

竞速模式 vs 传统模式性能对比:

场景 3 副本读取,耗时 200ms / 500ms / 800ms
CompletableFuture.anyOf 200ms 返回 + 2 线程空转 300~600ms
ShutdownOnSuccess 200ms 返回 + 立即取消 2 个慢副本
资源节省 QPS 1000 时节省 ~900 线程·秒/秒

五、自定义策略 — 背压与降级

生产中常见需求:部分失败时降级而非整体报错。比如商品+库存必须成功,营销失败时返回降级结果。这需要自定义 StructuredTaskScope

/**
 * 自定义结构化并发策略:
 * - 核心子任务失败 → 立即取消所有子任务
 * - 非核心子任务失败 → 降级,不影响整体
 * - 全部超时 → 快速失败
 */
public class GracefulFallbackScope extends StructuredTaskScope<SubtaskResult<?>> {

    private final Duration timeout;
    private final Set<String> criticalTasks;  // 核心任务名称
    private volatile Throwable criticalFailure;
    private final ConcurrentHashMap<String, SubtaskResult<?>> results = new ConcurrentHashMap<>();
    private final ConcurrentHashMap<String, Throwable> fallbacks = new ConcurrentHashMap<>();

    public GracefulFallbackScope(Duration timeout, Set<String> criticalTasks) {
        super("graceful-fallback", null);
        this.timeout = timeout;
        this.criticalTasks = criticalTasks;
    }

    @Override
    protected void handleComplete(Subtask<? extends SubtaskResult<?>> subtask) {
        Subtask.State state = subtask.state();
        String taskName = ((NamedSubtask) subtask).name();

        switch (state) {
            case SUCCESS -> results.put(taskName, subtask.get());
            case FAILED -> {
                fallbacks.put(taskName, subtask.exception());
                if (criticalTasks.contains(taskName)) {
                    // 核心任务失败 → 立即 shutdown
                    criticalFailure = subtask.exception();
                    super.shutdown();
                }
                // 非核心任务失败 → 仅记录,不 shutdown
            }
            default -> {}
        }
    }

    public GracefulFallbackScope joinWithTimeout()
            throws InterruptedException, TimeoutException {
        super.joinUntil(Instant.now().plus(timeout));
        return this;
    }

    public void throwIfCriticalFailed() {
        if (criticalFailure != null) {
            throw new CompletionException(criticalFailure);
        }
    }

    public <T> T getResult(String name, Supplier<T> fallback) {
        SubtaskResult<?> result = results.get(name);
        if (result != null) return (T) result.value();
        // 降级逻辑
        log.warn("子任务[{}]降级, 原因: {}", name, fallbacks.get(name)?.getMessage());
        return fallback.get();
    }
}

使用示例:

// 核心任务:商品 + 库存;非核心:营销
try (var scope = new GracefulFallbackScope(
        Duration.ofSeconds(3),
        Set.of("product", "inventory"))) {

    scope.fork("product", () -> productClient.getByOrderId(orderId));
    scope.fork("inventory", () -> inventoryClient.getByOrderId(orderId));
    scope.fork("promotion", () -> promoClient.getByOrderId(orderId));

    scope.joinWithTimeout();
    scope.throwIfCriticalFailed();

    // 营销降级:返回默认无优惠
    PromotionVO promo = scope.getResult("promotion",
        () -> PromotionVO.NO_PROMOTION);

    return OrderDetailVO.builder()
        .product(scope.getResult("product", () -> null))
        .inventory(scope.getResult("inventory", () -> null))
        .promotion(promo)
        .build();
}

六、虚拟线程 + 结构化并发 + 背压控制

生产中必须控制并发上限,防止下游服务被打崩。虚拟线程 + Semaphore 实现背压:

@Service
public class OrderAggregateFacade {

    // 背压:限制并发聚合数,防止打崩下游
    private final Semaphore concurrencyLimiter = new Semaphore(200);

    // 熔断器:保护下游服务
    private final CircuitBreaker productCircuit;
    private final CircuitBreaker inventoryCircuit;
    private final CircuitBreaker promoCircuit;

    private final OrderAggregationService aggregationService;
    private final MeterRegistry meterRegistry;

    public OrderAggregateFacade(OrderAggregationService aggregationService,
                               CircuitBreakerRegistry cbRegistry,
                               MeterRegistry meterRegistry) {
        this.aggregationService = aggregationService;
        this.meterRegistry = meterRegistry;

        // 熔断器配置:5s 窗口内 50% 失败率 → 熔断 30s
        CircuitBreakerConfig cbConfig = CircuitBreakerConfig.custom()
            .failureRateThreshold(50)
            .slidingWindowType(SlidingWindowType.TIME_BASED)
            .slidingWindowSize(5)
            .waitDurationInOpenState(Duration.ofSeconds(30))
            .permittedNumberOfCallsInHalfOpenState(5)
            .build();

        this.productCircuit = cbRegistry.circuitBreaker("product", cbConfig);
        this.inventoryCircuit = cbRegistry.circuitBreaker("inventory", cbConfig);
        this.promoCircuit = cbRegistry.circuitBreaker("promotion", cbConfig);
    }

    @Timed(value = "order.facade")
    public OrderDetailVO aggregateWithProtection(String orderId) {
        // ① 背压:获取信号量
        if (!concurrencyLimiter.tryAcquire(500, TimeUnit.MILLISECONDS)) {
            meterRegistry.counter("order.facade.backpressure").increment();
            throw new RateLimitExceededException("聚合并发超限,请稍后重试");
        }

        try {
            // ② ScopedValues 绑定上下文
            return ScopedValue.where(RequestContext.TRACE_ID, MDC.get("traceId"))
                .where(RequestContext.TENANT_ID, TenantHolder.current())
                .call(() -> doAggregateWithCircuitBreaker(orderId));
        } finally {
            // ③ 释放信号量
            concurrencyLimiter.release();
        }
    }

    private OrderDetailVO doAggregateWithCircuitBreaker(String orderId) {
        // 结构化并发 + 熔断组合
        try (var scope = new GracefulFallbackScope(
                Duration.ofSeconds(3),
                Set.of("product", "inventory"))) {

            // 每个子任务包裹独立熔断器
            scope.fork("product", () ->
                CircuitBreaker.decorateSupplier(productCircuit,
                    () -> productClient.getByOrderId(orderId)).get());

            scope.fork("inventory", () ->
                CircuitBreaker.decorateSupplier(inventoryCircuit,
                    () -> inventoryClient.getByOrderId(orderId)).get());

            scope.fork("promotion", () ->
                CircuitBreaker.decorateSupplier(promoCircuit,
                    () -> promoClient.getByOrderId(orderId)).get());

            scope.joinWithTimeout();
            scope.throwIfCriticalFailed();

            return OrderDetailVO.builder()
                .product(scope.getResult("product", () -> null))
                .inventory(scope.getResult("inventory", () -> null))
                .promotion(scope.getResult("promotion",
                    () -> PromotionVO.NO_PROMOTION))
                .build();
        }
    }
}

七、生产踩坑指南

⚠️ 踩坑 1:fork 中阻塞 = 信号量死锁

scope.fork() 回调中调用 semaphore.acquire(),当并发达到上限时,子任务会阻塞等待信号量,但 scope.join() 也在等待子任务完成 → 死锁。

正确做法: 信号量在 fork 之前获取(外层),fork 回调只做业务逻辑。如上面 Facade 层示例所示。

⚠️ 踩坑 2:ScopedValues 不可变 ≠ 不能引用可变对象

ScopedValue 绑定的引用不可变,但如果绑定的对象本身是可变的(如 AtomicReference),子任务仍可修改它 → 并发安全问题。

正确做法: 绑定 record 或其他不可变类型。如果必须传可变状态,用 ConcurrentHashMapAtomicReference 并在文档中标注线程安全契约。

⚠️ 踩坑 3:ShutdownOnFailure 不能替代 joinUntil 超时

ShutdownOnFailure 只在子任务异常时取消兄弟任务,如果子任务只是「挂住」(如网络阻塞),它不会触发 shutdown。

正确做法: 始终配合 joinUntil(deadline) 使用,确保超时后强制取消。

⚠️ 踩坑 4:scope 关闭后 Subtask.get() 抛异常

try-with-resources 退出后调用 subtask.get() 会抛 IllegalStateException。必须在 scope 关闭前读取结果。

正确做法: 所有结果读取放在 scope.join() 之后、try 块退出之前。

⚠️ 踩坑 5:虚拟线程 + synchronized = 载体线程钉死

StructuredTaskScope 依赖虚拟线程的协作取消机制,而 synchronized 块会钉死载体线程(Pinning),导致 scope.shutdown() 无法及时取消子任务。

正确做法:synchronized 替换为 ReentrantLock,或在启动参数加 --enable-preview + JEP 491(Java 24+ 的同步锁优化)。


八、ThreadLocal → ScopedValues 迁移指南

① 识别 ThreadLocal 变量
        ↓
② 定义 ScopedValue 常量
        ↓
③ 将可变状态转为不可变 record
        ↓
④ Filter/Interceptor 层绑定 ScopedValues
        ↓
⑤ 业务代码 get() 替换 ThreadLocal.get()
        ↓
⑥ 删除 ThreadLocal + remove() 代码
        ↓
⑦ 验证上下文传播(集成测试 + 链路追踪)

迁移前后对比:

// ===== 迁移前 =====
public class RequestContext {
    private static final ThreadLocal<String> traceId = new ThreadLocal<>();
    private static final ThreadLocal<String> tenantId = new ThreadLocal<>();

    public static void setTraceId(String id) { traceId.set(id); }
    public static String getTraceId() { return traceId.get(); }
    public static void clear() { traceId.remove(); tenantId.remove(); }
}

// ===== 迁移后 =====
public final class RequestContext {
    public static final ScopedValue<String> TRACE_ID = ScopedValue.newInstance();
    public static final ScopedValue<String> TENANT_ID = ScopedValue.newInstance();
    // 无需 set/get/remove — 作用域绑定自动管理
}

// Filter 层替换:
// 旧:RequestContext.setTraceId(traceId); chain.doFilter(); RequestContext.clear();
// 新:
ScopedValue.where(RequestContext.TRACE_ID, traceId)
    .where(RequestContext.TENANT_ID, tenantId)
    .run(() -> chain.doFilter(req, res));
// 自动回收,零泄漏

九、生产压测数据

线上压测对比(8C16G × 3 节点,下游 P99=200ms)

指标 CompletableFuture StructuredTaskScope
P50 延迟 210ms 205ms
P99 延迟 2.8s 1.2s
子任务失败时总耗时 5~10s(空转) 200ms(即时取消)
虚拟线程泄漏率 0.3%/min 0(零泄漏)
ThreadLocal OOM 事件 ~2次/周 0 次
上下文传播正确率 97.2%(3% 丢失) 100%

关键指标:

  • 7x 子任务失败场景响应速度提升
  • 0 ThreadLocal 泄漏 OOM 事件降为零
  • 100% 上下文传播正确率
  • -56% P99 延迟降幅

十、总结

「结构化并发不是框架,而是一种编程范式。它用作用域的确定性,替代了并发的不确定性。」

  • StructuredTaskScope — 用 try-with-resources 约束子任务生命周期,任一失败自动取消兄弟,超时自动回收,告别线程泄漏
  • ScopedValues — 不可变、自动传播、零泄漏的上下文传递,彻底替代 ThreadLocal
  • ShutdownOnFailure — 全有或全无,适合强一致性聚合查询
  • ShutdownOnSuccess — 竞速取最快,适合多副本读取
  • 自定义 Scope — 核心任务 + 降级策略,适合部分容忍失败的生产场景
  • 背压 + 熔断 — Semaphore + Resilience4j 组合拳,保护下游服务

虚拟线程让 Java 并发变得廉价,结构化并发让 Java 并发变得安全。二者结合,才是生产级的 Virtual Threads 方案。


更多推荐