生产级高并发实战

WEEKLY TECH · W16 · 2026
结构化并发 / 背压控制 / 熔断降级 / TraceId 传播 · 五个真实踩坑场景


01 AI 推理为什么特别适合虚拟线程?

一次大模型推理请求,真正花在 CPU 上的时间不超过 5%——剩下 95% 都在等:等 HTTP 响应、等 Token 流式输出、等数据库里的 Embedding 查询回来。

这种极端 I/O 密集的场景,OS 线程就是在烧钱。JDK 21 的虚拟线程(Virtual Thread)让 JVM 自己来调度挂起/恢复,一台 8 核机器可以同时挂起 十万级 虚拟线程,内存开销仅 ~1KB/个。

指标 OS 线程 虚拟线程
栈内存 ~1MB ~1KB
AI 推理 I/O 占比 - 95%
吞吐量提升 - ×14

⚠️ 误区提醒:虚拟线程不是银弹。CPU 密集型计算(如矩阵运算、图像处理)用虚拟线程没有收益,反而可能因调度开销略有损耗,请继续使用线程池 + ForkJoinPool。


02 生产架构:网关 → 并发推理 → 聚合

下图是一个真实的 AI 推理网关架构,每个请求需要并发调用 3 个模型(主力模型 + 两个备选),取最快返回的有效结果:

┌─────────────┐     ┌─────────────────┐     ┌─────────────────┐
│ 客户端请求   │ ──→ │ Spring Gateway  │ ──→ │ AI Inference    │
└─────────────┘     │ 限流 / 鉴权      │     │ Service         │
                    └─────────────────┘     └────────┬────────┘
                                                      │
              StructuredTaskScope.ShutdownOnSuccess 并发调用
                                                      ↓
        ┌─────────────┐    ┌─────────────┐    ┌─────────────┐
        │   GPT-4o    │    │ Claude 3.5  │    │ Gemini 1.5  │
        │   主力模型   │    │   备选 A    │    │   备选 B    │
        └─────────────┘    └─────────────┘    └─────────────┘
                                                      │
                             最快有效结果 → 取消其余调用
                                                      ↓
        ┌─────────────────┐    ┌─────────────┐    ┌─────────────┐
        │ 结果聚合 +      │ ──→│ Prometheus  │ ──→│ SSE 流式    │
        │ TraceId 注入    │    │ 指标上报     │    │ 返回客户端   │
        └─────────────────┘    └─────────────┘    └─────────────┘

03 结构化并发:谁先回来用谁,其余自动取消

JDK 21 引入的 StructuredTaskScope 是虚拟线程配套的"作用域线程管理"。ShutdownOnSuccess 策略:任何一个子任务成功返回,立即取消其他所有子任务——完美匹配"多模型竞速取最快"场景。

@Service
@Slf4j
public class MultiModelRaceService {

    private final List<AiInferenceClient> clients; // GPT / Claude / Gemini
    private final MeterRegistry   meterRegistry;
    private final Semaphore       globalSemaphore;   // 全局背压

    public MultiModelRaceService(List<AiInferenceClient> clients,
                                     MeterRegistry meterRegistry,
                                     @Value("${ai.max-concurrent:200}") int maxConcurrent) {
        this.clients        = clients;
        this.meterRegistry   = meterRegistry;
        // 最多同时 200 个推理请求,超出进入等待而不是拒绝
        this.globalSemaphore = new Semaphore(maxConcurrent, true);
    }

    /**
     * 并发调用多个模型,返回最快的有效结果。
     * 失败 / 超时的子任务不影响最终结果,全部失败才抛异常。
     */
    public InferenceResult raceInference(InferenceRequest req)
            throws InterruptedException {

        String traceId = MDC.get("traceId"); // 提前捕获,虚拟线程不自动继承
        Timer.Sample sample = Timer.start(meterRegistry);

        // ① 背压:获取令牌(最多等 3 秒,否则降级)
        boolean acquired = globalSemaphore.tryAcquire(3, TimeUnit.SECONDS);
        if (!acquired) {
            meterRegistry.counter("ai.inference.backpressure").increment();
            return fallbackResult(req, "backpressure");
        }

        try (
            // ② 结构化并发 —— 谁先成功就关闭其他任务
            var scope = new StructuredTaskScope.ShutdownOnSuccess<InferenceResult>()
        ) {
            // ③ 为每个模型 fork 一个虚拟线程(不用手动管线程池)
            for (AiInferenceClient client : clients) {
                scope.fork(() -> {
                    // 关键:手动将 traceId 传入子线程的 MDC
                    MDC.put("traceId", traceId);
                    try {
                        return client.callWithTimeout(req, Duration.ofSeconds(8));
                    } finally {
                        MDC.remove("traceId"); // 防止线程复用导致 MDC 污染
                    }
                });
            }

            // ④ 等待:直到有一个成功 or 全部完成(最长 10 秒)
            scope.joinUntil(Instant.now().plusSeconds(10));

            InferenceResult result = scope.result(); // 取最快成功结果
            recordMetrics(sample, result, "success");
            return result;

        } catch (ExecutionException e) {
            // 全部模型都失败,走降级逻辑
            log.error("[{}] All models failed", traceId, e.getCause());
            recordMetrics(sample, null, "all_failed");
            return fallbackResult(req, "all_models_failed");
        } finally {
            globalSemaphore.release();
        }
    }
}

