JDK 25 虚拟线程实战:60 万条抓取任务的线程模型迁移记录

2026-09-14 0 764

先交代背景。我们有一个商品价格刷新服务,每天凌晨跑一次全量,数据源是大约 60 万个商品详情页 URL。每个 URL 要做三件事:发一个 HTTP GET、从返回的 HTML 里抠出价格和库存、归一化以后写回 MySQL。

单次请求平均 120ms 左右,长尾能到好几秒。任务之间彼此独立,没有依赖关系。是很典型的 IO 密集型批处理。

原来这套东西跑在一个固定 200 线程的 ThreadPoolExecutor 上,队列用的是无界的 LinkedBlockingQueue。它在大多数时候确实能跑完,但每次大促前扩容、每次数据源变慢,都会以一种很难排查的方式出问题。这篇文章记录的是把它换成虚拟线程之后的完整过程,包括中间写错的第一版。

原来的写法,问题不在线程数

先看迁走之前的代码,几乎是教科书式的写法:

private static final ExecutorService POOL = new ThreadPoolExecutor(
        200, 200,
        0L, TimeUnit.MILLISECONDS,
        new LinkedBlockingQueue<>(),        // 没有容量上限
        new ThreadFactoryBuilder().setNameFormat("price-%d").build());

public void refreshAll(List<String> urls) {
    for (String url : urls) {
        POOL.execute(() -> refreshOne(url));
    }
}

第一眼看上去没什么毛病。真跑起来才知道问题出在哪:

提交侧完全没有反馈。60 万个任务会在几秒钟内全部塞进队列。因为队列是无界的,execute 从来不阻塞,也从来不抛拒绝异常。数据源已经开始变慢的时候,提交循环早就跑完了,你看到的只是”进程还活着,但指标不动”。等靠日志发现异常,队列里可能已经积压了几十万条任务。

延迟被队列吃掉了。我们当时统计过一个数:单个任务从提交到真正开始执行,P99 是 140 多秒。也就是说,一个本来 120ms 的请求,用户视角看是”两分半钟”。而这个时间不出现在任何一条 HTTP 超时日志里,因为请求根本还没发出去。排查的时候你会怀疑网络、怀疑对端、怀疑 DNS,唯独不会想到是队列。

并发度和资源上限被绑定在一起了。200 这个数字是拍脑袋定的,但它实际上同时决定了三件事:同时打开的 TCP 连接数、同时占用的 MySQL 连接数、以及堆上平台线程的栈空间开销。想提到 2000?每个平台线程默认 1MB 栈,光栈就是 2G,GC 和上下文切换都会跟着遭殃。

真正的问题不是”200 太小”,而是线程数这个词同时承担了并发控制、内存预算、资源配额三个职责,而你只能调一个数字。

换成虚拟线程,第一版是错的

JDK 21 正式发布的虚拟线程,到 JDK 25 已经是第二个 LTS 版本了。它的定位其实很朴素:把”阻塞”这件事的线程成本降到可以忽略。平台线程阻塞要付出一个栈的代价,虚拟线程阻塞只是把 Continuation 从载体线程上卸下来,堆上那点状态几百字节。

所以最自然的改写是这样:

try (ExecutorService pool = Executors.newVirtualThreadPerTaskExecutor()) {
    for (String url : urls) {
        pool.submit(() -> refreshOne(url));
    }
}

这段代码能跑,而且在压测里看起来很快。但它是错的,错在它把原来”无界队列”的问题换了个形式保留了下来——只不过这次堆积的不是任务对象,而是同时打开的连接。

newVirtualThreadPerTaskExecutor 不是线程池,它没有核心线程数、没有最大线程数、没有队列。每个任务一条虚拟线程,立刻开始运行。60 万个任务就是 60 万条虚拟线程同时奔向网络栈。

虚拟线程本身很便宜,但你的下游不便宜。DNS 解析器、对端服务的限流策略、容器里的文件描述符上限、公司的出口网关,没有一个欢迎 60 万并发。压测跑到第三分钟的时候我们的出口带宽被占满,整栋楼的同事开始抱怨网络卡。这不是虚拟线程的问题,是我们把”没有上限”误当成了”不需要上限”。

加上背压之后的完整实现

正确的思路是:虚拟线程负责降低阻塞成本,背压必须自己管。Java 里做这件事最直接的工具就是 Semaphore。

这里有一个非常容易写错的地方。很多人会把 acquire 放在任务体内部:

pool.submit(() -> {
    inFlight.acquire();          // 错:任务已经提交了,压力还在队列里
    try { refreshOne(url); }
    finally { inFlight.release(); }
});

这不叫背压。任务已经被提交给执行器,虚拟线程已经创建,只是卡在 acquire 上排队。内存里照样躺着 60 万个待执行的任务,和原来无界队列的毛病一模一样。

