1. 现象

用户在一个 JSON-contract Agent(策略卡生成类)会话中观察到:

  • 思考(thinking)token 实时流出到 …现在我需要基于这些数据生成策略突然冻住
  • 冻住约 48 秒后,思考从 卡。任务…爆发式续上——注意 策略 + = 策略卡,是同一句话被拦腰截断,说明模型没停,是管道断了;
  • 日志里 21:33:48 这一秒内涌出 4134 行,之前 33:00→33:48 几乎完全静默。

前端表现为"分析中"假死,实际模型一直在正常生成。

2. 排查时间线与证据链

证据 1:卡死期间模型并没有停

读流线程(boundedElastic-3)上的 "chunk 到达即打印" 日志显示,21:33:15→21:33:48 模型在以正常节奏持续生成 JSON 内容。数据在到达后端,只是没有流向下游。

这一条直接排除了"模型推理慢 / 上游网关卡住"类假设。

证据 2:下游处理线程与读流线程不同

  • 读流(readLine 循环 + sink.next):boundedElastic-3
  • 下游所有 Advisor 日志与 SSE 出口日志:boundedElastic-1

同一条流跨了线程 → 中间必有带队列的异步边界。反编译 Spring AI 的 client-chat 包确认:ChatModelStreamAdvisor 内部是 chatModel.stream(prompt).publishOn(Schedulers.boundedElastic()).map(...) —— 这个 publishOn 就是那个边界,队列在它手里,drain 在它的 worker(线程 -1)上跑。

证据 3:卡死点精确命中 prefetch=256

统计卡死前第 3 轮 LLM 调用放行的 chunk 数:247 个非空 token(加上若干空白/元数据 chunk ≈ 256)。publishOn 默认预取 Queues.SMALL_BUFFER_SIZE = 256,消费到 75% 时向上游补 request(192)

"恰好 256 个后断流"是预取额度耗尽、补货请求丢失的标准签名。

证据 4:恢复时刻与流结束时刻只差 1ms

  • 21:33:48.451:读流线程收到最后一个 chunk },随后 [DONE] → 读循环退出;
  • 21:33:48.452:积压 token 开始在下游 worker 上倾出。

不是巧合,是因果:读循环退出 ⇒ 占用的 worker 释放 ⇒ 排队的 request 任务终于执行 ⇒ 缓冲倾泻。

证据 5:为什么只有第 3 轮卡

第 1、2 轮(工具调用轮)各自不足 256 个 chunk 就结束了,预取额度没用完;第 3 轮是"长思考 + 长 JSON",超过 256 后必卡。短回复永远复现不了——这也解释了为什么这个 bug 能潜伏这么久。

3. 根因机制

完整因果链(Spring AI + Reactor):

ToolCallingAdvisor 工具递归:toolCallFlux.subscribeOn(boundedElastic)
   │  第 N+1 轮的"执行工具 + 订阅下一轮流"作为一个任务在 worker W 上执行
   ▼
ChatModel.stream 的 Flux.create consumer 在订阅时【同步执行】
   │  consumer 里是 while(readLine()) 阻塞循环 —— 任务不返回,W 被占死直到 [DONE]
   ▼
FluxSubscribeOn 把下游所有 request(n) 信号 worker.schedule 到【同一个 W】
   │  publishOn 消费 192 个后发出的补货 request(192) 排在阻塞任务后面,永远执行不到
   ▼
publishOn 预取的 256 个发完 → 上游断粮
   │  后续 token 全部堆在 Flux.create(BUFFER) 的无界缓冲里
   ▼
[DONE] → 读循环退出 → W 释放 → 排队 request 执行 → FluxCreate 缓冲倾泻
   → publishOn 在自己的 worker(另一线程)上一次性 drain 全部积压

三个必要条件,缺一不可:

  1. Flux.create consumer 内同步阻塞读(我们自己的代码);
  2. 上游被 subscribeOn 包裹,request 信号被 marshal 到同一 worker(Spring AI 工具递归的实现);
  3. 中间有固定 prefetch 的 publishOn(Spring AI 的 ChatModelStreamAdvisor)。

首轮调用不触发是因为订阅线程是 reactor-http-nio(Netty event loop),request 走的路径不同;只有工具递归后的轮次满足条件 2。这也是为什么"普通对话不卡、用了工具的长回复才卡"。

关键点在于 Reactor 的一个容易忽略的行为:subscribeOn(scheduler) 不只是把订阅动作搬到 scheduler 上——它会把所有向上游传播的 request(n) 信号也作为任务 schedule 到同一个 worker。当那个 worker 被一个永不返回的阻塞任务占住时,整条链路的背压协议实质死亡:上游以为下游不要数据,下游以为上游没数据,两边都在等一个永远不会执行的任务。

