虚拟线程实战:6个下游聚合接口从1.2秒降到200毫秒的完整改造

2026-09-17 0 516

我们有个接口叫 /api/offer/aggregate,商品详情页比价用的。一次请求要问6个下游:主站库存、第三方渠道价、优惠券、评价摘要、物流时效、推荐兜底。每个下游平均响应120~250毫秒,串行调完P99是1.2秒。这个接口每天1.4亿次调用,是我们流量前五的接口。

这篇文章讲的是我怎么把它改成虚拟线程方案,以及踩到的那些坑。不是”虚拟线程真香”的软文,里面有至少两个坑我改了一整天才定位到。

改造前的样子

最早的实现就是最朴素的串行:

public List<Offer> aggregate(String sku) {
    List<Offer> result = new ArrayList<>(6);
    for (OfferClient client : clients) {
        result.add(client.query(sku));
    }
    return result;
}

问题很直白。6个下游全部串行,总耗时是6个响应时间之和。中间任何一个慢,整条链路跟着慢。

第一次优化用的是固定线程池,200个平台线程,每个请求把6个下游提交进去然后等结果:

private static final ExecutorService POOL = Executors.newFixedThreadPool(200);

public List<Offer> aggregate(String sku) {
    List<Future<Offer>> futures = clients.stream()
        .map(c -> POOL.submit(() -> c.query(sku)))
        .toList();
    // ...依次 get
}

效果有,但没有想象中那么好。Tomcat 本身有 200 个工作线程,业务又开 200 个,机器上 400 个平台线程在跑。每个线程栈默认 1MB,光是线程栈就占掉几百兆。更麻烦的是这些线程 90% 的时间在等网络 IO,什么都没干,但栈内存和上下文切换的成本一直在付。

压测把并发拉到 2000 的时候,线程池的队列开始堆积,请求排队时间比实际下游耗时还长。加线程数到 500?上下文切换的损耗已经吃掉收益了。

为什么换虚拟线程

这里的瓶颈不是 CPU,是等待。6个下游里最快的一个平均 110 毫秒,最慢的 380 毫秒,CPU 在这段时间里基本什么都没做。平台线程的问题是”一个线程只能承载一个等待”,而等待这件事本身不需要那么贵的资源。

虚拟线程把”线程”和”操作系统线程”解耦了。JVM 调度器把成千上万个虚拟线程映射到少量的载体线程(carrier thread,默认数量等于 CPU 核数)上。当虚拟线程执行阻塞 IO 时,JVM 会把它从载体线程上摘下来,让载体线程去跑别的虚拟线程。栈被挪到堆里,按需增长,不是一开始就分配 1MB。

换成虚拟线程之后,代码长得和串行版本几乎一样,因为不需要池化,也不需要考虑队列:

public List<Offer> aggregate(String sku) {
    try (var executor = Executors.newVirtualThreadPerTaskExecutor()) {
        List<Future<Offer>> futures = new ArrayList<>(clients.size());
        for (OfferClient client : clients) {
            futures.add(executor.submit(() -> client.query(sku)));
        }
        List<Offer> result = new ArrayList<>(clients.size());
        for (int i = 0; i < futures.size(); i++) {
            result.add(futures.get(i).get());
        }
        return result;
    }
}

newVirtualThreadPerTaskExecutor() 是 JDK 21 转正的 API,每个任务创建一条新虚拟线程。别看到”每个任务一条线程”就觉得贵,创建一条虚拟线程的成本大概是几百字节的堆分配,比创建平台线程便宜两三个数量级。

生产版本:超时、降级、取消都要有

上面那段代码能跑,但不能上生产。真实场景里必须处理三件事:单个下游超时怎么办、一个下游挂了要不要拖累整个接口、请求被取消时怎么收尾。

先定义返回值,带上降级标记:

public record Offer(String source, long priceCent, int stock, boolean degraded) {

    public static Offer fallback(String source) {
        return new Offer(source, -1L, 0, true);
    }
}

public interface OfferClient {
    String source();
    Offer query(String sku) throws Exception;
}

聚合器的完整实现:

public final class OfferAggregator {

    private static final long TIMEOUT_NANOS = Duration.ofMillis(800).toNanos();

    private final List<OfferClient> clients;

    public OfferAggregator(List<OfferClient> clients) {
        this.clients = List.copyOf(clients);
    }

