Administrator
发布于 2025-01-21 / 3346 阅读
87

虚拟线程在 AI 高并发调用中的实践

5 万条评测数据,从 80 分钟压到 9 分钟

今年 1 月要做一次大模型效果评测:5 万条标注问题,分别跑三个候选模型,记录每条的输出、耗时和 token 数。总共 15 万次调用。

第一版脚本用 200 线程的固定线程池,跑了 82 分钟还没跑完(因为触发限流重跑了一批)。我花了一天改成虚拟线程版本,最后 9 分 20 秒跑完 5 万条。这篇记录中间踩的坑和实测数据。

先看原来的瓶颈在哪

ExecutorService pool = Executors.newFixedThreadPool(200);
List<CompletableFuture<EvalResult>> fs = questions.stream()
        .map(q -> CompletableFuture.supplyAsync(() -> callModel(q), pool))
        .toList();

跑起来之后看监控:

# 压测客户端的数据
active_threads            200 / 200        # 线程池打满
cpu_usage                 0.14             # 但 CPU 只有 14%
llm_call_inflight         200              # 在飞的请求数被线程池卡死
llm_response_p95          1.82 s
throughput                约 620 条/分钟

CPU 14%,线程 100% 占用。这是典型的 IO 密集型负载:所有线程都在等 HTTP 响应,什么也没干,但线程池没有更多线程可用了。

这时候加线程是对的,但加到 2000 个平台线程就要 2 GB 栈内存(默认 1MB/线程),而且上下文切换成本会吃掉收益。这正是虚拟线程的设计场景。

第一次改:直接换虚拟线程,效果一般

try (var executor = Executors.newVirtualThreadPerTaskExecutor()) {
    List<Future<EvalResult>> fs = questions.stream()
            .map(q -> executor.submit(() -> callModel(q)))
            .toList();
    for (Future<EvalResult> f : fs) results.add(f.get());
}

结果让我意外:

llm_call_inflight         201          # ???
throughput                约 630 条/分钟

并发度几乎没变。查了一圈才反应过来:并发上限被连接池卡住了,跟线程数无关

我们用的是 JDK 自带的 java.net.http.HttpClient,它内部对同一主机的并发请求数没有硬上限,但用的是一个共享的 selector 线程池。更关键的是我在外层套了个自己写的限流器,里面有个 synchronized

public class RateLimiter {
    private final AtomicInteger inflight = new AtomicInteger();
    private final int maxConcurrent;

    public synchronized void acquire() throws InterruptedException {
        while (inflight.get() >= maxConcurrent) {
            wait(50);        // ← synchronized 里 wait,虚拟线程被 pin
        }
        inflight.incrementAndGet();
    }
}

所有虚拟线程都挤在这个 synchronized 块上,被钉在 8 个载体线程上(8 核机器),实际并发度就是 8 到 200 之间抖动。

用 JFR 确认:

$ jfr summary vt.jfr | grep -i pinned
 jdk.VirtualThreadPinned    3,412,908 events
$ jfr print --events jdk.VirtualThreadPinned vt.jfr | grep -A3 "Stack Trace" | head -20
   at app//RateLimiter.acquire(RateLimiter.java:42)
   at java.base/java.lang.Object.wait(Object.java:366)

第二次改:换掉阻塞原语,用 Semaphore

public class ConcurrencyGate {
    private final Semaphore permits;

    public ConcurrencyGate(int max) { this.permits = new Semaphore(max); }

    public void acquire() throws InterruptedException {
        permits.acquire();      // 虚拟线程里是协作式挂起,不占载体线程
    }

    public void release() { permits.release(); }
}
int concurrency = 2000;
var gate = new ConcurrencyGate(concurrency);

try (var executor = Executors.newVirtualThreadPerTaskExecutor()) {
    List<Future<EvalResult>> fs = questions.stream()
        .map(q -> executor.submit(() -> {
            gate.acquire();
            try {
                return callModel(q);
            } finally {
                gate.release();
            }
        }))
        .toList();
    ...
}