🔥 生产踩坑 #1:虚拟线程不继承父线程的 MDC(SLF4J ThreadLocal)。直接 fork 子任务,所有日志的 traceId 会消失,分布式链路追踪全部断掉。必须在 fork 时手动传入 traceId,在 finally 块清理,防止线程复用后污染下一个请求。


04 背压控制 + 熔断降级:防止上游雪崩

虚拟线程虽然轻量,但下游 AI API 有并发限制(大多数模型 API QPS 上限在 100~500)。没有背压机制,流量洪峰会直接打爆下游,触发 429 或超时雪崩。

正确姿势是:Semaphore 信号量控制最大并发,超限请求等待而不是直接报错;配合 Resilience4j 熔断器,当错误率过高时自动开启熔断。

@Component
public class CircuitBreakerInferenceClient implements AiInferenceClient {

    private final CircuitBreaker  circuitBreaker;
    private final RateLimiter    rateLimiter;    // 令牌桶限速
    private final OpenAiClient   delegate;

    public CircuitBreakerInferenceClient(CircuitBreakerRegistry cbRegistry,
                                              RateLimiterRegistry  rlRegistry,
                                              OpenAiClient         delegate) {
        this.delegate = delegate;
        // 熔断器配置:60 秒窗口,错误率 > 50% 则打开,15 秒后半开探测
        this.circuitBreaker = cbRegistry.circuitBreaker("openai",
            CircuitBreakerConfig.custom()
                .slidingWindowType(TIME_BASED)
                .slidingWindowSize(60)
                .failureRateThreshold(50.0f)
                .waitDurationInOpenState(Duration.ofSeconds(15))
                .permittedNumberOfCallsInHalfOpenState(3)
                // 429 / 503 视为失败;4xx 业务错误不计入熔断
                .recordException(e -> e instanceof RateLimitException
                                    || e instanceof ServiceUnavailableException)
                .build());

        // 限速器:每秒最多 80 次调用(对齐 OpenAI tier-2 限制)
        this.rateLimiter = rlRegistry.rateLimiter("openai",
            RateLimiterConfig.custom()
                .limitForPeriod(80)
                .limitRefreshPeriod(Duration.ofSeconds(1))
                .timeoutDuration(Duration.ofMillis(500))
                .build());
    }

    @Override
    public InferenceResult callWithTimeout(InferenceRequest req, Duration timeout) {
        return RateLimiter.decorateCheckedSupplier(rateLimiter,
            CircuitBreaker.decorateCheckedSupplier(circuitBreaker,
                () -> delegate.call(req, timeout)
            )
        ).get();
    }

    /**
     * 降级兜底:返回本地小模型的简短回答,或缓存结果。
     * 绝对不能让熔断异常透传到用户界面。
     */
    @Recover
    public InferenceResult fallback(CallNotPermittedException ex,
                                       InferenceRequest req) {
        log.warn("Circuit OPEN for openai, using local fallback");
        return localFallbackModel.quickAnswer(req);
    }
}

💡 关键配置recordException 一定要精细配置——只有真正的服务故障(429 限流、503 不可用)才计入熔断统计。用户输入导致的 400 参数错误不应该触发熔断,否则会误伤正常请求。


05 上下文传播:TraceId 跨虚拟线程的正确姿势

这是团队升级 JDK 21 后最高频的问题:切到虚拟线程之后,Sleuth / Micrometer Tracing 的 TraceId 链路断了。根因是 MDC 基于 ThreadLocal,虚拟线程 fork 不会自动拷贝父线程的 ThreadLocal 值。

生产推荐方案:统一用 ScopedValue(JDK 21 preview,22 正式)替代 ThreadLocal,或者封装一个 ContextCarrier 工具类在任务提交时显式传递。

/**
 * 支持 MDC 上下文传播的虚拟线程执行器。
 * 替换项目中所有 Executors.newVirtualThreadPerTaskExecutor() 的调用。
 */
public class VirtualThreadContextExecutor implements Executor {

    private static final Executor DELEGATE =
        Executors.newVirtualThreadPerTaskExecutor();

    @Override
    public void execute(Runnable command) {
        // 在提交时快照当前线程的完整上下文
        Map<String, String> mdcCopy         = MDC.getCopyOfContextMap();
        Observation           parentObservation = ObservationThreadLocalAccessor
                                                        .getValue();

        DELEGATE.execute(() -> {
            // 恢复 MDC(含 traceId / spanId / userId 等所有 key)
            if (mdcCopy != null) MDC.setContextMap(mdcCopy);

            // 恢复 Micrometer Tracing 的 Observation(链路追踪不断链)
            Observation obs = null;
            if (parentObservation != null) {
                obs = parentObservation.createChildObservation(
                    "virtual-thread-task");
                obs.start();
                ObservationThreadLocalAccessor.setValue(obs);
            }

            try {
                command.run();
            } finally {
                if (obs != null) obs.stop();
                MDC.clear(); // 务必清理,防止线程复用污染
            }
        });
    }
}