    public List<Offer> aggregate(String sku) {
        try (var executor = Executors.newVirtualThreadPerTaskExecutor()) {
            List<Future<Offer>> futures = new ArrayList<>(clients.size());
            for (OfferClient client : clients) {
                futures.add(executor.submit(() -> client.query(sku)));
            }

            List<Offer> result = new ArrayList<>(clients.size());
            long deadline = System.nanoTime() + TIMEOUT_NANOS;

            for (int i = 0; i < futures.size(); i++) {
                OfferClient client = clients.get(i);
                Future<Offer> future = futures.get(i);
                long remain = deadline - System.nanoTime();
                try {
                    result.add(future.get(Math.max(remain, 0L), TimeUnit.NANOSECONDS));
                } catch (TimeoutException e) {
                    future.cancel(true);
                    result.add(Offer.fallback(client.source()));
                } catch (ExecutionException e) {
                    result.add(Offer.fallback(client.source()));
                } catch (InterruptedException e) {
                    for (Future<Offer> f : futures) {
                        f.cancel(true);
                    }
                    Thread.currentThread().interrupt();
                    return List.of();
                }
            }
            return result;
        }
    }
}

几个设计点值得说清楚:

  • 所有下游共享一个 800 毫秒的绝对 deadline,而不是每个下游各给 800 毫秒。否则最坏情况下总耗时是 800 乘 6。第一个下游如果吃掉 700 毫秒,后面的只剩 100 毫秒,这符合”接口整体响应时间”的预期。
  • 单个下游失败或超时降级为空 Offer,不影响其他 5 个。比价这种场景,少一路数据用户感知不强,整个接口 500 才是事故。
  • 捕获 InterruptedException 时要把所有还没完成的任务取消掉,然后把中断标志位恢复。Tomcat 在客户端断开连接时会中断工作线程,这时候如果不把子任务收干净,虚拟线程会继续往下游发请求。
  • try-with-resources 关闭 executor 时,close() 内部会执行 shutdown 并等待所有任务结束。这一点很关键,见后面的坑。

压测数据

环境:JDK 21.0.3,16 核 32G,下游用 Mock 服务模拟(固定 150 毫秒延迟),压测工具 wrk2,持续 5 分钟。Tomcat 最大线程数统一 200。

方案 额外线程/连接 P50 P99 QPS 错误率
串行调用 1180ms 1350ms 168 0
固定线程池 200 200 平台线程 780ms 1140ms 253 0.31%
CompletableFuture + 公共 ForkJoinPool FJP 公共池 265ms 490ms 618 0.09%
虚拟线程 0(复用载体线程) 212ms 268ms 1140 0

P99 从 1350 毫秒降到 268 毫秒。注意到虚拟线程的 P99 和 P50 差距很小,因为排队的成分基本消失了,每个请求都是 6 条虚拟线程立刻开始跑。

固定线程池那一行的错误率来自队列溢出后的 RejectedExecutionException,这个是压测时故意打出来的,不是意外。CompletableFuture 的表现已经不错,但用公共 ForkJoinPool 有个隐患:它的并行度默认是 CPU 核数减一,如果有别的业务也在用公共池,就会互相干扰,而且这个池的大小你没法针对下游特性调优。

五个真实的坑

坑一:synchronized 会把虚拟线程钉在载体线程上

这是 JDK 21 到 JDK 23 上最要命的问题。虚拟线程在 synchronized 块里执行阻塞操作时,JVM 没法把它从载体线程上摘下来,只能连载体线程一起阻塞住。这个现象叫 pinning(钉住)。

我们代码里有一处日志切面用了 synchronized 锁住一个共享的缓冲区,里面调用了下游 HTTP 接口:

synchronized (bufferLock) {
    // 错误示范:这里会钉住载体线程
    String body = httpClient.send(request, BodyHandlers.ofString()).body();
    buffer.append(body);
}

16 核的机器,载体线程池只有 16 条。如果几十个虚拟线程同时卡在这个同步块里,16 条载体线程全被钉死,整个调度器就瘫痪了,表现为接口整体卡住不动,CPU 使用率却很低。

排查办法是加上 JVM 参数 -Djdk.tracePinnedThreads=full,它会在检测到钉住时打出完整堆栈。我们在 JDK 21.0.3 上用这个参数两分钟就定位到了。

修复很简单,换成 ReentrantLock

private final ReentrantLock bufferLock = new ReentrantLock();

bufferLock.lock();
try {
    String body = httpClient.send(request, BodyHandlers.ofString()).body();
    buffer.append(body);
} finally {
    bufferLock.unlock();
}