JFR 再抓一次,pin 事件变成 0,并发度稳定在 2000。

HTTP 客户端的选型

顺带对比了三种客户端在虚拟线程下的表现(2000 并发,目标服务 P95 响应 1.8 秒):

客户端配置pin 事件吞吐说明
Apache HttpClient 5maxConnPerRoute=200018 万/min4100 条/分钟连接池内部有 synchronized
OkHttp 4maxRequests=200005200 条/分钟Dispatcher 用 synchronized 但很快
JDK HttpClient默认05400 条/分钟异步内核,最适合虚拟线程

最终选了 JDK 自带的 HttpClient。send() 这个同步方法在虚拟线程里不会 pin,因为它的阻塞发生在 JDK 内部的 NIO 层,虚拟线程能正常挂起。

private static final HttpClient HTTP = HttpClient.newBuilder()
        .connectTimeout(Duration.ofSeconds(5))
        .version(HttpClient.Version.HTTP_2)      // 供应商支持 HTTP/2 多路复用
        .build();

EvalResult callModel(String q) throws Exception {
    HttpRequest req = HttpRequest.newBuilder(URI.create(ENDPOINT))
            .timeout(Duration.ofSeconds(60))
            .header("Content-Type", "application/json")
            .header("Authorization", "Bearer " + API_KEY)
            .POST(HttpRequest.BodyPublishers.ofString(buildBody(q)))
            .build();

    HttpResponse<String> resp = HTTP.send(req, HttpResponse.BodyHandlers.ofString());
    if (resp.statusCode() == 429) throw new RateLimitedException();
    return parse(resp.body());
}

HTTP/2 在这里帮了大忙。2000 个并发请求复用 4 条 TCP 连接,省掉了 2000 次 TLS 握手和连接建立的开销。从 HTTP/1.1 切到 HTTP/2,吞吐从 3900 涨到 5400 条/分钟。

第三次改:并发度不是越大越好

我试了从 500 到 8000 的六档并发,看吞吐和错误率的变化:

并发度吞吐(条/分钟)429 错误率P95 单次耗时客户端内存
200(原方案)6200.02%1.82 s0.6 GB
50014200.03%1.85 s0.9 GB
100031800.11%1.91 s1.4 GB
200054100.48%2.13 s2.3 GB
400059803.71%3.42 s4.1 GB
8000542011.2%6.80 s7.9 GB

拐点在 2000 到 4000 之间。超过 2000 之后吞吐增长放缓,429 错误率开始指数上升,单次耗时也因为服务端排队而变长。8000 并发的吞吐反而比 2000 低,因为大量请求在重试。

最后定在 2500,配合指数退避重试:

private EvalResult callWithRetry(String q, int maxRetry) throws Exception {
    long backoff = 500;
    for (int i = 0; i < maxRetry; i++) {
        try {
            return callModel(q);
        } catch (RateLimitedException e) {
            Thread.sleep(backoff + ThreadLocalRandom.current().nextInt(200));
            backoff = Math.min(backoff * 2, 8000);
        }
    }
    throw new IllegalStateException("retry exhausted: " + q);
}

注意这里的 Thread.sleep 在虚拟线程里是协作式挂起,不会占用载体线程。这是虚拟线程相对平台线程又一个优势——退避等待也是"免费"的。

内存:容易被忽略的成本

2500 并发下客户端内存 2.3 GB。拆开看,主要是三块:

# 用 JFR 的对象统计
jfr print --events ObjectCount vt.jfr | head -12
  byte[]            1.42 GB    (HTTP 响应缓冲 + 请求体)
  VirtualThread     0.38 GB    (25 万条任务分批,峰值 2500 个线程对象)
  String            0.29 GB

每条约 48 KB 的响应,2500 条在飞就是 120 MB。加上重试队列里攒的,峰值到 1.4 GB 的字节数组。

优化手段是流式消费响应体,不要一次性读成 String:

