Spring Boot + SSE 流式接口实战
作者:程序员马丁
Ragent AI —— 从 0 到 1 纯手工打造企业级 Agentic RAG,拒绝 Demo 玩具!AI 时代,助你拿个offer。
上一篇大家已经能消费大模型的流式 API 了。SseStreamClient 类封装好了,onToken 回调里能实时拿到每个 token,onComplete 回调里能拿到完整内容和 Token 统计。
但现在有个问题:这些 token 只在你的 Java 进程里。用户在浏览器里什么都看不到。
你需要把后端收到的流式数据实时转发给前端,这个时候 Spring Boot 服务要扮演双重角色——对大模型 API 来说它是 SSE 客户端,对前端浏览器来说它是 SSE 服务端。
这篇就来讲怎么用 Spring Boot 搭建 SSE 服务端,实现前端 → Spring Boot → 大模型 API → 前端的完整流式转发链路,以及生产环境绕不开的那些工程细节。
为什么需要后端中间层
1. 前端直连大模型 API 的三个问题
你可能会想:前端直接调大模型 API 不就行了?何必多一层后端中间层?
不行。有三个硬伤:
-
API Key 暴露:前端代码跑在用户的浏览器里,你的 JavaScript 代码、网络请求对用户是完全透明的。打开 F12 → Network 面板,请求头里的
Authorization: Bearer sk-xxx一目了然。API Key 泄露意味着任何人都能用你的额度调用大模型,账单直接爆炸。 -
跨域限制:大模型 API 的域名(比如
api.siliconflow.cn)和你的前端域名(比如www.yourapp.com)不同,浏览器的同源策略会拦截请求。虽然可以通过 CORS 解决,但 CORS 的响应头需要服务端配置——你控制不了第三方大模型 API 的 CORS 策略。 -
无法加业务逻辑:用户是否登录了?是否有调用权限?单个用户每分钟最多调几次?对话记录要不要存?敏感词要不要过滤?这些业务逻辑都需要在后端处理。前端直连等于跳过了整个业务层,你的系统就是一个没有门禁的大楼。
所以架构必须是这样的:
前端浏览器 → 你的 Spring Boot 后端 → 大模型 API
后端负责鉴权、限流、日志、对话历史存储,然后代理调用大模型 API,把流式响应转发给前端。
2. 完整链路概览
完整的流式转发链路用时序图展示:

