先交代背景。我们有一个商品价格刷新服务,每天凌晨跑一次全量,数据源是大约 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 上是真实存在的。