4. 两次误诊史(本次一并纠正)

时间 误诊 错误补丁 实际情况
第一次 "一次 read() 预读太多 SSE 行导致积压" BufferedReader 缓冲改成 64B 64B 只是字符缓冲,下层 InputStreamReader 自带 8KB 字节缓冲,补丁从未生效;积压也不在这层
第二次 "delayElements(30ms) 造成思考 token 积压卡死" 移除 delayElements delayElements 只是让消费变慢、更快暴露 256 额度耗尽;根因是背压饿死

共性教训:症状是"积压后爆发",两次都在"让上游慢一点"上做文章,而真正的问题是"下游的需求信号死了"。 Reactor 里"数据不动"要先问"request 在哪、在谁的线程上"。

还有一个组织层面的失误:两次"修复"后都没有把复现路径(长回复 + 工具递归轮)固化成验证手段,所以"看起来缓解了"的错觉得以存续。

5. 修复

5.1 根因修复:阻塞读循环移入虚拟线程

自定义 ChatModel 的 stream() 改为:阻塞读循环放入独立虚拟线程Flux.create consumer 立即返回:

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();
    });
});

效果:

  • worker 不再被占用 → request 信号正常流转 → publishOn 持续补货 → token 全程实时;
  • 顺带修复一个既有泄漏:此前客户端断开后读循环仍会把整条 LLM 流读完(onDispose 清理逻辑之前不存在)。

用虚拟线程而非另一个 boundedElastic 任务,是因为虚拟线程天然适合"一个连接一个阻塞读循环"的模型,不占用 Reactor 调度器容量,也不会再引入新的调度边界。

5.2 恢复输出平滑节流

根因修复后,节流不再有积压风险,按产品需求加回(15~30ms 随机间隔,先到的数据在上游 Flux.create(BUFFER) 自然缓冲,出口匀速下发):

var chatFlux = buildChatFlux(...)
        .concatMap(ev -> Mono.just(ev).delayElement(
                Duration.ofMillis(ThreadLocalRandom.current().nextLong(15, 31)),
                Schedulers.boundedElastic()), 32)
        .publish(1).refCount(1);

注意这里用 concatMap + delayElement 而不是 delayElements:前者只延迟发射、并发预取上限 32,不会在背压协议上叠加额外的不确定性;publish(1).refCount(1) 替代 share(),避免其内部 prefetch=256 重新攒批。

5.3 排查期间沉淀的辅助改动

  • SSE 出口统一 tapFlux.merge(...).doOnNext(e -> log.debug("[SSE] >> session={} {}", sid, e))——本次定位的关键观测点,予以保留;
  • 本地/开发环境日志同时落文件,事后统计(如"每秒日志行数直方图")全靠它;
  • 工具回调加容错壳:模型输出非法 JSON 参数时把错误文本还给模型重试,不再打断整条流(另一个独立 bug,同日修复);
  • 流错误优雅收尾:异常时改为正常 complete 事件收集器与 SSE emitter,不再向已 committed 的 SSE 响应传播异常(传播了前端也收不到,只会刷一条无意义的 error 日志)。

6. 经验与守则

  1. Flux.create 的 consumer 里绝不能同步阻塞。 它在订阅线程上执行;一旦上游链路有 subscribeOn,你阻塞的就是背压信号的调度线程。阻塞 I/O 一律放虚拟线程/专用线程,consumer 只做启动与注册清理。
  2. "精确卡在 2 的幂次附近"是背压额度耗尽的签名。 卡死前数一下放行了多少元素:≈32/256 → 查 prefetch;任意值 → 查阻塞。
  3. "恢复时刻与某事件只差 1ms"不是巧合。 先找与恢复精确同时发生的事件,再倒推谁在等它。
  4. 跨线程的流日志要带线程名分析。 本次"读在 -3、下游在 -1"直接指出了队列边界的存在;grep | awk '{print $3,$5}' | sort | uniq -c 五分钟出结论。
  5. 补丁式修复要留"根因存疑"标记。 64B 缓冲和移除 delayElements 都"看起来缓解了",实则掩盖问题。症状复现路径没被固化成验证手段,是两次误诊的直接原因。
  6. 依赖框架的 reactive 内部行为要读实现。 ChatModelStreamAdvisorpublishOn、工具递归的 subscribeOn 都不在官方文档里;反编译(javap -c | grep invoke.*Flux)比猜快。