我们的内部 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.deferinterval 这类操作符构造即"待命",defer 保证只有 WebFlux 真正订阅时才执行内部代码,此时依赖对象(listener)已就位。
  • merge vs concatmerge 是并行汇流,谁先产出谁先到;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 根因机制

我们的自定义 ChatModelFlux.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 清理是这次才补上的)。


五、守则

  1. Flux.create 的 consumer 里绝不能同步阻塞。 它在订阅线程上执行;一旦上游链路有 subscribeOn,你阻塞的就是背压信号的调度线程。阻塞 I/O 一律放虚拟线程/专用线程,consumer 只做启动与注册清理(onDispose)。
  2. "精确卡在 2 的幂次附近"是背压额度耗尽的签名(≈32/256 → 查 prefetch;任意值 → 查阻塞)。
  3. "数据不动"时先问 request 在哪、在谁的线程上,而不是给上游降速——给上游降速的"补丁"只会掩盖下游需求信号已死的事实。
  4. 跨线程的流日志带线程名grep | awk | sort | uniq -c 数一下放行元素个数,五分钟出结论。
  5. 依赖框架的 reactive 内部行为要读实现(反编译比猜快):publishOnsubscribeOn 递归这些都不在官方文档里。
  6. 虚拟线程不是银弹,是分层工具:请求热路径(Netty 线程、背压管道)保持 reactive 心智;后台长生命周期执行体(Agent 循环、阻塞读泵)用虚拟线程换可读性。两者的边界——也就是 Flux.create 这种桥接点——是最需要小心对待的地方。

本文涉及的代码已脱敏为示意结构,类名/包名与线上实现无关。