有一个关键点需要注意:Controller 线程不会被长时间占用。Controller 方法创建 SseEmitter 对象后立即返回,线程就释放了。真正的数据推送是在另一个异步线程里进行的。这意味着即使流式生成需要 30 秒,Tomcat 的 Controller 线程也只占用了几毫秒。
Spring Boot SseEmitter 核心机制
1. SseEmitter 是什么
SseEmitter 是 Spring MVC 提供的异步响应工具,专门用于 SSE 服务端推送。
普通的 Controller 返回值(比如 ResponseEntity、String、Map)是同步响应——方法执行完,响应就发了,HTTP 连接就关了。一问一答,干脆利落。
SseEmitter 是异步响应——Controller 方法返回一个 SseEmitter 对象后,HTTP 连接不会关闭。你可以在其他线程里通过这个 SseEmitter 对象持续向客户端推送事件,直到你调用 emitter.complete() 主动关闭连接。
打个比方:普通 Controller 像自动售货机——投币、出货、走人。SseEmitter 像外卖平台——你下了单(发了请求),然后骑手持续给你推送状态更新(已接单、正在配送、即将到达),直到送达(complete())。
2. SseEmitter 的核心 API
SseEmitter 的 API 不多,核心就这几个:
| API | 作用 | 说明 |
|---|---|---|
new SseEmitter(timeout) | 创建实例,设置超时 | 超时单位毫秒, 超时后自动关闭连接 |
send(Object data) | 推送默认类型事件 | 事件类型为 message,data 会被序列化为字符串 |
send(SseEventBuilder event) | 推送自定义事件 | 可指定 event 名称、id、data、comment |
complete() | 正常结束 | 关闭 SSE 连接 |
completeWithError(Throwable ex) | 异常结束 | 关闭连接并触发 onError 回调 |
onCompletion(Runnable) | 连接关闭回调 | 不管正常还是异常关闭都会触发 |
onTimeout(Runnable) | 超时回调 | 超时时触发 |
onError(Consumer<Throwable>) | 错误回调 | 出错时触发 |
其中 send(SseEventBuilder) 是最常用的,因为它可以指定事件类型:
// 推送一个 token 事件
emitter.send(SseEmitter.event()
.name("token") // event: token
.data("{\"content\":\"你好\"}")); // data: {"content":"你好"}
// 推送一个 done 事件
emitter.send(SseEmitter.event()
.name("done")
.data("{\"usage\":{\"total_tokens\":128}}"));
// 推送一个注释(心跳保活)
emitter.send(SseEmitter.event()
.comment("keepalive"));
还记得上一篇讲的 SSE event: 字段吗?大模型 API 通常不用 event: 字段,所有事件都是默认的 message 类型。但你自己搭建 SSE 服务端时,event: 字段就很有用了——前端可以根据不同的事件类型做不同的处理。
3. SseEmitter 的线程模型
SseEmitter 的异步本质是理解它的关键。整个过程涉及两个线程:
Controller 线程(Tomcat 线程池里的线程):
- 接收 HTTP 请求
- 创建
SseEmitter对象 - 启动异步任务(把
SseEmitter对象传给异步线程) - 返回
SseEmitter对象 - Controller 线程结束,释放回 Tomcat 线程池
推送线程(你自己创建的异步线程):
- 拿着
SseEmitter对象 - 调用大模型流式 API
- 在回调里调用
emitter.send()逐 Token 推送 - 调用
emitter.complete()结束
这个设计的好处是 Controller 线程不会被长时间占用。如果不用 SseEmitter,而是在 Controller 里直接写一个循环发送数据,那这个请求会占住一个 Tomcat 线程直到流结束——30 秒的流式生成就占住一个线程 30 秒。Tomcat 默认 200 个线程,200 个并发用户就把线程池打满了。
用 SseEmitter,Controller 线程几毫秒就释放了。推送操作在独立的线程中进行,不占 Tomcat 线程池。
如果你在 Controller 线程里直接调
send()发完所有数据再complete(),那和普通接口没什么区别——线程还是被占住了。SseEmitter的价值在于跨线程推送:Controller 线程返回,推送线程接管。
Java 实战:三个层次的 SSE 服务端
1. 最简版:Hello SSE
下述为示例代码,大家可自行创建项目进行运行验证,也可仅关注其中的核心思想。关于基础的 Spring Boot 框架代码逻辑,本文在概念章节中不再做过多展开。
先用 最少的代码跑通一个 SSE 接口,理解 SseEmitter 的基本套路。
@RestController
@RequestMapping("/api/sse")
public class SimpleSseController {
@GetMapping("/hello")
public SseEmitter hello() {
// 创建 SseEmitter,超时 60 秒
SseEmitter emitter = new SseEmitter(60_000L);
// 启动一个新线程做数据推送(不能在 Controller 线程里做)
new Thread(() -> {
try {
for (int i = 1; i <= 5; i++) {
// 推送事件:event 名称为 "message"(默认),data 是字符串
emitter.send(SseEmitter.event()
.name("token")
.data("第 " + i + " 条消息"));
Thread.sleep(1000); // 每秒推送一条
}
// 推送完毕,关闭连接
emitter.send(SseEmitter.event()
.name("done")
.data("{\"msg\":\"推送完毕\"}"));
emitter.complete();
} catch (Exception e) {
emitter.completeWithError(e);
}
}).start();
// Controller 线程立即返回,不等推送完成
return emitter;
}
}
启动 Spring Boot 之后,用浏览器直接访问 http://localhost:8080/api/sse/hello,你会看到数据一条一条地出现——每秒出一条,5 秒后结束。
也可以用下面这段 HTML 在浏览器里测试(后面完整版也用这种方式验证):
<!DOCTYPE html>
<html>
<body>
<h3>SSE 测试</h3>
<div id="output"></div>
<script>
const output = document.getElementById('output');
const eventSource = new EventSource('/api/sse/hello');
eventSource.addEventListener('token', function(e) {
output.innerHTML += '<p>收到 token: ' + e.data + '</p>';
});
eventSource.addEventListener('done', function(e) {
output.innerHTML += '<p style="color:green">完成: ' + e.data + '</p>';
eventSource.close(); // 收到 done 事件后关闭连接
});
eventSource.onerror = function() {
output.innerHTML += '<p style="color:red">连接断开</p>';
eventSource.close();
};
</script>
</body>
</html>
这个最简版能跑通,但离生产可用还有距离——没有接入大模型、没有心跳、没有断连检测。接下来一步步加。
2. 完整版:接入大模型流式转发
完整版要实现的链路是:前端发请求 → Spring Boot 接收 → 调用大模型流式 API → 逐 Token 转发给前端。
事件类型设计:
| 事件名称 | 用途 | data 格式 |
|---|---|---|
token | 内容增量 | {"content": "一段文字"} |
done | 流结束 | {"usage": {"prompt_tokens":15, "completion_tokens":12}} |
error | 错误通知 | {"code": 500, "message": "错误描述"} |
前端根据事件类型做不同处理——收到 token 就拼接内容,收到 done 就结束,收到 error 就展示错误信息。
2.1 请求和响应结构
/** 前端请求体 */
public class ChatRequest {
private String question;
// getter/setter 省略
}
2.2 Controller 层
@RestController
@RequestMapping("/api/chat")
public class ChatController {
private final ChatService chatService;
public ChatController(ChatService chatService) {
this.chatService = chatService;
}
@GetMapping("/stream")
public SseEmitter stream(@RequestParam String question) {
// 创建 SseEmitter,超时 3 分钟
SseEmitter emitter = new SseEmitter(180_000L);
// 设置回调(连接关闭时的清理逻辑)
emitter.onCompletion(() ->
System.out.println("SSE 连接关闭"));
emitter.onTimeout(() ->
System.out.println("SSE 连接超时"));
emitter.onError(e ->
System.err.println("SSE 连接异常: " + e.getMessage()));
// 异步线程中调用大模型并转发
chatService.streamChat(question, emitter);
return emitter;
}
}
这里用
@GetMapping+@RequestParam是为了兼容浏览器原生的EventSourceAPI(它只支持 GET 请求)。如果你的前端用fetch来消费 SSE(后面会讲),可以改成@PostMapping+@RequestBody,支持更复杂的请求体。
2.3 Service 层:流式转发核心逻辑
@Service
public class ChatService {
private static final String API_URL = "https://api.siliconflow.cn/v1/chat/completions";
private static final String API_KEY = "sk-xxx"; // 替换为你的 API Key
private static final String MODEL = "Qwen/Qwen3-32B";
private final OkHttpClient httpClient = new OkHttpClient.Builder()
.connectTimeout(15, TimeUnit.SECONDS)
.readTimeout(60, TimeUnit.SECONDS)
.writeTimeout(15, TimeUnit.SECONDS)
.build();
private final Gson gson = new Gson();
/**
* 异步调用大模型流式 API,逐 Token 通过 SseEmitter 转发给前端
*/
public void streamChat(String question, SseEmitter emitter) {
// 在新线程中执行(生产环境应使用线程池,后面会讲)
new Thread(() -> doStreamChat(question, emitter)).start();
}
private void doStreamChat(String question, SseEmitter emitter) {
// 1. 构建大模型请求体
JsonObject requestBody = new JsonObject();
requestBody.addProperty("model", MODEL);
requestBody.addProperty("temperature", 0.7);
requestBody.addProperty("max_tokens", 2048);
requestBody.addProperty("stream", true);
JsonArray messages = new JsonArray();
JsonObject systemMsg = new JsonObject();
systemMsg.addProperty("role", "system");
systemMsg.addProperty("content", "你是一个技术专家,回答简洁清晰。");
messages.add(systemMsg);
JsonObject userMsg = new JsonObject();
userMsg.addProperty("role", "user");
userMsg.addProperty("content", question);
messages.add(userMsg);
requestBody.add("messages", messages);
// 2. 发起 HTTP 请求
Request request = new Request.Builder()
.url(API_URL)
.addHeader("Authorization", "Bearer " + API_KEY)
.addHeader("Content-Type", "application/json")
.addHeader("Accept", "text/event-stream")
.post(RequestBody.create(requestBody.toString(),
MediaType.parse("application/json")))
.build();
// 3. 解析 SSE 流并转发
try (Response response = httpClient.newCall(request).execute()) {
if (!response.isSuccessful()) {
sendError(emitter, response.code(),
"大模型 API 调用失败: HTTP " + response.code());
return;
}
BufferedReader reader = new BufferedReader(
new InputStreamReader(response.body().byteStream(),
StandardCharsets.UTF_8));
String line;
while ((line = reader.readLine()) != null) {
if (line.isEmpty() || line.startsWith(":")) {
continue;
}
if (!line.startsWith("data:")) {
continue;
}
String data = line.substring(5);
if (data.startsWith(" ")) {
data = data.substring(1);
}
if ("[DONE]".equals(data)) {
// 流结束,发送 done 事件
emitter.send(SseEmitter.event()
.name("done")
.data("{\"msg\":\"completed\"}"));
emitter.complete();
return;
}
// 解析 JSON
JsonObject chunk;
try {
chunk = JsonParser.parseString(data).getAsJsonObject();
} catch (Exception e) {
continue; // JSON 解析失败,跳过
}
// 提取 content
JsonArray choices = chunk.getAsJsonArray("choices");
if (choices == null || choices.isEmpty()) {
// 可能是只有 usage 的 chunk
if (chunk.has("usage") && !chunk.get("usage").isJsonNull()) {
JsonObject usage = chunk.getAsJsonObject("usage");
emitter.send(SseEmitter.event()
.name("done")
.data(gson.toJson(usage)));
emitter.complete();
return;
}
continue;
}
JsonObject choice = choices.get(0).getAsJsonObject();
JsonObject delta = choice.getAsJsonObject("delta");
if (delta != null && delta.has("content")) {
JsonElement contentElement = delta.get("content");
if (!contentElement.isJsonNull()) {
String token = contentElement.getAsString();
if (!token.isEmpty()) {
// 核心:通过 SseEmitter 把 token 推送给前端
emitter.send(SseEmitter.event()
.name("token")
.data("{\"content\":\"" +
escapeJson(token) + "\"}"));
}
}
}
// 检查 finish_reason,提取 usage
JsonElement finishElement = choice.get("finish_reason");
if (finishElement != null && !finishElement.isJsonNull()) {
if (chunk.has("usage") && !chunk.get("usage").isJsonNull()) {
JsonObject usage = chunk.getAsJsonObject("usage");
emitter.send(SseEmitter.event()
.name("done")
.data(gson.toJson(usage)));
}
}
}
// readLine 返回 null 但没收到 [DONE],连接异常
sendError(emitter, 500, "流式响应异常结束");
} catch (IOException e) {
sendError(emitter, 500, "连接异常: " + e.getMessage());
}
}
/** 向前端推送错误事件 */
private void sendError(SseEmitter emitter, int code, String message) {
try {
emitter.send(SseEmitter.event()
.name("error")
.data("{\"code\":" + code + ",\"message\":\"" +
escapeJson(message) + "\"}"));
emitter.complete();
} catch (IOException e) {
emitter.completeWithError(e);
}
}
/** 简单的 JSON 字符串转义 */
private String escapeJson(String s) {
return s.replace("\\", "\\\\")
.replace("\"", "\\\"")
.replace("\n", "\\n")
.replace("\r", "\\r")
.replace("\t", "\\t");
}
}
核心逻辑就是一句话:在大模型 SSE 客户端的回调里,调用 emitter.send() 把数据转发给前端。
整个链路是这样的:
- 前端发请求 → Controller 创建
SseEmitter,启动异步线程,立即返回 - 异步线程调用大模型流式 API → 逐行读取
data:行 - 解析出
delta.content→ 通过emitter.send()推送token事件给前端 - 收到
[DONE]→ 推送done事件 →emitter.complete()关闭连接 - 出错 → 推送
error事件 →emitter.complete()关闭连接
注意
sendError方法的设计——出错时不是直接completeWithError(),而是先通过 SSE 推送一个error事件让前端知道出了什么问题,然后再complete()