// 之前:整个响应读进内存
HttpResponse<String> resp = HTTP.send(req, BodyHandlers.ofString());

// 之后:流式解析,边收边处理
HttpResponse<Stream<String>> resp = HTTP.send(req, BodyHandlers.ofLines());
try (Stream<String> lines = resp.body()) {
    return lines.filter(l -> l.startsWith("data: "))
                .map(this::parseChunk)
                .reduce(new EvalResult(), EvalResult::merge);
}

内存降到 1.1 GB。对评测这种"跑完就丢"的任务意义不大,但如果是在线服务的批量接口,这个差别就决定了会不会 OOM。

流式输出场景:虚拟线程还香不香

评测任务是拿到完整响应就行,但在线服务大多是流式输出(SSE),要一直读流直到结束。这种场景虚拟线程的表现我单独测了一遍。

EvalResult callStreaming(String q) throws Exception {
    HttpRequest req = HttpRequest.newBuilder(URI.create(ENDPOINT))
            .header("Accept", "text/event-stream")
            .POST(BodyPublishers.ofString(buildBody(q)))
            .build();

    // 用 InputStream 逐行读,虚拟线程在 read() 上会协作式挂起
    HttpResponse<InputStream> resp = HTTP.send(req, BodyHandlers.ofInputStream());
    try (BufferedReader r = new BufferedReader(
            new InputStreamReader(resp.body(), UTF_8))) {
        String line;
        StringBuilder acc = new StringBuilder();
        while ((line = r.readLine()) != null) {
            if (line.startsWith("data: ")) acc.append(parseDelta(line));
        }
        return parse(acc.toString());
    }
}

关键在于 InputStream.read() 在虚拟线程里阻塞时,JDK 会把底层的 socket 注册到 NIO selector,虚拟线程挂起、载体线程释放。实测 2500 个并发流式请求,载体线程占用只有 8 个(就是核数),CPU 22%。

场景并发载体线程占用CPU内存
非流式(完整响应)2500834%1.1 GB
流式(平均 42 秒)2500822%0.7 GB

流式场景内存反而更低,因为响应是边读边处理,不用攒完整的 buffer。这一点跟很多人的直觉相反——大家觉得流式连接持有时间长,资源占用应该更高,实际上只要阻塞操作是协作式的,持有连接本身几乎不消耗资源

跟 CompletableFuture 异步写法对比

不用虚拟线程的话,同样的并发要写成响应式风格:

// 异步写法:可读性差,调试困难
List<CompletableFuture<EvalResult>> fs = questions.stream()
    .map(q -> HTTP.sendAsync(buildReq(q), BodyHandlers.ofString())
                .thenApply(r -> parse(r.body()))
                .exceptionally(this::fallback))
    .toList();

// 虚拟线程写法:同步代码,同样的并发度
try (var executor = Executors.newVirtualThreadPerTaskExecutor()) {
    List<Future<EvalResult>> fs = questions.stream()
        .map(q -> executor.submit(() -> callWithRetry(q, 3)))
        .toList();
}

两种方式的吞吐实测差距在 3% 以内(异步 5240 条/分钟,虚拟线程 5410 条/分钟)。但异步写法下,一次请求的逻辑被拆成了多个 lambda,异常堆栈里全是无意义的 CompletableFuture 帧,加个重试逻辑要嵌套两层 handle。虚拟线程版本就是普通的 for 循环加 try-catch。

这个差距在小脚本里不明显,但我们的在线批量接口用的是同一套代码,那个接口有 400 行业务逻辑,异步版本当初写了 700 多行,改需求的时候谁都不想碰。

结构化并发:让失败可控

我们用 JDK 21 的 StructuredTaskScope(预览特性,评测脚本这种非生产场景用)管理批次:

