我们的内部 Agent 平台后端是 Spring WebFlux 全栈响应式的:入口是 Netty 的 NIO 线程,出口是 SSE 流,中间是 Reactor 的 Flux 管道。但 LLM 调用这个场景和典型的响应式网关不太一样——一次调用动辄 10~60 秒,模型边想边吐 token,中间还可能穿插多轮工具调用。纯响应式写法在这里会非常拧巴,而 Java 21 的虚拟线程恰好是那个泄压阀。
这篇文章讲三件事:纯响应式在 LLM 场景为什么拧巴、我们如何用"主链路 reactive + 后台执行体虚拟线程"的混用模型解决它、以及混用时踩过的一个典型 Reactor 背压大坑。
一、为什么纯响应式在 LLM 场景拧巴
经典 WebFlux 的心智模型是:请求快进快出,I/O 等待全部交给事件循环,线程数 ≈ CPU 核数。这套模型在"网关转发 / CRUD"场景无敌,但 LLM 场景有三个错位:
1. LLM 客户端 SDK 本身是"半响应式"的。 以 Spring AI 为例,ChatClient.stream() 返回 Flux<ChatResponse>,看起来很 reactive,但其内部——尤其是工具调用递归(模型要求调工具 → 执行工具 → 把结果喂回模型 → 继续流式输出)——是靠 subscribeOn(boundedElastic) 层层嵌套订阅实现的。工具回调(查库、调 MCP、读 Redis)大多是同步阻塞代码,SDK 把它们统统甩到 boundedElastic 弹性线程池上执行。你以为自己在写纯响应式,其实底层已经在偷偷开线程池了。
2. Agent 执行体天然是"循环 + 阻塞"结构。 一个子 Agent 的生命周期是:while (step < maxSteps) { 等消息 → 调 LLM → 处理结果 }。这里面有两个本质阻塞点:等消息(轮询收件箱,可能等 60 秒)和等 LLM(流式收完一整轮,可能 10~60 秒)。把这种结构硬翻成响应式,就得把 while 拆成 expand/repeatWhen 之类的操作符套娃,代码可读性直线下降,而收益为零——因为它根本不在请求热路径上。
3. 阻塞等待的并发量级可能很大。 多 Agent 协作时,一个主 Agent 可以同时唤醒几十个子 Agent,每个都在"阻塞"等 LLM 或等消息。用平台线程扛,每个 1MB 栈;用固定线程池扛,排队就是延迟。
结论:LLM 场景需要的是"出口响应式、内部随便阻塞"——对浏览器的 SSE 出口必须是非阻塞流(连接多、长尾、需要背压),而内部执行体只想写最朴素的同步代码。虚拟线程让后者重新变得可行。
二、虚拟线程:阻塞代码的泄压阀
虚拟线程(Java 21 正式 GA)由 JVM 调度,M:N 映射到少量平台线程(Carrier)上:
平台线程: 虚拟线程:
OS 线程 1 ──1:1── Java Thread 1 OS 线程 (Carrier)
OS 线程 2 ──1:1── Java Thread 2 │
┌────┼────┐ ← JVM 内部调度
│ │ │
VT 1 VT 2 VT 3 ... VT N
虚拟线程执行阻塞 I/O(网络读、sleep、.block())时,JVM 把它从 Carrier 上卸载,Carrier 去跑别的虚拟线程;I/O 就绪后再挂回某个空闲 Carrier 继续。成本对比:
| 平台线程 | 虚拟线程 | |
|---|---|---|
| 创建时间 | ~1ms | ~1µs |
| 栈内存 | ~1MB | 堆上按需增长 |
| 10 万个同时阻塞 | ~100GB,不可行 | 几乎只占堆内存,可行 |
这意味着:"每个执行体一个线程、线程内随便阻塞" 这个最朴素的编程模型,在 Java 21 之后重新成立了。
三、我们的混用模型:出口 reactive,内部虚拟线程
3.1 主链路:保持 Flux 管道
主 Agent 的 SSE 流跑在 reactor-http-nio 线程上,必须保持非阻塞。核心模式是 Flux.create 桥接——把 LLM SDK 的回调式流接进 Reactor 管道:
var aiFlux = Flux.<String>create(emitter -> {
chatClient.prompt()
.system(ctx.systemPrompt)
.messages(historyMsgs)
.user(message)
.stream().chatResponse() // Spring AI 的流式回调
.flatMapSequential(cr -> { ... })
.subscribe(
emitter::next, // LLM 出一帧 → 推一帧进 Flux
emitter::error,
emitter::complete
);
}, FluxSink.OverflowStrategy.BUFFER); // 下游慢了就缓冲
然后用心跳流、子 Agent 消息轮询流与之合并,一条 SSE 发出去:
var heartbeat = Flux.interval(Duration.ofSeconds(10))
.map(i -> formatSSEEvent("h", Map.of()))
.takeUntilOther(chatFlux.then().flux()); // 主流结束,心跳自动停
Flux<String> tailEvents = Flux.defer(() -> // defer:订阅时才执行,避免构造时 listener 还没就位
Flux.interval(Duration.ofSeconds(2))
.handle((tick, sink) -> {
var msgs = messageBus.readInbox(sessionId, agentId);
if (msgs != null) for (var m : msgs) sink.next(format(m));
if (allSubAgentsExited) sink.complete();
}));
return Flux.merge(chatFlux, heartbeat, tailEvents)
.doFinally(signal -> listenerRegistry.unregister(sessionId));
几个要点:
Flux.defer:interval这类操作符构造即"待命",defer保证只有 WebFlux 真正订阅时才执行内部代码,此时依赖对象(listener)已就位。mergevsconcat:merge是并行汇流,谁先产出谁先到;concat是严格串行。SSE 多路事件汇流用merge。share()/publish(1).refCount(1):LLM 流是冷流,被订阅两次就会触发两次 LLM 调用(一次几百秒、一份账单),必须转热。
3.2 后台执行体:虚拟线程里写同步代码
子 Agent 不在请求热路径上,直接一人一虚拟线程:
// 消息监听器收到唤醒通知后
Thread.ofVirtual()
.name("sub-agent-" + agentId + "-" + sessionId)
.start(() -> {
var runner = new SubAgentRunner(/* 配置、总线、工具、emitter... */);
runner.run(); // 同步阻塞方法,跑在虚拟线程上
});
执行体内部就是最朴素的 while + sleep + block:
public void run() {
RequestContext.setCurrentSessionId(sessionId); // ① ThreadLocal 上下文
var memory = new SessionMemory();
int step = 0;
while (step < config.maxSteps() && !shouldShutdown) {
List<AgentMessage> inbox = messageBus.readInbox(sessionId, agentId);
if (inbox.isEmpty()) {
inbox = waitForMessage(); // ② 内部 sleep(1s) 轮询,最多 60s
}
var result = runSingleRound(systemPrompt, memory, inboxText, tools); // ③ 同步调 LLM
// ...处理结果、发消息回主 Agent...
}
messageBus.send(/* final result */);
RequestContext.clear(); // ④ 清理 ThreadLocal
}
同步调 LLM 长这样——注意 .block() 写在虚拟线程上,"阻塞"的只是虚拟线程,Carrier 早被释放了:
var response = chatClientBuilder.build().prompt()
.system(systemPrompt).messages(history).user(userMessage)
.tools(tools)
.stream().chatResponse()
.collectList() // 流式收完一整轮
.block(LLM_TIMEOUT); // 通常 10~60 秒,期间不占平台线程
3.3 全景:三种线程各司其职
用户: "帮我分析 Q2 数据,画两张图表"
│
▼
主 Agent (reactor-http-nio 线程,Flux 管道)
│
├─ send_message(to=alice) ─→ Virtual Thread sub-agent-alice
│ while → .block() 等 LLM (10s)
├─ send_message(to=bob) ─→ Virtual Thread sub-agent-bob
│ while → .block() 等 LLM + 生图 (18s)
└─ send_message(to=charlie)─→ Virtual Thread sub-agent-charlie
while → .block() 等 LLM + 生图 (20s)
总耗时 = max(10, 18, 20) = 20s ← 真并行
平台线程占用 ≈ 0 ← 都在等 I/O,Carrier 全释放
| 线程类型 | 用途 | 数量 |
|---|---|---|
reactor-http-nio-* |
WebFlux 网络线程,SSE 出口 | ≈ CPU 核数 |
boundedElastic-* |
Spring AI 工具回调、SDK 内部调度 | 动态(有界弹性) |
Virtual Thread (msg-listener-*) |
每会话的消息监听 | 1 / 会话 |
Virtual Thread (sub-agent-*) |
每子 Agent 的运行循环 | 1 / 子 Agent |
parallel-* |
Reactor 定时调度(interval 等) | 全局共享 |
对比传统线程池方案:
| 场景 | 固定线程池 (10) | 虚拟线程 |
|---|---|---|
| 10 个子 Agent 并发 | 刚好够用 | 无压力 |
| 50 个子 Agent 并发 | 后 40 个排队 | 全部立即并行 |
| 100 个空闲等待中 | 线程池耗尽 | Carrier 全释放,零压力 |
ThreadLocal 注意:每个虚拟线程有独立 ThreadLocal。执行体开头 RequestContext.set(...)、结束 clear()。虚拟线程每次新建,不存在线程池复用导致的 ThreadLocal 残留问题——这也是它比线程池省心的一点。
四、踩坑实录:Flux.create 里阻塞读,背压饿死 48 秒
混用模型最大的坑,出在两种世界的交界处。我们曾遇到一个线上症状:长回复轮次的思考流中途冻结 48 秒,随后瞬间爆发式刷出。前端看是"分析中"假死,实际模型一直在正常生成。
4.1 证据链
- 读流线程(
boundedElastic-3)的 chunk 日志显示冻结期间模型在正常吐 token——数据到了后端,只是没流向下游; - 下游所有处理日志在另一个线程(
boundedElastic-1)——同一条流跨了线程,中间必有带队列的异步边界(查实现确认是 SDK 内部的.publishOn(boundedElastic())); - 卡死前放行的 chunk 数精确命中 256——
publishOn默认预取额度Queues.SMALL_BUFFER_SIZE = 256,消费 75% 后向上游补request(192)。"恰好 2 的幂次后断流"是预取额度耗尽、补货请求丢失的标准签名; - 恢复时刻与流结束时刻只差 1ms:读循环收到
[DONE]退出 → 占用的 worker 释放 → 积压 token 倾泻。不是巧合,是因果。
4.2 根因机制
我们的自定义 ChatModel 在 Flux.create 的 consumer 里同步 while(readLine()) 阻塞读 HTTP 流:
工具递归:toolCallFlux.subscribeOn(boundedElastic)
│ 第 N+1 轮"执行工具 + 订阅下一轮流"作为一个任务在 worker W 上执行
▼
Flux.create 的 consumer 在订阅时【同步执行】
│ consumer 里是阻塞读循环 —— 任务不返回,W 被占死直到 [DONE]
▼
subscribeOn 把下游所有 request(n) 信号调度到【同一个 W】
│ publishOn 发出的补货 request(192) 排在阻塞任务后面,永远执行不到
▼
publishOn 预取的 256 个发完 → 上游断粮,后续 token 堆在 Flux.create(BUFFER) 的无界缓冲里
▼
[DONE] → 读循环退出 → W 释放 → 排队 request 执行 → 缓冲倾泻 → 前端爆发式刷屏
三个必要条件缺一不可:consumer 内同步阻塞(我们的代码)+ 上游被 subscribeOn 包裹使 request 信号被 marshal 到同一 worker(SDK 工具递归的实现)+ 中间有固定 prefetch 的 publishOn(SDK 内部)。首轮调用不触发,是因为订阅线程是 Netty 的 NIO 线程,request 走的路径不同;只有工具递归后的轮次才满足全部条件——所以短回复永远复现不了,前两次排查都误诊成了"上游太快/节流太慢"。
4.3 修复:阻塞读移进虚拟线程
return Flux.create(sink -> {
var connRef = new AtomicReference<HttpURLConnection>();
Thread reader = Thread.startVirtualThread(() -> {
// 原 HTTP 请求 + readLine 循环
// sink.next / sink.complete / sink.error
});
// 客户端断开时强制断开连接,使阻塞中的 readLine 抛异常退出
sink.onDispose(() -> {
HttpURLConnection c = connRef.get();
if (c != null) c.disconnect();
reader.interrupt();
});
});
consumer 立即返回 → worker 不再被占用 → request 信号正常流转 → publishOn 持续补货 → token 全程实时。顺带还修了一个既有泄漏:此前客户端断开后读循环仍会把整条 LLM 流读完(onDispose 清理是这次才补上的)。
五、守则
Flux.create的 consumer 里绝不能同步阻塞。 它在订阅线程上执行;一旦上游链路有subscribeOn,你阻塞的就是背压信号的调度线程。阻塞 I/O 一律放虚拟线程/专用线程,consumer 只做启动与注册清理(onDispose)。- "精确卡在 2 的幂次附近"是背压额度耗尽的签名(≈32/256 → 查 prefetch;任意值 → 查阻塞)。
- "数据不动"时先问 request 在哪、在谁的线程上,而不是给上游降速——给上游降速的"补丁"只会掩盖下游需求信号已死的事实。
- 跨线程的流日志带线程名,
grep | awk | sort | uniq -c数一下放行元素个数,五分钟出结论。 - 依赖框架的 reactive 内部行为要读实现(反编译比猜快):
publishOn、subscribeOn递归这些都不在官方文档里。 - 虚拟线程不是银弹,是分层工具:请求热路径(Netty 线程、背压管道)保持 reactive 心智;后台长生命周期执行体(Agent 循环、阻塞读泵)用虚拟线程换可读性。两者的边界——也就是
Flux.create这种桥接点——是最需要小心对待的地方。
本文涉及的代码已脱敏为示意结构,类名/包名与线上实现无关。