有效的做法是在提交之前 acquire,让提交循环自己停住:

for (String url : urls) {
    inFlight.acquire();          // 拿不到许可,这个循环就走不下去
    pool.submit(() -> {
        try { refreshOne(url); }
        finally { inFlight.release(); }
    });
}

这样在途任务数量被硬性卡在 MAX_IN_FLIGHT 附近,内存占用可预测,而且背压会一路传回给生产端——如果数据源在变慢,扫描 URL 的循环会自己慢下来。

完整实现如下。多了一个分批,原因是如果对 60 万个 URL 一次性循环,提交循环虽然被限住了,但整个 try 块要等到所有任务结束才退出,中间没法观测进度,也没法做分批提交。每批 4000 条,批内用同一个执行器,批结束关闭时自然形成一个同步点。

import java.io.IOException;
import java.net.URI;
import java.net.http.HttpClient;
import java.net.http.HttpRequest;
import java.net.http.HttpResponse;
import java.time.Duration;
import java.util.List;
import java.util.concurrent.ExecutorService;
import java.util.concurrent.Executors;
import java.util.concurrent.Semaphore;
import java.util.concurrent.atomic.AtomicInteger;

public class PriceRefresher {

    private static final int MAX_IN_FLIGHT = 800;
    private static final int BATCH_SIZE = 4_000;

    private final HttpClient http = HttpClient.newBuilder()
            .connectTimeout(Duration.ofSeconds(3))
            .followRedirects(HttpClient.Redirect.NORMAL)
            .build();

    private final Semaphore inFlight = new Semaphore(MAX_IN_FLIGHT);
    private final AtomicInteger okCount = new AtomicInteger();
    private final AtomicInteger failCount = new AtomicInteger();

    public void refreshAll(List<String> urls) throws InterruptedException {
        int total = urls.size();
        for (int from = 0; from < total; from += BATCH_SIZE) {
            int to = Math.min(from + BATCH_SIZE, total);
            refreshBatch(urls.subList(from, to));
            System.out.printf("进度 %d/%d  成功 %d  失败 %d%n",
                    to, total, okCount.get(), failCount.get());
        }
    }

    private void refreshBatch(List<String> batch) throws InterruptedException {
        try (ExecutorService pool = Executors.newVirtualThreadPerTaskExecutor()) {
            for (String url : batch) {
                inFlight.acquire();          // 背压发生在这里,不是任务体里
                try {
                    pool.submit(() -> {
                        try {
                            refreshOne(url);
                            okCount.incrementAndGet();
                        } catch (Exception e) {
                            failCount.incrementAndGet();
                        } finally {
                            inFlight.release();
                        }
                    });
                } catch (RuntimeException e) {
                    inFlight.release();      // 提交失败必须还回许可,否则信号量会慢慢泄漏
                    throw e;
                }
            }
        }   // close() 会阻塞到这批任务全部结束
    }

    private void refreshOne(String url) throws Exception {
        HttpRequest req = HttpRequest.newBuilder(URI.create(url))
                .timeout(Duration.ofSeconds(5))
                .header("User-Agent", "PriceBot/1.0")
                .GET()
                .build();

        HttpResponse<String> resp = http.send(req, HttpResponse.BodyHandlers.ofString());
        if (resp.statusCode() != 200) {
            throw new IOException("HTTP " + resp.statusCode() + " from " + url);
        }

        priceRepository.upsert(HtmlPriceParser.parse(url, resp.body()));
    }
}

有两个细节值得单独说。

第一,任务体内部把异常全吃了。这是故意的。批量任务里,一个 URL 抓失败不应该影响其他任何一个,异常往上冒只会污染整批的语义。失败计数才是我们关心的东西,具体某一条的错误日志走异步队列。如果你确实需要拿到每个任务的结果,那就保留 Future 列表,但要注意 4000 个 Future 对象本身也是内存。

第二,ExecutorService.close() 不能被中断。它的契约决定了它会一直等到所有任务结束。所以每个任务自己必须带超时——HTTP 请求有 timeout、连接有 connectTimeout、数据库有查询超时。把这些兜底做全,close 才是安全的。少一个,进程就可能永远挂在那里。

迁完之后踩到的几个坑

ThreadLocal 的副本数量变了

我们原来在过滤器里往 ThreadLocal 里塞了一个 traceId,方便日志串联,用了几年一直没问题。切到虚拟线程之后,在途虚拟线程的数量从 200 涨到了几千,而且每批结束就换一批新的。

ThreadLocal 的开销本身不算大,真正难受的是两点:一是如果用了 InheritableThreadLocal,父线程在创建子线程时会复制整张 map,虚拟线程是每个任务新建一条,等于每个任务复制一次;二是那些被 ThreadLocal 隐式持有的对象(比如连接、ByteBuffer)会一直挂到虚拟线程结束才回收,堆上会出现大量短命的中等对象。