ReentrantLock 在等待时是 park,虚拟线程会被正常摘下来。顺便说一句,JDK 24 的 JEP 491 已经修掉了 synchronized 的钉住问题,但还是建议统一用 ReentrantLock,因为你在 JDK 21 上维护的代码不一定马上能升上去。

坑二:ThreadLocal 的成本被放大了

以前用固定线程池的时候,200 个线程复用到死,ThreadLocal 里的对象最多也就 200 份。换成虚拟线程之后,每个请求 6 条虚拟线程,QPS 1000 就是每秒 6000 条线程被创建。如果 ThreadLocal 在初始化时塞了一个 64KB 的字符串缓冲区,那每秒要分配 384MB。

我们有个 MDC 的 traceId 透传就是通过 ThreadLocal 做的。压测跑了十分钟之后老年代增长速度明显变快,GC 日志里全是这些短命的大对象。

处理方式有两个方向。短期是把 ThreadLocal 的初始值改小,或者改成懒加载。长期是换 ScopedValue(JDK 21 预览,JDK 25 转正路径上),它是不可变的、按作用域绑定,虚拟线程继承语义清晰,也没有”每条线程一份拷贝”的问题。如果用 Spring,MDC 那套适配需要自己封装一层,官方的 TaskDecorator 方案在虚拟线程下依然有效,但要注意复制的内容别太大。

坑三:真正的瓶颈在连接池,不在线程数

虚拟线程把线程问题解决了,然后你会发现请求全堵在数据库连接池上。

我们有个下游是查本地缓存表,走 HikariCP,maximumPoolSize 是 10。改造前 Tomcat 200 线程、业务池 200 线程,其实已经有不少请求在等了,但因为线程池本身也在排队,问题被掩盖了。换成虚拟线程之后,并发瞬间拉满,10 个连接根本不够用,P99 反而涨了。

这里必须建立正确的心智模型:虚拟线程解决的是”等待不占资源”,不解决”下游容量”。连接池大小、下游 QPS 上限、限流阈值,这些是下游的容量约束,本质上和线程模型无关。虚拟线程只是把等待的成本降下来了,让你能更清楚地看到真正的瓶颈在哪。

我们的处理是把连接池调到 32 并配合下游的 QPS 限制,同时对访问数据库的那条支路单独加了信号量:

private final Semaphore dbPermits = new Semaphore(24);

Offer queryDb(String sku) throws Exception {
    dbPermits.acquire();
    try {
        return doQuery(sku);
    } finally {
        dbPermits.release();
    }
}

注意 Semaphore.acquire() 在虚拟线程上是安全的,它内部是 park,不会钉住载体线程。如果你想限制下游并发,用信号量,不要想着”给虚拟线程池设个大小”——虚拟线程池本来就不该有固定大小,设了等于退化成平台线程池。

坑四:try-with-resources 的 close 会一直等

这是我自己踩的。前面那段代码里 try (var executor = ...) 关闭时的语义是:shutdown,然后等待所有已提交任务结束。如果某个下游的 HTTP 客户端没有设 readTimeout,连接挂在那里十分钟,那么 close() 就会等十分钟,即使上层已经超时降级返回了。

换句话说,聚合逻辑里的 800 毫秒 deadline 保证了”返回给用户的数据在 800 毫秒内组装完”,但不保证”所有虚拟线程在 800 毫秒内结束”,更不保证”方法在 800 毫秒内返回”。

两个补救措施。第一,所有 HTTP 客户端必须设 connectTimeout 和 readTimeout,超时值要略大于聚合层 deadline,让取消能正常传导。第二,future.cancel(true) 只是发送中断信号,如果下游客户端不响应中断(比如用了一些老版本的阻塞 IO 库),线程照样不退出。要验证这一点,最简单的方法是把下游 Mock 改成无限睡眠,然后观察接口的方法返回时间。

如果你的业务对方法返回时间极其敏感,可以考虑用 StructuredTaskScope,它在 JDK 21 是预览 API,需要 --enable-preview,写法更干净:

try (var scope = new StructuredTaskScope.ShutdownOnFailure()) {
    List<StructuredTaskScope.Subtask<Offer>> subtasks = clients.stream()
        .map(c -> scope.fork(() -> c.query(sku)))
        .toList();

    scope.joinUntil(Instant.now().plusMillis(800));

    List<Offer> result = new ArrayList<>(subtasks.size());
    for (int i = 0; i < subtasks.size(); i++) {
        var st = subtasks.get(i);
        result.add(st.state() == StructuredTaskScope.Subtask.State.SUCCESS
            ? st.get()
            : Offer.fallback(clients.get(i).source()));
    }
    return result;
}