List<EvalResult> runBatch(List<String> batch) throws Exception {
    try (var scope = new StructuredTaskScope.ShutdownOnFailure()) {
        List<Subtask<EvalResult>> tasks = batch.stream()
                .map(q -> scope.fork(() -> callWithRetry(q, 3)))
                .toList();

        scope.joinUntil(Instant.now().plusSeconds(120));

        List<EvalResult> ok = new ArrayList<>();
        for (Subtask<EvalResult> t : tasks) {
            if (t.state() == Subtask.State.SUCCESS) {
                ok.add(t.get());
            } else {
                failed.add(q);      // 失败的单条记录下来,最后统一补偿
                log.warn("failed: {}", t.exception().getMessage());
            }
        }
        return ok;
    }
}

好处是批次的生命周期是封闭的:退出 try 块时,这一批的所有虚拟线程必然已经结束,不会有孤儿任务泄漏到下一批。之前用裸的 executor.submit 时,我们出现过几次"脚本跑完了但 JVM 不退出"的情况。

最终结果

方案5 万条耗时CPU内存峰值失败率
200 固定线程池82 分钟(未跑完)14%0.6 GB1.2%
虚拟线程 + Semaphore(2000)11 分 40 秒31%2.3 GB0.48%
+ HTTP/2 + 流式响应9 分 20 秒34%1.1 GB0.51%

15 万次调用(三个模型)总共跑了 28 分钟,之前预计要 4 小时以上。

结果落盘:别在虚拟线程里写文件

15 万条结果要写 CSV。第一版我在每个任务里直接 append 到文件,结果 2500 个虚拟线程抢同一个 FileWriter,加了 synchronized 之后又回到 pinning 的老路。

改成生产者消费者模式:虚拟线程只负责把结果放进队列,一个专门的写线程消费:

BlockingQueue<EvalResult> sink = new ArrayBlockingQueue<>(10_000);

// 写线程,只有一个
Thread writer = Thread.ofPlatform().daemon().start(() -> {
    try (var out = new BufferedWriter(Files.newBufferedWriter(path), 1 << 16)) {
        while (true) {
            EvalResult r = sink.take();
            if (r == EvalResult.POISON) break;
            out.write(r.toCsv());
        }
    } catch (Exception e) { log.error("writer failed", e); }
});

// 虚拟线程任务里只做 put
executor.submit(() -> {
    sink.put(callWithRetry(q, 3));
});

注意写线程用的是平台线程Thread.ofPlatform())而不是虚拟线程。因为它是长时间运行且大部分时间在阻塞 IO,用虚拟线程没有收益,反而多一层调度。虚拟线程适合"短任务、高并发",常驻的后台线程用平台线程就行

这个改动之后落盘的速度从 3400 条/秒提到 12000 条/秒,而且没有再出现文件锁竞争。

几条值得记的结论

  1. 换虚拟线程之前,先把限流和连接池里的 synchronized 清掉。 我第一次改完并发度几乎没变,就是被这个坑的。用 -Djdk.tracePinnedThread=full 能在开发期直接定位。
  2. 并发度要实测拐点,不能拍脑袋定。我们的拐点在 2000 到 4000 之间,超过之后吞吐反而下降,因为触发了供应商限流。不同供应商的配额不一样,这个值要自己测。
  3. 连接池确实"不再是瓶颈",但指的是它不再需要按线程数配置。HTTP/2 多路复用下,4 条连接就能撑 2500 并发。 注意这里的前提是供应商支持 HTTP/2,不支持的话还是要按并发数配连接池。
  4. 内存是新的约束。2500 并发下响应缓冲能吃掉 1.4 GB。批量场景用流式响应消费,别一次性读成 String。
  5. Thread.sleep 在虚拟线程里是免费的,这让指数退避重试的实现简单了很多,不用再搞额外的定时任务线程池。
  6. 用结构化并发管批次,避免任务泄漏。这一点在小脚本上看着不起眼,但对在线服务是必需的。

虚拟线程没有让系统"更快",它让系统"更能等"。IO 等待期间的线程成本降到接近于零之后,真正的瓶颈就变成了下游配额、内存和网络协议——这些是换线程模型解决不了的。

参考