// Spring Bean 注册:让 @Async 和 TaskExecutor 都走这个执行器
@Configuration
public class ExecutorConfig {
    @Bean("taskExecutor")
    public AsyncTaskExecutor taskExecutor() {
        return new TaskExecutorAdapter(new VirtualThreadContextExecutor());
    }
}

🔥 生产踩坑 #2:只清理 MDC 还不够,SecurityContextHolder(Spring Security)、RequestContextHolder(Spring MVC)同样基于 ThreadLocal,切虚拟线程后都需要显式传播。建议统一封装一个 ContextCarrier,一次快照所有需要传播的上下文。


06 Spring Boot 3 一键启用 + 监控配置

开启虚拟线程只需一行配置,但生产环境还需要配好监控指标,否则出问题根本不知道瓶颈在哪。

# ── Spring Boot 虚拟线程开关(3.2+ 正式支持)──
spring:
  threads:
    virtual:
      enabled: true       # 一行开启,Tomcat/Undertow 全部走虚拟线程

# ── AI 推理并发控制 ──
ai:
  max-concurrent: 200  # 全局 Semaphore 上限,按下游 QPS 限制调整
  model:
    timeout-seconds: 8  # 单次推理超时,防止慢请求占用连接
    race-timeout-seconds: 10  # 竞速总超时

# ── Resilience4j 熔断配置 ──
resilience4j:
  circuitbreaker:
    instances:
      openai:
        sliding-window-type: time_based
        sliding-window-size: 60
        failure-rate-threshold: 50
        wait-duration-in-open-state: 15s
        permitted-calls-in-half-open-state: 3
        register-health-indicator: true

# ── Micrometer 指标暴露(Prometheus 拉取)──
management:
  endpoints:
    web:
      exposure:
        include: health, prometheus, metrics
  metrics:
    tags:
      application: ${spring.application.name}
    distribution:
      percentiles-histogram:
        ai.inference.latency: true   # P50/P95/P99 延迟分布
      percentiles:
        ai.inference.latency: 0.5, 0.95, 0.99

Grafana 推荐监控指标(Prometheus PromQL)

# 1. 推理延迟 P99(告警阈值:> 5s)
histogram_quantile(0.99,
  rate(ai_inference_latency_seconds_bucket[5m]))

# 2. 背压触发率(突增说明下游扛不住)
rate(ai_inference_backpressure_total[1m])

# 3. 熔断器状态(0=CLOSED 1=OPEN 2=HALF_OPEN)
resilience4j_circuitbreaker_state{name="openai"}

# 4. 虚拟线程挂起数(JVM 内部指标,需开启 JFR)
jvm_threads_virtual_mounted

07 压测数据:升级前后对比

以下数据来自相同硬件(8C16G)、相同负载(1000 并发 AI 推理请求)的对比测试:

指标 OS 线程池 (200线程) 虚拟线程 变化
吞吐量 (RPS) 312 4,380 ↑ ×14
P99 延迟 18.4s 5.8s ↓ 68%
JVM 堆外内存 ~2.1GB ~360MB ↓ 83%
线程数峰值 200 (硬上限) 100,000+ 无上限
CPU 使用率 62% 58% 持平
GC 停顿 频繁 Full GC ZGC <1ms ↓ 显著

结论:在 AI 推理这种极端 I/O 密集场景下,切换虚拟线程是几乎零成本的改造(改一行配置),收益极其显著。但要在生产用好,还需要把背压、熔断、上下文传播这三个配套机制一起做到位。


08 生产踩坑清单:升级前必读

踩坑场景 根因 解决方案
日志 traceId 消失 MDC ThreadLocal 不继承 封装 VirtualThreadContextExecutor,显式传播
Spring Security 鉴权失效 SecurityContextHolder ThreadLocal 使用 InheritableThreadLocal 模式或显式传递
数据库连接池耗尽 虚拟线程太多,HikariCP 连接不够 HikariCP maximum-pool-size 按实际 DB 并发上调
synchronized 锁死 虚拟线程遇 synchronized 会 pin 住载体线程 换用 ReentrantLock / StampedLock
CPU 密集任务性能下降 调度开销,虚拟线程不减少 CPU 竞争 CPU 密集任务保留 ForkJoinPool 线程池
AI API 429 雪崩 无背压,虚拟线程并发打爆下游 Semaphore + Resilience4j 双重保护

一句话总结

I/O 密集用虚拟线程,CPU 密集仍用平台线程。

虚拟线程不是银弹,但在 AI 推理这个场景,它是真正的游戏规则改变者。

Logo

汇聚全球AI编程工具,助力开发者即刻编程。

更多推荐