从做AI应用开始我几乎天天和SSE打交道。Java生态里聊SSE的人不少但大多数人用在AI流式对话上还是靠复制粘贴从服务端SseEmitter到客户端逐行读流中间还有一堆连接超时、报文截断、线程打满的坑。这篇我把一整条链路拆开讲清楚先理解SSE在AI场景里到底承担什么角色再讲怎么从显式调用的笨拙写法过渡到隐式封装最后落到虚拟线程——为什么它能让我们在Java里轻松扛住成千上万条SSE长连接。这篇内容适合两类人一类是刚接触AI应用开发、想搞懂流式响应原理的Java后端另一类是已经在写SSE但总觉得代码难维护、并发一上来就出问题的同学。我会把协议细节、代码示例、踩坑记录都放进来照着做基本能少走两三个月的弯路。1. 为什么AI服务偏爱SSE从“轮询”到“长连接”的必然进化1.1 AI流式输出的本质需求不是实时而是“打字机”先想一个问题大模型生成一段回答通常要几秒到几十秒如果等服务端把完整内容算完再一次性返回用户看到的就是一个漫长的转圈等待。而流式输出的好处是模型每生成一个词、一个标点就立刻推给前端屏幕上出现打字机效果。用户的第一体验是“它已经在回答了”等待焦虑被拆成了持续的反馈。这个场景有两个硬性要求一是数据必须增量到达不能等全量二是连接要能长期保持因为生成一句几百字的回答可能需要几十秒。这不叫“实时通信”严格说叫“服务器主动按需推送”。在HTTP生态里能做服务器推送的候选主要有三个WebSocket、轮询、SSE。轮询的问题最明显客户端每隔一秒问一次“有结果了吗”连接频繁建立销毁空请求占满服务端资源。WebSocket倒是双向的但它是一个独立的TCP协议需要协议升级握手服务端要额外维护状态对纯单向的AI流式输出来说是大材小用。SSE走的是普通HTTP客户端用EventSource或普通HTTP客户端就能连服务端只要设置Content-Type: text/event-stream然后一直往响应体里写数据就行。协议简单代理穿透性也远比WebSocket好。1.2 text/event-stream协议细节别被“长连接”三个字骗了SSE全称Server-Sent Events本质是HTTP响应体被拆成了一连串事件。服务端维持一个不关闭的HTTP响应每有新数据就按照特定格式写几行文本客户端按行解析。所谓长连接指的不是HTTP/2多路复用那种连接而是服务端故意不结束响应。协议格式我举个例子HTTP/1.1 200 OK Content-Type: text/event-stream Cache-Control: no-cache Connection: keep-alive data: {content:你好} data: {content:世界} data: [DONE]这里有几个关键字需要理解data:前缀的每一行代表一条数据如果一条事件由多行data:组成它们会被合并成一条用换行分隔的事件空行代表事件结束id:字段用于客户端断线重连时提交Last-Event-IDretry:是客户端自动重连的间隔毫秒数以冒号开头的行是注释服务端可以定期发注释行当心跳避免代理或防火墙把闲置连接关掉。实际开发里最容易踩的坑是把SSE当成WebSocket来用去做双向收发。SSE是单向的客户端想传数据只能另发HTTP请求。但在AI聊天场景这根本不算问题——用户发送一句话走普通POST模型生成内容走SSE单向推流。方向完全匹配。2. 显式调用手写SSE的原始阶段痛在哪里2.1 服务端实现从原始Servlet到Spring的SseEmitter如果完全不借助框架服务端手写SSE用的是Servlet异步输出。核心思路是拿到HttpServletResponse设置响应头然后在异步线程里不断写响应体WebServlet(/stream) public class SseServlet extends HttpServlet { Override protected void doGet(HttpServletRequest req, HttpServletResponse resp) throws IOException { resp.setContentType(text/event-stream); resp.setCharacterEncoding(UTF-8); resp.setHeader(Cache-Control, no-cache); resp.setHeader(Connection, keep-alive); for (int i 0; i 100; i) { resp.getWriter().write(data: {\index\: i }\n\n); resp.getWriter().flush(); try { Thread.sleep(1000); } catch (InterruptedException e) { Thread.currentThread().interrupt(); break; } } } }这种写法在并发低时没有问题但它有一个致命缺陷如果每个SSE请求都占一个Servlet线程而生成内容很慢线程就一直被挂着。Servlet 3.1的异步上下文可以把业务执行放到别的线程池让容器线程先回收本质上还是需要一个执行线程在跑。Spring的SseEmitter就是在这个基础上的封装GetMapping(value /api/ai/chat, produces MediaType.TEXT_EVENT_STREAM_VALUE) public SseEmitter chat(RequestParam String prompt) { SseEmitter emitter new SseEmitter(0L); CompletableFuture .runAsync(() - streamFromModel(prompt, emitter)) .whenComplete((r, e) - { if (e ! null) { emitter.completeWithError(e); } else { emitter.complete(); } }); return emitter; }new SseEmitter(0L)里的0表示不自动超时否则默认30秒连接就会被收掉。这里有一个新手常犯的错在Controller方法里用Thread.sleep模拟流式输出结果发现Tomcat线程全被sleep占满。正确的做法是把耗时任务丢到独立线程池或者响应式订阅里让Controller方法马上返回SseEmitter对象。2.2 客户端手写逐行读流的原始滋味服务端写完之后客户端如果要手写最直接的方式是用Java 11的HttpClient获取InputStream然后BufferedReader逐行读HttpClient client HttpClient.newBuilder() .connectTimeout(Duration.ofSeconds(10)) .build(); HttpRequest request HttpRequest.newBuilder() .uri(URI.create(http://localhost:8080/api/ai/chat?prompt你好)) .GET() .build(); HttpResponseInputStream response client.send( request, HttpResponse.BodyHandlers.ofInputStream() ); try (BufferedReader reader new BufferedReader( new InputStreamReader(response.body(), StandardCharsets.UTF_8))) { String line; while ((line reader.readLine()) ! null) { if (line.startsWith(data: )) { String data line.substring(6); if ([DONE].equals(data.trim())) { break; } System.out.println(收到增量: data); } } }这段代码能把一个SSE流完整读完但它的问题一眼就能看出来。逻辑全堆在while循环里数据格式解析、连接断开处理、重连机制完全靠手写而且每来一个流就要写一遍。如果项目中同时有对话流、通知推送流、日志流三个SSE场景这段代码至少要复制三份改一处要同步改三处。2.3 显式阶段绕不开的三座大山第一座是超时控制。SseEmitter默认超时、Nginx的proxy_read_timeout默认60秒、客户端HTTP连接超时设置三个超时叠加在一起哪个先到都会掐断连接。AI生成一句话动辄几十秒如果哪个中间层的超时漏配用户看到的就是“刚说两个字就断了”。第二座是线程占用。一个SSE连接从建立到结束服务端至少有一个线程在持续输出数据。传统Tomcat默认最大线程200意味着最多同时维持200条流式对话。这在个人项目里无所谓但一旦做产品化就立刻撞墙。第三座是解析逻辑混乱。现实中AI平台的SSE返回并不只是一个data:字段还穿插着注释行、空行、event:类型以及像[DONE]这样的结束标记。手动解析一次没问题每次接到新需求都要重新解析一遍就会出事。3. 隐式封装把SSE从“网络细节”变成“业务方法”3.1 封装设计的核心目标调用方不该知道“流”的存在我做了两三个流式对接项目之后得出一个结论对上层业务来说SSE不应该以“我看得见流”的方式存在。业务方想要的只有一个东西给我一句话我把结果增量地回调给你。连接怎么建、怎么重试、怎么解析、怎么结束这些都应该被封装到黑盒里。所以封装层的接口长这样public interface StreamChatClient { void chat(String prompt, ChatCallback callback); interface ChatCallback { void onDelta(String token); void onComplete(String fullText); void onError(Throwable t); } }先不用Reactor那种Flux接口因为很多同事还不熟响应式用简单的回调反而更容易让队内推广。后续如果需要对接到WebFlux或前端EventSource只要再加一个把回调转成Flux的适配器。3.2 一个可以直接抄的SSE流式客户端封装我用OkHttp的sse扩展来实现因为它的EventSourceListener已经处理了连接生命周期比手写HttpClient.readLine省事得多。依赖先加上dependency groupIdcom.squareup.okhttp3/groupId artifactIdokhttp-sse/artifactId version4.12.0/version /dependency核心实现类如下public class OkHttpStreamChatClient implements StreamChatClient { private final OkHttpClient httpClient; public OkHttpStreamChatClient() { this.httpClient new OkHttpClient.Builder() .connectTimeout(Duration.ofSeconds(10)) .readTimeout(Duration.ZERO) .build(); } Override public void chat(String prompt, ChatCallback callback) { Request request new Request.Builder() .url(https://your-ai-endpoint/api/chat) .post(RequestBody.create( ({\prompt\:\ prompt \}).getBytes(StandardCharsets.UTF_8), MediaType.parse(application/json; charsetutf-8))) .build(); EventSource.Factory factory EventSources.createFactory(httpClient); factory.newEventSource(request, new EventSourceListener() { private final StringBuilder fullText new StringBuilder(); Override public void onEvent(EventSource eventSource, String id, String type, String data) { if ([DONE].equals(data.trim())) { callback.onComplete(fullText.toString()); return; } String token parseDelta(data); if (token ! null) { fullText.append(token); callback.onDelta(token); } } Override public void onFailure(EventSource eventSource, Throwable t, Response response) { callback.onError(t); } }); } private String parseDelta(String data) { // 根据对接的大模型返回结构解析增量内容字段 // 例如 OpenAI 风格{choices:[{delta:{content:...}}]} JsonNode root JsonParser.readTree(data); JsonNode content root.at(/choices/0/delta/content); return content.isMissingNode() ? null : content.asText(); } }这里我把readTimeout设成Duration.ZERO目的是禁止OkHttp的读超时自动掐断流。SSE流的特点是两次data之间可能隔很久传统读超时会把长思考过程当成无响应而断开。这也是显式调用阶段最容易出的问题。它的onEvent回调已经帮我们处理了SSE的分行和事件边界我们不关心协议细节只关心业务逻辑。如果有断线重连的需求可以在onFailure里判断错误类型并配合请求头里的Last-Event-ID重新发起连接。这就是隐式封装的意义把一类偶发、需要反复处理的边界情况收口到一个地方。3.3 增量JSON解析半包数据才是真正的难点对接不同大模型平台时有一个细节会让人抓狂有些服务端不是按事件边界输出完整JSON而是把一个JSON对象拆成多个data:事件依次推送。比如第一次data:里的内容是{choices:[{delta:{content:你}}]}第二次就变成了{choices:[{delta:{content:好}}]}这还算好的。糟糕的情况是一次事件里出现两个不完整的嵌套JSON片段或者一个事件本身是半个字符串。处理这个问题的标准做法是引入增量JSON解析而不是每次收到事件后重新readTree。Jackson的JsonParser支持流式读取即使字节流还没完整也可以边读边解析public class IncrementalJsonParser { private final JsonFactory factory new JsonFactory(); private final StringBuilder buffer new StringBuilder(); public OptionalString parseNextToken(String chunk) { buffer.append(chunk); try (JsonParser parser factory.createParser(buffer.toString())) { while (parser.nextToken() ! null) { if (parser.currentToken() JsonToken.FIELD_NAME content.equals(parser.getCurrentName())) { parser.nextToken(); return Optional.of(parser.getText()); } } } catch (JsonParseException ignored) { // 说明当前buffer里的JSON还不完整等下一个chunk拼起来再解析 } return Optional.empty(); } }这段代码的做法是把新收到的chunk追加到buffer里然后从头解析。如果解析成功就返回content字段如果抛JsonParseException说明JSON不完整保持buffer等下一个chunk。在token级流式输出场景里content字段通常在一个chunk里就能拿到不会出现真正需要跨chunk拼字段的情况但如果你对接的是自定义协议这套增量解析兜底就很有必要。封装层的完整链路线是底层OkHttp负责连接和事件边界IncrementalJsonParser负责把不规范的数据变成干净的token流上层回调把所有token拼成完整文本。调用方哪里还看得到SSE的影子它只知道“传了prompt回调被调用了几十次”。4. 虚拟线程把“线程不够用”彻底变成一个伪命题4.1 长连接场景下的“线程饿死”问题比想象中更早到来传统Java服务端处理并发请求靠的是线程池。Tomcat默认最大线程200也就是说同一时刻最多只有200个请求在执行。如果200个请求全都是SSE长连接每个连接都在等待模型生成数据那么第201个普通请求就会被拒绝或者排队。我在一个内部工具项目里真实的体验是白天团队十来个人同时用AI问答每个页面开了两三条流连接数也就三四十一切正常。但到了下午测试自动化脚本并发跑任务时每条任务都开SSE连接连接数瞬间冲到150紧接着新请求就开始报Connection refused。当时的排查结论不是流量大而是Tomcat线程被SSE连接“睡”满了。为什么说是“睡”满因为每个SSE连接在等待数据时线程是阻塞在read或者write上的。它不干活但占着位置。传统的IO密集场景里线程大多数时间都在休眠系统依然不自由——线程的创建和栈空间本身就吃内存数量上不去。4.2 虚拟线程把“线程”变成廉价消耗品Java 21正式发布的虚拟线程Virtual Threads本质上把“线程”从系统资源变成了一个JVM管理的对象。虚拟线程不绑定操作系统线程它在阻塞时会被调度器挂起让出载体线程当数据可达时再重新挂到某个载体线程上继续执行。对开发者来说代码写法和平台线程一模一样不需要改Thread、ExecutorService的调用方式只需要换一个Executor。真正重要的区别是平台线程数量受限于内核和内存虚拟线程则可以以百万计地创建。这意味着SSE长连接完全不需要复用线程了每条连接独自占一个虚拟线程线程睡多久都是零成本因为睡的时候不占OS线程。一个直观的对比// 传统方式 ExecutorService executor Executors.newFixedThreadPool(200); // 虚拟线程方式 ExecutorService executor Executors.newVirtualThreadPerTaskExecutor();就改了这一行系统能同时维持的长连接数量从几百上升到数千甚至更多。我在一次压测里用原来会崩溃的并发场景去跑虚拟线程下连接全都能建立每连一个占一个虚拟线程内存增长几乎可以忽略。4.3 Spring Boot中启用虚拟线程的两种姿势Spring Boot 3.2以上的项目最省事的配置是在application.yml里加一行spring.threads.virtual.enabled: true这会让Tomcat的请求处理线程、Async异步任务的线程池都切换到虚拟线程。做这一件事之前你要确认项目里没有处在关键路径上的ThreadLocal依赖——虚拟线程的ThreadLocal虽然能用但成千上万个虚拟线程每个都带一份ThreadLocal副本内存会快速膨胀。如果想更精细地控制在自己的代码里指定虚拟线程执行器ExecutorService executor Executors.newVirtualThreadPerTaskExecutor(); GetMapping(value /api/ai/chat, produces MediaType.TEXT_EVENT_STREAM_VALUE) public SseEmitter chatWithVirtualThread(RequestParam String prompt) { SseEmitter emitter new SseEmitter(0L); executor.submit(() - { try { // 模拟模型流式输出 for (int i 0; i 50; i) { emitter.send(SseEmitter.event().data(chunk- i)); Thread.sleep(100); } emitter.complete(); } catch (Exception e) { emitter.completeWithError(e); } }); return emitter; }这种方式的好处是业务线程被限制在虚拟线程执行器里Tomcat容器线程只负责接收请求和返回SseEmitter对象整个请求链路都不会有占用平台线程的阻塞点。4.4 虚拟线程的隐藏坑pinning现象和池化执念虚拟线程并非完全没有上限。有一个经典问题叫“pinning”当虚拟线程进入synchronized代码块或调用某些原生方法时它会被钉在当前的载体线程上JVM无法再调度它让出。如果高流量的SSE推送代码里有synchronized保护共享资源就可能导致载体线程被占满虚拟线程的优势瞬间消失。排查方法是观察日志或JFRJava Flight Recorder里有没有Virtual Thread Pinned事件。一旦发现用ReentrantLock替换关键路径上的synchronized即可。另一个反直觉的坑是虚拟线程压根不需要线程池复用它本来就是轻量到可以不停创建和销毁的。我一个朋友把newVirtualThreadPerTaskExecutor()替换成newFixedThreadPool(2000)的虚拟线程版本然后还去调池大小、调队列长度完全是多此一举。虚拟线程的正确用法就是“每任务一线程”让操作系统和JVM去调度。还有一个环境层面的前提必须强调虚拟线程是Java 21的正式特性如果项目还在Java 8、Java 11那就没办法用。升级JDK是启动这个方案的第一步也是很多公司最卡壳的一步。如果确实升不了就只能沿用传统线程池调整参数但性能天花板就是换汤不换药。5. 实战问题排查SSE和虚拟线程最容易翻车的地方5.1 连接中断、数据不全、超时问题的对照排查表我在对接外部大模型平台、自建服务端、前端联调三个环节里收集了一些出现频率最高的异常现象。整理成一张表排查速度能快很多。现象可能原因处理办法连接十几秒后断开没有任何报错服务端SseEmitter默认超时构造时传0L或足够大的超时值数据到达不连续经常卡住中间Nginx的proxy_read_timeout过短调大proxy_read_timeout到300s以上流式效果变成了“一次性全部出现”Nginx开启了响应缓冲加proxy_buffering off;中文全部乱码服务端报文编码不是UTF-8响应头设置Content-Type: text/event-stream; charsetUTF-8客户端永远收不到最后[DONE]服务端没有主动complete调用emitter.complete()或者关闭EventSource大量请求报线程池满Tomcat默认200线程被阻塞占用启用spring.threads.virtual.enabled或调大传统线程池虚拟线程模式下偶发超高延迟存在synchronized导致的pinning用JFR确认后替换为ReentrantLock连接断线重连后重复收到旧消息未处理Last-Event-ID客户端带上Last-Event-ID服务端根据ID断点续传这里最隐蔽的是Nginx缓冲区问题。如果不关掉proxy_bufferingNginx会把后端吐出来的SSE数据攒进缓冲区直到攒满或者超时才一次性发给浏览器。表面上接口调通了但用户看不到打字机效果反而以为页面卡死了。只要在location里加上location /api/ai/ { proxy_buffering off; proxy_cache off; proxy_set_header Connection ; proxy_http_version 1.1; proxy_read_timeout 300s; }基本就能解决转发层的问题。proxy_set_header Connection 是为了让上游连接保持长连接配合SSE的最基本要求。5.2 虚拟线程配合SSE时的三个建议第一不要过度使用ThreadLocal。虚拟线程的数量可以很大每个虚拟线程带着独立的ThreadLocal如果存的是大对象GC压力会上升。能用方法参数传递的数据尽量不要塞到ThreadLocal里。第二关掉所有无意义的“连接池”。SSE是长连接每个请求独占一个虚拟线程不需要像普通HTTP那样拼命复用连接。把OkHttpClient的connectionPool调小甚至忽略它因为占用的连接本来就该长期保持。第三留意监控指标的变化。启用虚拟线程后原来“线程数上升”这个告警可能不再敏感需要额外监控“虚拟线程创建速率”和“载体线程利用率”。如果载体线程利用率一直很高说明调度或者pinning有问题而不是业务真的这么吃CPU。6. 最后说一点个人体会从我最早用Servlet写SSE逐行flush到现在用Spring Boot 虚拟线程 OkHttp封装一条完整的AI流式链路最大的感受是技术选型不是越新越好而是越能消灭“重复处理细节”越好。SSE的协议本身很简单难的是它和线程模型、超时策略、代理配置搅在一起之后呈现出来的那些诡异现象。遇到SSE相关问题别急着改代码先画一条从浏览器到模型服务的完整数据链路逐个检查每一层的超时时间、缓冲开关、线程模型。大多数坑不是代码逻辑错误而是某一层的“自动驾驶”设置不符合流式场景。把封装层做好把虚拟线程打开你会发现所谓的SSE并发压力很大一部分根本是纸老虎。我现在做AI后端接口凡是涉及模型流式输出的地方一律走封装好的StreamChatClient连接细节不外泄凡是新起的服务直接Java 21 虚拟线程。这套组合跑了大半年性能和可维护性都超出预期至少不用再半夜被“连接断开”告警叫醒了。