我们的做法是把 traceId 改成方法参数显式往下传。听起来很土,但改动量比想象中小,而且比任何隐式上下文都好调试。JDK 25 里的 ScopedValue 是更优雅的答案,但它目前还是预览特性,需要 –enable-preview,生产环境我们暂时没有开。

数据库连接池成了新的瓶颈

这个坑最典型。MAX_IN_FLIGHT 设成 800,HikariCP 的 maximumPoolSize 还是默认的 10。结果就是 800 条虚拟线程里有 790 条阻塞在获取连接上。

阻塞本身不贵,问题是默认的 connectionTimeout 是 30 秒。790 条虚拟线程各等 30 秒才拿到异常,这个等待时间在监控面板上完全看不出来——你只会看到”任务失败率突然变高”。

两处调整:连接获取超时从 30 秒压到 1.5 秒,快速失败;同时把 MAX_IN_FLIGHT 从 800 降到 240,让它更接近下游真实能承受的量。

这件事让我意识到一个判断标准:信号量的许可数不应该按”我想跑多快”来定,而应该按”链路里最弱的一环能承受多少”来定。如果 HTTP 抓取和数据库写入的容量差了十倍,更好的做法是把它们拆成两段流水线,各配各的信号量,而不是让一条流水线用同一个数字去约束两种资源。

CPU 密集的那段代码没有变快

迁完之后总耗时确实降了很多,但单任务的 CPU 时间几乎没变。原因是我们解析 HTML 用的是一堆正则,其中一条在异常输入下会回溯,平均解析耗时 8ms 左右。

虚拟线程解决的是阻塞成本,不解决计算成本。后来我们把那段正则改写成了一个简单的状态机,解析耗时从 8ms 降到 1.2ms。这一项的收益,说实话比并发模型本身的改动还大。

所以如果任务是纯计算的,虚拟线程帮不上忙,用 ForkJoinPool 或者并行流更合适。虚拟线程的适用边界很清楚:任务大部分时间在等,而且等待的东西没有比线程更紧的容量限制。

监控面板要重新做

以前看线程池,几个指标就够了:getActiveCount、getQueue().size()、getCompletedTaskCount。换成虚拟线程之后这些全都没了。

现在看的是这么几个:Semaphore 的 availablePermits(还有多少并发余量)、成功失败计数器、以及 JFR 里的 jdk.VirtualThreadStart 和 jdk.VirtualThreadEnd 事件。另外提醒一句,Thread.activeCount() 不统计虚拟线程,用它做健康检查会得到一个很有欺骗性的数字。

实测数据和几句不中听的话

环境是 4 核 8G 的容器,60 万个 URL,两版跑在同一批数据上:

指标 200 平台线程 + 无界队列 虚拟线程 + Semaphore(240) + 分批
墙钟耗时 3 小时 47 分 58 分
单任务 P99 耗时 142 秒 6.2 秒
堆峰值 4.8G,频繁 Full GC 860M

数字看着很漂亮,但我得说清楚两件事,不然容易误导人。

第一,这个提升里只有一部分是虚拟线程的功劳。P99 从 142 秒降到 6.2 秒,绝大部分来自”把提交循环卡住”这件事——原来的 142 秒里,绝大多数时间是在队列里干等,而不是真正在跑。哪怕你继续用平台线程,只要把队列换成有界队列加上 CallerRunsPolicy,P99 也会大幅下降。虚拟线程的贡献是把并发度从 200 提到 240 而不用付出额外的内存和切换成本,这是真实的,但没到十倍。

第二,堆峰值从 4.8G 降到 860M,主要是背压的功劳,不是虚拟线程更省内存。在途任务少了,自然堆上就干净了。

如果你现在正在做类似的迁移,我的建议顺序是:先把背压做对,用有界队列或者 Semaphore,观察一段时间;确认瓶颈真的在”线程不够”而不是”下游不够”之后,再换虚拟线程。这两步的收益是能分开衡量的,混在一起改,出了问题你根本不知道该回滚哪一步。

最后一点:synchronized 导致的线程固定(pinning)问题在 JDK 24 的 JEP 491 里已经基本解决,跑在 JDK 25 上不用再为此改代码。但如果你的服务还在 JDK 21 上,老代码里那些包着 IO 的 synchronized 块最好换成 ReentrantLock,这个坑在 21 上是真实存在的。

JDK 25 虚拟线程实战:60 万条抓取任务的线程模型迁移记录
收藏 (0) 打赏

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

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

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

淘吗网 java JDK 25 虚拟线程实战:60 万条抓取任务的线程模型迁移记录 https://www.taomawang.com/server/java/2757.html

常见问题

相关文章

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

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