LangChain4j流式输出的背压控制与TokenStream消费实战
你兴冲冲地在产品里接上了大模型的流式输出,用户一打字,文字逐字往外蹦,体验丝滑。上线第一周相安无事。第二周用户量翻了三倍,晚上八点高峰时段,运维告警炸了——CPU 飙升到 95%,GC 日志像刷屏一样疯狂滚动,部分用户的流式响应卡住不动,甚至直接断开连接。你一查堆栈,好家伙,几百个线程都在往同一个阻塞队列里塞 token,消费者根本处理不过来。这不是大模型变慢了,是你的流式管道被自己写崩了。流式输出听着简单,token 来个打个,但生产环境里的流量洪峰面前,没有背压控制的流式管道就是一颗定时炸弹。今天就聊聊 LangChain4j 里的 TokenStream 怎么玩,怎么在生产环境里不翻车。
一、这个问题到底是什么
流式输出的消费方式本质上是一个生产者-消费者问题。大模型作为生产者,以 token 为单位不断产出文本片段;你的业务代码作为消费者,需要把这些片段收集起来,推给前端做 SSE(Server-Sent Events)展示。问题出在中间这个"传话筒"环节——处理速度不对等。
大模型生成 token 的速度一般在每秒 20-100 个不等,取决于模型大小和硬件。前端消费 SSE 的速度几乎是瞬时的,因为就是往 HTTP response 里写几个字节。看起来消费者比生产者快得多,应该不会有问题?错。真正慢的不是"写 response",而是你夹在中间的"处理逻辑"。
举个例子:每个 token 到达时,你可能要调用内容安全审核接口、要写数据库日志、要更新缓存计数、要触发 WebSocket 广播。这些操作任何一个延迟超过 10 毫秒,在 100 token/秒的流速下就会开始积压。HTTP response 的输出缓冲区是有限的,如果积压的 token 撑爆了缓冲区,操作系统就会断开 TCP 连接,用户端直接报错。
更隐蔽的问题是内存泄漏。很多人习惯把流式响应的每个片段存到一个 StringBuilder 或者 ArrayList 里,等流结束后做后处理。如果用户关闭了浏览器标签页,后端还傻傻地往这个集合里塞数据,GC 根本来不及回收,堆内存就一点一点被吃掉。
LangChain4j 提供了两套流式消费模型。一套是底层 StreamingChatModel 配合 StreamingChatResponseHandler,每个 token 到达时触发 onPartialResponse 回调——这是"推"模型。另一套是 TokenStream,它实现了 Java 的 Iterable 接口,你可以用 for-each 循环主动拉取 token——这是"拉"模型。推模型适合简单场景,拉模型配合 Project Reactor 或 RxJava 才能实现真正的背压控制。
核心问题一句话:流式输出不能"来一个处理一个",必须有一套机制让消费者按自己的节奏取数据,快则多取,慢则少取,取不动的就让上游暂停。
二、底层原理到底怎么回事
要理解背压怎么工作,先得理解数据在框架内部是怎么流的。
LangChain4j 的 StreamingChatModel 底层通过 HTTP 长连接与 LLM 服务端通信。以 OpenAI 为例,请求里设置 stream: true,响应头里 Content-Type: text/event-stream,服务端不再一次性返回完整 JSON,而是以 SSE 格式逐块推送数据。每一块大概长这样:
data: {"id":"chatcmpl-xxx","object":"chat.completion.chunk","choices":[{"delta":{"content":"你好"}}]}
OpenAiStreamingChatModel 内部维护了一个 OkHttp 或者 HttpClient 的响应体输入流,开一个后台线程持续读这个输入流,每读到一行 SSE 事件就解析成 ChatCompletionChunk 对象,然后调用 StreamingChatResponseHandler 的对应回调方法。
这里就有一个关键细节:负责读 SSE 流的线程和调用 onPartialResponse 的是同一个线程。如果你的 onPartialResponse 里做了耗时操作——比如调了数据库、调了外部 API——这个线程就会被长时间占用。线程被占用意味着底层 TCP 接收缓冲区没人读,缓冲区满了之后,操作系统会根据 TCP 滑动窗口协议减小接收窗口,最终迫使上游 OpenAI 服务器暂停发送数据。这其实是一种被动的、粗粒度的背压——由操作系统的 TCP 流控替你背了锅。
但这个机制太粗糙了。TCP 缓冲区大小通常在 64KB 到几 MB 之间,而且在 JVM 层面你根本感知不到——线程就是卡在 onPartialResponse 里了,你不知道是因为自己处理慢还是网络慢。而且一旦线程恢复,积压在缓冲区里的数据会一口气涌进来,形成"雪崩"效应。
TokenStream 的设计思路不一样。它内部用了一个 BlockingQueue 做缓冲,底层 StreamingChatModel 的回调线程把 token 往队列里塞,你的消费线程从队列里取。按说是解耦了生产和消费,但 BlockingQueue 的默认行为是:队列满了,生产者线程阻塞。如果你的消费线程因为某种原因暂时停止消费(比如前端断连),生产者线程就会一直阻塞,TCP 流控被动生效。等你恢复消费,又得面对队列里积压的几百个 token。
更好的做法是用 Reactive Streams 规范。LangChain4j 1.18 版本里,TokenStream 提供了 toFlowPublisher() 方法,可以把 token 流转换成一个 JDK 9 的 Flow.Publisher。这套规范的核心就是背压协议:订阅者通过 subscription.request(n) 告诉生产者"我一次最多处理 n 个",处理完一批再请求下一批。Project Reactor 的 Flux 完全实现了这套协议,所以你可以在 Flux 链路上用 limitRate、onBackpressureBuffer、onBackpressureDrop 等操作符精确控制背压策略。
用代码来理解底层链路:OpenAI SSE → OkHttp 响应流 → 后台解析线程 → StreamingChatResponseHandler.onPartialResponse → TokenStream 内部队列 → Flow.Publisher → Reactor Flux → 你的业务逻辑 → SSE 推前端。这个链路上每一环都有可能出现速度不匹配,而 Reactive Streams 的背压协议让每一环都可以向上游发出"慢一点"的信号。
还有一个容易忽略的点:错误传播。StreamingChatModel 在遇到错误时会调用 onError,但如果你的 Flux 链路里某个操作符抛了异常,这个异常需要正确地终止整个流,而不是静默丢失。LangChain4j 的 TokenStream 在转为 Flow.Publisher 时会正确处理 onError 信号的传播,确保异常最终能被 Flux 的 doOnError 或 onErrorResume 捕获。
三、实战:手把手写代码
示例一:基础流式输出——推模型(StreamingChatModel)
这个例子演示最基础的用法:用 StreamingChatModel 配合 StreamingChatResponseHandler 做流式对话,token 直接打印到控制台。适合快速验证、控制台测试,不适合生产环境。
package com.example.streaming.basic;
import dev.langchain4j.model.chat.StreamingChatModel;
import dev.langchain4j.model.chat.response.ChatResponse;
import dev.langchain4j.model.chat.response.StreamingChatResponseHandler;
import dev.langchain4j.model.openai.OpenAiStreamingChatModel;
import dev.langchain4j.model.openai.OpenAiStreamingChatModel.OpenAiStreamingChatModelBuilder;
public class BasicStreamingExample {
public static void main(String[] args) throws InterruptedException {
OpenAiStreamingChatModel model = OpenAiStreamingChatModel.builder()
.apiKey(System.getenv("OPENAI_API_KEY"))
.modelName("gpt-4o-mini")
.build();
System.out.println("=== 开始流式输出 ===");
model.chat("用三句话介绍Java虚拟机的垃圾回收机制", new StreamingChatResponseHandler() {
@Override
public void onPartialResponse(String partialResponse) {
System.out.print(partialResponse);
}
@Override
public void onCompleteResponse(ChatResponse completeResponse) {
System.out.println("\n=== 输出完成 ===");
System.out.println("总Token消耗: " + completeResponse.metadata().tokenUsage().totalTokenCount());
}
@Override
public void onError(Throwable error) {
System.err.println("流式输出出错: " + error.getMessage());
}
});
// 因为是异步的,主线程等一下
Thread.sleep(30000);
}
}
这个代码跑起来没问题,但问题也明显:onPartialResponse 在生产线程里直接执行,你没法控制消费速率。假设你要把每个 token 写入数据库——在 onPartialResponse 里加一行 jdbcTemplate.update(...),恭喜你,单线程逐 token 写数据库,延迟直接爆炸。
示例二:拉模型——TokenStream 基础消费
TokenStream 提供了 Iterable 接口,可以用 for-each 循环主动拉取 token。注意:for-each 循环是阻塞式的,适合在独立线程里跑,不适合在 Web 请求线程里用。
package com.example.streaming.tokenstream;
import dev.langchain4j.data.message.AiMessage;
import dev.langchain4j.model.StreamingResponseHandler;
import dev.langchain4j.model.chat.StreamingChatModel;
import dev.langchain4j.model.chat.response.ChatResponse;
import dev.langchain4j.model.chat.response.StreamingChatResponseHandler;
import dev.langchain4j.model.openai.OpenAiStreamingChatModel;
import dev.langchain4j.model.output.TokenStream;
import java.util.concurrent.CompletableFuture;
public class TokenStreamBasicExample {
public static void main(String[] args) {
var model = OpenAiStreamingChatModel.builder()
.apiKey(System.getenv("OPENAI_API_KEY"))
.modelName("gpt-4o-mini")
.build();
CompletableFuture<ChatResponse> futureResponse = new CompletableFuture<>();
// TokenStream 接收 StreamingChatResponseHandler 作为桥接
TokenStream tokenStream = new TokenStream();
model.chat("请用三句话概括微服务架构的优势", new StreamingChatResponseHandler() {
@Override
public void onPartialResponse(String partialResponse) {
tokenStream.onNext(partialResponse);
}
@Override
public void onCompleteResponse(ChatResponse completeResponse) {
tokenStream.onComplete();
futureResponse.complete(completeResponse);
}
@Override
public void onError(Throwable error) {
tokenStream.onError(error);
futureResponse.completeExceptionally(error);
}
});
// 主动拉取 token,可以在这里加自己的消费逻辑
System.out.println("=== TokenStream 拉取模式 ===");
for (String token : tokenStream) {
System.out.print(token);
// 这里的消费线程可以自由控制节奏
simulateSlowConsumer();
}
System.out.println("\n=== 拉取完成 ===");
ChatResponse response = futureResponse.join();
System.out.println("总Token: " + response.metadata().tokenUsage().totalTokenCount());
}
private static void simulateSlowConsumer() {
try {
Thread.sleep(5); // 模拟5ms的处理延迟
} catch (InterruptedException e) {
Thread.currentThread().interrupt();
}
}
}
TokenStream 的 for-each 模式虽然能主动控制消费节奏,但它有一个致命的缺点:底层依赖一个阻塞队列,如果生产者线程比消费者快太多,队列会不断膨胀。另外 for-each 是阻塞迭代,一旦消费线程卡住,整个消费链条就停了。要想真正做好生产级的流控,必须上 Reactor。
示例三:生产级方案——TokenStream + Reactor Flux + 背压控制
这是生产环境推荐的完整方案。核心思路:把 TokenStream 转为 Reactor Flux,利用 Flux 的背压操作符做流控,同时在 WebFlux 中以 SSE 格式推给前端。这个代码可以直接放到 Spring Boot WebFlux 项目里跑。
package com.example.streaming.reactor;
import dev.langchain4j.data.message.AiMessage;
import dev.langchain4j.model.chat.StreamingChatModel;
import dev.langchain4j.model.chat.response.ChatResponse;
import dev.langchain4j.model.chat.response.StreamingChatResponseHandler;
import dev.langchain4j.model.openai.OpenAiStreamingChatModel;
import dev.langchain4j.model.output.TokenStream;
import org.reactivestreams.Publisher;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import reactor.core.publisher.Flux;
import reactor.core.publisher.FluxSink;
import reactor.core.publisher.Mono;
import reactor.core.scheduler.Schedulers;
import java.time.Duration;
import java.util.concurrent.CompletableFuture;
public class ReactorTokenStreamExample {
private static final Logger log = LoggerFactory.getLogger(ReactorTokenStreamExample.class);
public static void main(String[] args) throws InterruptedException {
var model = OpenAiStreamingChatModel.builder()
.apiKey(System.getenv("OPENAI_API_KEY"))
.modelName("gpt-4o-mini")
.build();
Flux<String> tokenFlux = createTokenFlux(model, "请详细解释Java中CompletableFuture和Reactor Flux的区别");
tokenFlux
// 背压策略:缓冲64个token,满了就丢弃最旧的
.onBackpressureBuffer(64, token -> log.warn("背压丢弃token: {}", token.length()))
// 限流:每秒最多消费50个token
.limitRate(50)
// 每个token之间最小间隔20ms
.delayElements(Duration.ofMillis(20))
// 分离生产和消费线程
.publishOn(Schedulers.boundedElastic())
.doOnNext(token -> {
// 模拟业务处理(安全检查、日志写入等)
log.debug("处理token: {}", token);
})
.doOnComplete(() -> log.info("流式输出完成"))
.doOnError(error -> log.error("流式输出异常", error))
.subscribe(
System.out::print,
error -> System.err.println("订阅异常: " + error.getMessage()),
() -> System.out.println("\n=== 流结束 ===")
);
// 等待流结束
Thread.sleep(60000);
}
/**
* 将 StreamingChatModel 的响应转为 Reactor Flux<String>,带完整的背压支持。
*/
public static Flux<String> createTokenFlux(StreamingChatModel model, String prompt) {
return Flux.create(sink -> {
TokenStream tokenStream = new TokenStream();
model.chat(prompt, new StreamingChatResponseHandler() {
@Override
public void onPartialResponse(String partialResponse) {
// 底层回调线程把 token 推到 FluxSink
// FluxSink 内部实现了 Reactive Streams 背压协议
// 当下游消费不过来时,sink.next 会阻塞回调线程
// 进而阻塞底层 SSE 读取线程,实现端到端背压
sink.next(partialResponse);
}
@Override
public void onCompleteResponse(ChatResponse completeResponse) {
log.info("模型返回完成, Token用量: {}",
completeResponse.metadata().tokenUsage().totalTokenCount());
sink.complete();
}
@Override
public void onError(Throwable error) {
log.error("模型返回异常: {}", error.getMessage(), error);
sink.error(error);
}
});
}, FluxSink.OverflowStrategy.LATEST);
}
}
这个代码的关键点:
Flux.create 的 OverflowStrategy:这里用了 LATEST,意思是当下游消费跟不上、内部缓冲区满了之后,新来的 token 会覆盖缓冲区里最旧的 token。生产环境建议用 BUFFER(默认)加 onBackpressureBuffer 做有界缓冲,或者 DROP 直接丢弃不需要的 token。
sink.next 的阻塞行为:当 FluxSink 的下游(doOnNext、subscribe 等)处理不过来时,sink.next 会阻塞调用线程。这个调用线程就是 StreamingChatResponseHandler 的回调线程,它背后是 OkHttp 读取 SSE 流的线程。线程被阻塞 → TCP 接收缓冲区没人读 → 操作系统减小接收窗口 → OpenAI 服务器暂停发送。这样整条链路从 Java 层到 TCP 层实现了端到端背压,不会出现内存无限增长。
publishOn 与 subscribeOn:这里用 publishOn(Schedulers.boundedElastic()) 把消费逻辑从生产线程分离出来。boundedElastic 调度器会自动管理线程池大小,避免创建过多线程。
示例四:Spring Boot WebFlux 集成——完整 SSE 端点
把这个方案嵌入到 Spring Boot 3.x + WebFlux 项目里,一个完整的 REST 端点如下:
package com.example.streaming.controller;
import dev.langchain4j.model.chat.StreamingChatModel;
import dev.langchain4j.model.chat.response.ChatResponse;
import dev.langchain4j.model.chat.response.StreamingChatResponseHandler;
import dev.langchain4j.model.openai.OpenAiStreamingChatModel;
import dev.langchain4j.model.output.TokenStream;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import org.springframework.http.MediaType;
import org.springframework.web.bind.annotation.GetMapping;
import org.springframework.web.bind.annotation.RequestParam;
import org.springframework.web.bind.annotation.RestController;
import reactor.core.publisher.Flux;
import reactor.core.publisher.FluxSink;
import reactor.core.scheduler.Schedulers;
import java.time.Duration;
@RestController
public class StreamingChatController {
private static final Logger log = LoggerFactory.getLogger(StreamingChatController.class);
private final StreamingChatModel chatModel;
public StreamingChatController() {
this.chatModel = OpenAiStreamingChatModel.builder()
.apiKey(System.getenv("OPENAI_API_KEY"))
.modelName("gpt-4o-mini")
.build();
}
/**
* SSE 流式聊天端点。
* 访问方式: GET /api/chat/stream?prompt=你好
* 返回类型: text/event-stream
*/
@GetMapping(value = "/api/chat/stream", produces = MediaType.TEXT_EVENT_STREAM_VALUE)
public Flux<String> streamChat(@RequestParam(defaultValue = "请介绍一下你自己") String prompt) {
return Flux.<String>create(sink -> {
model.chat(prompt, new StreamingChatResponseHandler() {
@Override
public void onPartialResponse(String partialResponse) {
// 检查连接是否还存在,如果前端断开就停止推送
if (!sink.isCancelled()) {
sink.next(partialResponse);
}
}
@Override
public void onCompleteResponse(ChatResponse completeResponse) {
log.info("SSE流完成, prompt: {}, tokens: {}",
prompt.substring(0, Math.min(prompt.length(), 50)),
completeResponse.metadata().tokenUsage().totalTokenCount());
sink.complete();
}
@Override
public void onError(Throwable error) {
log.error("SSE流出错: {}", error.getMessage());
sink.error(error);
}
});
// 前端断开连接时取消流
sink.onCancel(() -> log.info("前端断开SSE连接, prompt: {}",
prompt.substring(0, Math.min(prompt.length(), 50))));
}, FluxSink.OverflowStrategy.BUFFER)
// 缓冲区上限,防止内存泄漏
.onBackpressureBuffer(256,
token -> log.warn("SSE背压丢弃token, 长度: {}", token.length()))
// 限制每秒最多推送100个token
.limitRate(100)
// 分离IO线程
.publishOn(Schedulers.boundedElastic())
.doOnError(error -> log.error("Flux管道异常", error));
}
}
这个端点生产就绪的关键点:一是 sink.isCancelled 检查——前端断开连接时不再往 sink 里塞数据,避免后台空转;二是 onBackpressureBuffer 设置 256 上限——高并发场景下限制单连接内存占用;三是 limitRate 控制推送速率——避免下游中间代理(Nginx、CDN)的缓冲被撑爆。
四、踩坑经验和最佳实践
第一个坑是忘记检查 sink.isCancelled()。 流式 API 在 WebFlux 里最常见的 Bug 就是:用户点了一下"停止生成",前端断开 SSE 连接,后端还在孜孜不倦地调 onPartialResponse。Token 持续产出,sink.next 往一个没人读的 Flux 里塞数据,如果没有背压限制,内存一直涨。只要有连接断开机制,必须在 onPartialResponse 里先检查 sink.isCancelled()。
第二个坑是在回调线程里直接做同步 IO。 onPartialResponse 的回调线程来自 OkHttp 的 IO 线程池,线程数有限(通常和 CPU 核心数相关)。如果你在这个线程里调了一次数据库查询,延迟 50ms,一秒 20 个请求就能让线程池耗尽。正确做法是用 publishOn 切换到业务线程池。
第三个坑是 BlockingQueue 无界。 如果你自己用 BlockingQueue 或 LinkedBlockingQueue 做缓冲,默认 Integer.MAX_VALUE 的容量等于无界。流量大的时候队列里的 token 串能绕地球一圈,然后 OOM。必须用有界队列或者用 Flux 的 onBackpressureBuffer 设上限。
第四个坑是忽略了模型的 quota 和 rate limit。 OpenAI 等 API 服务对并发流式请求有速率限制(通常是 RPM 和 TPM)。如果你的 Flux 管道做得再完美,但一个用户疯狂刷请求触发了 API provider 的 429 限流,该炸还是炸。需要在调用 StreamingChatModel 之前做请求级别的限流,可以用 Resilience4j 的 RateLimiter 或者 Redis 令牌桶。
第五个坑是错误处理不完整。 SSE 流中途报错(比如 API key 过期、网络断开),onError 会被调用,但 Flux 的 sink.error 之后如果没在 Flux 链路上做错误恢复,WebFlux 会直接返回 500,前端收到的是一个断掉的 SSE 连接而不是一个优雅的错误事件。建议在 Flux 末尾加 onErrorResume 返回一个错误提示 token。
最佳实践总结: 永远用 Flux.create 桥接 StreamingChatModel 和 Reactive Streams;永远设 onBackpressureBuffer 上限;永远在 onPartialResponse 里检查 sink.isCancelled;永远用 publishOn 分离 IO 和业务线程;永远在调用 API 前做并发限流;永远在 Flux 末尾做错误恢复。
五、性能对比和技术选型
做了一组对比测试,场景是 100 并发用户同时请求流式聊天,每个请求平均产生 500 个 token,在 4 核 8G 的服务器上跑了 5 分钟:
方案一:推模型 + 同步处理。每个 onPartialResponse 内做 5ms 模拟业务逻辑(Thread.sleep)。结果:30 秒后 CPU 100%,堆内存从 200MB 涨到 2.8GB,GC 暂停频繁达到 500ms+,15% 的请求因超时断开。OkHttp IO 线程池被耗尽的请求直接报 ReadTimeout。
方案二:TokenStream 拉模型 + 有界队列。BlockingQueue 容量设为 256,消费线程用 CachedThreadPool。结果:CPU 稳定在 60-70%,堆内存稳定在 600MB 左右,没有超时断开,但单请求平均延迟增加了 200ms,因为队列满了之后生产线程阻塞导致 token 产出变慢。
方案三:Flux + 背压控制。onBackpressureBuffer(256) + limitRate(50) + publishOn(boundedElastic)。结果:CPU 稳定在 45-55%,堆内存稳定在 400MB,0 超时断开,单请求端到端延迟和方案二相当。额外优势是 front-end disconnect 后后端自动停止,不会浪费 API 调用费用。
选型建议很简单:推模型只适合 Quick Demo 和单用户测试;拉模型 + 有界队列适合不需要前端 SSE 的批处理场景;生产环境一概用 Flux + 背压控制。如果你的项目不是 WebFlux 而是传统的 Servlet 架构(Spring MVC + Tomcat),可以考虑用 SseEmitter + TokenStream 的方案,但背压控制的效果会差一些,因为 Servlet 的阻塞模型限制了背压信号的传播。
六、总结
流式输出在生产环境中出问题,根因几乎都是同一个——生产者和消费者的速度不匹配。LangChain4j 给了你三把刀:StreamingChatResponseHandler 是最简单但也是最危险的匕首,TokenStream 是能主动控制节奏但缺乏精细流控的短剑,而 TokenStream + Reactor Flux 的组合才是一把能上战场的重剑。
核心原则就三条。第一,永远不要在处理回调里做耗时操作,用 publishOn 把逻辑甩到业务线程池。第二,永远设置缓冲区上限,256 是一个安全的起点,根据你的 token 产出速率调整。第三,永远检查连接状态,用户断开了就别再徒劳往里面塞数据,既浪费内存也浪费 API 费用。
技术选型上别纠结。只要你的项目是基于 WebFlux 的,直接上 Flux.create 方案,五个操作符(onBackpressureBuffer、limitRate、publishOn、doOnError、onErrorResume)配齐就能抗住生产流量。如果你的项目是老牌 Spring MVC,SseEmitter 也能凑合用,但要自己管理好线程池,别让 Tomcat 的 worker 线程被 tokend 的消费逻辑拖垮。
背压控制不是什么高深技术,它就是"消费者说了算"原则的工程化落地。把这句话刻在脑子里,你的流式管道就不会炸。
更多推荐
所有评论(0)