它的好处是超时后 joinUntil 返回,”结构化”地保证作用域退出时所有子任务都已经结束或取消。但因为它还是预览特性,我们线上暂时用的还是前面那套普通 executors 方案,等转正再换。

坑五:CPU 密集型任务不要用虚拟线程

我们还有一个下游是本地规则引擎,纯计算,单次 150 毫秒 CPU 时间。这条支路换成虚拟线程之后,QPS 掉了 15%。

原因不复杂。虚拟线程只在阻塞时让出载体线程,纯计算场景下它一直占着载体线程不放,调度器多了一层间接性却没有收益。而且载体线程数默认等于 CPU 核数,并发计算的吞吐上限就是核数,虚拟线程改变不了这个事实。

这条支路我们改回了固定大小的平台线程池,大小设成核数,效果和新版本持平但更省内存。

什么时候不该用虚拟线程

  • CPU 密集型任务。并行度受限于核数,虚拟线程只是加了一层调度开销。
  • 需要精确限制并发的场景。虚拟线程池没有固定大小,想限并发得显式用 Semaphore 或者连接池,别指望线程数能约束住。
  • 依赖 ThreadLocal 传递上下文的框架。先评估一下清空和重建的成本,有些框架的 ThreadLocal 初始化很重。
  • 调用了大量 native 库或者 JNI 的代码。native 阻塞调用会钉住载体线程,JVM 管不了。
  • 持锁粒度过粗的代码。锁竞争在虚拟线程下会变得更明显,因为并发度上去了,临界区的争抢更激烈。

迁移清单

如果你打算在现有项目里上虚拟线程,我建议按这个顺序走:

  1. 先把 JDK 升到 21 或以上。如果升到 24,synchronized 钉住问题自动消失,但代码里还是建议统一用 ReentrantLock。
  2. 如果用 Spring Boot 3.2+,可以先把 spring.threads.virtual.enabled 打开,让 Tomcat 的请求处理线程变成虚拟线程,观察一圈。这一步改动最小,风险也最低。
  3. 全量搜索代码里所有的 synchronized 关键字,逐个检查里面有没有阻塞调用,有的话先改掉。
  4. 清点 ThreadLocal 的使用,特别是存放集合、缓冲区、大字符串的那种。
  5. 检查所有连接池的配置,包括数据库、Redis、HTTP 客户端。虚拟线程下这些配置才是真正的并发上限。
  6. 给所有外部调用加上超时,connect、read、write 都要设。这一步在任何线程模型下都是必须的,只是虚拟线程下更明显。
  7. 压测的时候重点看 P99 和 P999,别只看平均值。虚拟线程的优势主要体现在尾延迟上,平均值可能看不出差别。

最后说一句感受。虚拟线程最大的价值不是让代码变快,而是让并发代码变简单。原来那套”提交任务、组合 Future、处理超时、管理线程池生命周期”的模板代码,现在可以退化成很接近串行的写法,而这种简单是能被 review 的简单。

但它不是银弹。把虚拟线程接上去之前,先想清楚你的瓶颈到底在哪。如果瓶颈是数据库的 10 个连接,那换什么线程模型都没用。

虚拟线程实战:6个下游聚合接口从1.2秒降到200毫秒的完整改造
收藏 (0) 打赏

感谢您的支持,我会继续努力的!

打开微信/支付宝扫一扫,即可进行扫码打赏哦,分享从这里开始,精彩与您同在
点赞 (0)

版权声明:
本站资源有的来自互联网收集整理,本站纯免费分享提供学习使用,如果侵犯了您的合法权益,请发送邮件1506151422@qq.com联系,将会及时下架删除。
本站资源仅供研究、学习交流之用,免费开源项目不代表完全可商用,若商业用途请先咨询开发企业能否商用,否则产生的一切后果将由下载用户自行承担。
原创板块未经允许不得转载,否则将追究法律责任。

淘吗网 java 虚拟线程实战:6个下游聚合接口从1.2秒降到200毫秒的完整改造 https://www.taomawang.com/server/java/2768.html

常见问题

相关文章

猜你喜欢
发表评论
暂无评论
官方客服团队

为您解决烦忧 - 24小时在线 专业服务