Spring Boot 3 SSE 流式输出:AI 对话打字机效果的完整实现(附 Vue 3 案例)
Spring Boot 3 SSE流式输出实现原理
引言
在AI对话场景中,大语言模型(LLM)的响应通常需要数秒甚至数十秒才能完整生成。如果等待全部内容生成完毕再返回,用户体验极差——用户面对长时间空白页面,不知道系统是否在正常工作。
SSE(Server-Sent Events)协议允许服务端向客户端单向推送数据流,是实现AI对话"打字机效果"的核心技术。本文将深入分析SSE协议原理、Spring Boot 3中的实现方案,以及AI对话流式输出的完整架构。
一、SSE协议原理
1.1 SSE vs WebSocket vs 轮询
| 特性 | HTTP轮询 | WebSocket | SSE |
|---|---|---|---|
| 通信方向 | 客户端→服务端 | 双向 | 服务端→客户端 |
| 协议 | HTTP | WS | HTTP |
| 连接开销 | 每次新建连接 | 一次握手 | 一次握手 |
| 断线重连 | 不需要 | 需手动实现 | 浏览器自动 |
| 数据格式 | 自定义 | 自定义 | 文本(text/event-stream) |
| 浏览器支持 | 全部 | 现代浏览器 | 现代浏览器 |
| 适用场景 | 低频查询 | 实时双向通信 | 服务端推送 |
SSE的优势在于:基于标准HTTP协议,无需额外握手,浏览器原生支持断线重连,天然适合服务端向客户端单向推送数据的场景。
1.2 SSE消息格式
SSE遵循text/event-stream MIME类型,消息格式如下:
field: value
支持的字段:
event:事件类型(默认为message)data:数据内容(可多行)id:事件ID(用于断线重连)retry:重连间隔(毫秒)
id: 1
event: message
data: {"content": "你", "index": 0}
id: 2
event: message
data: {"content": "好", "index": 1}
id: 3
event: done
data: [DONE]
二、Spring Boot 3 SSE实现
2.1 SseEmitter核心类
Spring Boot 3提供了SseEmitter类,封装了SSE协议的细节:
sequenceDiagram
participant Client as 客户端
participant Controller as Spring Controller
participant Emitter as SseEmitter
participant Service as AI服务
Client->>Controller: GET /ai/chat/stream
Controller->>Emitter: 创建SseEmitter(timeout)
Controller-->>Client: 返回SseEmitter(Content-Type: text/event-stream)
loop 流式输出
Service->>Emitter: emitter.send(data)
Emitter-->>Client: 推送SSE事件
end
Service->>Emitter: emitter.complete()
Emitter-->>Client: 关闭连接
2.2 基础实现
/**
* AI对话控制器
* 提供SSE流式对话接口
*/
@RestController
@RequestMapping("/ai/chat")
public class AiChatController {
@Autowired
private AiChatService aiChatService;
/** SSE超时时间:5分钟 */
private static final long SSE_TIMEOUT = 5 * 60 * 1000L;
/**
* 流式对话接口
* @param message 用户消息
* @return SseEmitter
*/
@GetMapping(value = "/stream", produces = MediaType.TEXT_EVENT_STREAM_VALUE)
public SseEmitter streamChat(@RequestParam String message) {
// 创建SseEmitter,设置超时时间
SseEmitter emitter = new SseEmitter(SSE_TIMEOUT);
// 注册超时回调
emitter.onTimeout(() -> {
emitter.complete();
});
// 注册错误回调
emitter.onError(ex -> {
emitter.complete();
});
// 异步执行AI对话
CompletableFuture.runAsync(() -> {
try {
aiChatService.streamChat(message, emitter);
emitter.complete(); // 正常完成
} catch (Exception e) {
emitter.completeWithError(e); // 异常完成
}
});
return emitter;
}
}
2.3 AI服务流式调用
/**
* AI对话服务
* 调用大模型API并流式输出
*/
@Service
public class AiChatService {
@Autowired
private LlmClient llmClient;
/**
* 流式对话
* @param message 用户消息
* @param emitter SSE发射器
*/
public void streamChat(String message, SseEmitter emitter) throws IOException {
// 构建AI请求
ChatRequest request = ChatRequest.builder()
.message(message)
.stream(true) // 启用流式输出
.build();
// 调用大模型API,获取流式响应
Flux<ChatChunk> responseFlux = llmClient.streamChat(request);
// 订阅流式响应,逐块推送给客户端
responseFlux.subscribe(
chunk -> {
try {
// 构建SSE事件
SseEmitter.SseEventBuilder event = SseEmitter.event()
.id(String.valueOf(chunk.getIndex()))
.data(chunk.getContent(), MediaType.APPLICATION_JSON);
emitter.send(event);
} catch (IOException e) {
emitter.completeWithError(e);
}
},
error -> {
try {
// 发送错误事件
emitter.send(SseEmitter.event()
.name("error")
.data("AI服务异常"));
} catch (IOException ignored) {
}
emitter.completeWithError(error);
},
() -> {
try {
// 发送完成事件
emitter.send(SseEmitter.event()
.name("done")
.data("[DONE]"));
} catch (IOException ignored) {
}
}
);
}
}
三、前端EventSource对接
3.1 原生EventSource
/**
* 使用原生EventSource接收SSE流
*/
function startSSEChat(message, onMessage, onDone, onError) {
const url = `/api/ai/chat/stream?message=${encodeURIComponent(message)}`;
const eventSource = new EventSource(url);
// 监听消息事件
eventSource.addEventListener('message', (event) => {
const data = JSON.parse(event.data);
onMessage(data.content);
});
// 监听完成事件
eventSource.addEventListener('done', (event) => {
eventSource.close();
onDone();
});
// 监听错误事件
eventSource.addEventListener('error', (event) => {
eventSource.close();
onError(event);
});
// 连接错误处理
eventSource.onerror = () => {
eventSource.close();
onError(new Error('SSE连接异常'));
};
return eventSource;
}
3.2 Vue 3组合式封装
/**
* AI对话SSE组合式函数
* 封装SSE连接、消息接收、状态管理
*/
import { ref } from 'vue';
export function useAiChat() {
const messages = ref<ChatMessage[]>([]);
const isStreaming = ref(false);
const currentContent = ref('');
let eventSource: EventSource | null = null;
/**
* 发送消息并接收流式响应
*/
function sendMessage(content: string) {
if (isStreaming.value) return;
// 添加用户消息
messages.value.push({ role: 'user', content });
isStreaming.value = true;
currentContent.value = '';
const url = `/api/ai/chat/stream?message=${encodeURIComponent(content)}`;
eventSource = new EventSource(url);
eventSource.addEventListener('message', (event) => {
const data = JSON.parse(event.data);
currentContent.value += data.content;
});
eventSource.addEventListener('done', () => {
// 流式输出完成,将完整AI消息加入列表
messages.value.push({
role: 'assistant',
content: currentContent.value
});
currentContent.value = '';
isStreaming.value = false;
eventSource?.close();
});
eventSource.onerror = () => {
isStreaming.value = false;
eventSource?.close();
};
}
/**
* 停止生成
*/
function stopGeneration() {
eventSource?.close();
if (currentContent.value) {
messages.value.push({
role: 'assistant',
content: currentContent.value
});
}
currentContent.value = '';
isStreaming.value = false;
}
return { messages, isStreaming, currentContent, sendMessage, stopGeneration };
}
四、进阶:支持POST请求的SSE
原生EventSource只支持GET请求,无法携带请求体。在实际项目中,对话请求通常需要携带复杂的参数(会话ID、上下文、模型配置等),需要使用fetch API实现POST方式的SSE。
4.1 服务端:基于ResponseBodyEmitter
/**
* 支持POST请求的流式对话接口
*/
@PostMapping(value = "/stream/post", produces = MediaType.TEXT_EVENT_STREAM_VALUE)
public SseEmitter streamChatPost(@RequestBody ChatRequestDTO request) {
SseEmitter emitter = new SseEmitter(SSE_TIMEOUT);
emitter.onTimeout(emitter::complete);
emitter.onError(ex -> emitter.complete());
CompletableFuture.runAsync(() -> {
try {
aiChatService.streamChat(request.getMessage(), emitter);
emitter.complete();
} catch (Exception e) {
emitter.completeWithError(e);
}
});
return emitter;
}
4.2 前端:fetch + ReadableStream
/**
* 使用fetch API实现POST方式的SSE接收
* 支持携带请求体和自定义Header
*/
async function fetchSSE(
url: string,
body: object,
onMessage: (data: string) => void,
onDone: () => void,
onError: (error: Error) => void
) {
const response = await fetch(url, {
method: 'POST',
headers: {
'Content-Type': 'application/json',
'Authorization': `Bearer ${getToken()}`
},
body: JSON.stringify(body)
});
if (!response.ok) {
onError(new Error(`HTTP ${response.status}`));
return;
}
const reader = response.body!.getReader();
const decoder = new TextDecoder();
let buffer = '';
while (true) {
const { done, value } = await reader.read();
if (done) break;
buffer += decoder.decode(value, { stream: true });
// 按SSE格式解析(双换行分隔事件)
const events = buffer.split('
');
buffer = events.pop() || ''; // 保留未完成的部分
for (const event of events) {
const lines = event.split('
');
let eventType = 'message';
let data = '';
for (const line of lines) {
if (line.startsWith('event:')) {
eventType = line.substring(6).trim();
} else if (line.startsWith('data:')) {
data = line.substring(5).trim();
}
}
if (eventType === 'done') {
onDone();
return;
}
if (eventType === 'error') {
onError(new Error(data));
return;
}
if (data) {
onMessage(data);
}
}
}
onDone();
}
五、完整时序图
sequenceDiagram
actor User as 用户
participant FE as 前端Vue
participant API as 后端API
participant Emitter as SseEmitter
participant LLM as 大模型API
User->>FE: 输入消息并发送
FE->>API: POST /ai/chat/stream (SSE请求)
API->>Emitter: 创建SseEmitter(5min)
API-->>FE: 返回SseEmitter (text/event-stream)
API->>LLM: 发起流式请求 (stream=true)
loop 流式响应
LLM-->>API: ChatChunk: "你"
API->>Emitter: send(event: message, data: "你")
Emitter-->>FE: SSE: data: "你"
FE->>FE: 追加显示"你"
LLM-->>API: ChatChunk: "好"
API->>Emitter: send(event: message, data: "好")
Emitter-->>FE: SSE: data: "好"
FE->>FE: 追加显示"好"
LLM-->>API: ChatChunk: "!"
API->>Emitter: send(event: message, data: "!")
Emitter-->>FE: SSE: data: "!"
FE->>FE: 追加显示"!"
end
LLM-->>API: 流结束
API->>Emitter: send(event: done, data: [DONE])
Emitter-->>FE: SSE: event: done
API->>Emitter: complete()
FE->>FE: 关闭EventSource
FE->>User: 显示完整AI回复
六、关键问题与解决方案
6.1 连接超时
SSE连接可能因网络问题或代理超时断开。解决方案:
// 设置合理的超时时间
SseEmitter emitter = new SseEmitter(5 * 60 * 1000L);
// 定期发送心跳包保持连接
ScheduledExecutorService scheduler = Executors.newSingleThreadScheduledExecutor();
scheduler.scheduleAtFixedRate(() -> {
try {
emitter.send(SseEmitter.event().name("heartbeat").data(""));
} catch (IOException e) {
scheduler.shutdown();
}
}, 0, 30, TimeUnit.SECONDS);
6.2 背压(Backpressure)控制
当客户端处理速度慢于服务端推送速度时,需要背压控制:
// 使用ResponseBodyEmitter的alreadyComplete状态判断
if (emitter instanceof SseEmitter) {
try {
emitter.send(data);
} catch (IllegalStateException e) {
// 连接已关闭,停止发送
return;
}
}
6.3 Nginx代理配置
Nginx默认会缓冲代理响应,导致SSE消息无法实时推送:
# Nginx SSE代理配置
location /ai/chat/ {
proxy_pass http://backend;
proxy_http_version 1.1;
# 禁用响应缓冲,确保SSE实时推送
proxy_buffering off;
proxy_cache off;
# 关闭连接超时
proxy_read_timeout 300s;
proxy_send_timeout 300s;
# 不使用chunked编码
proxy_set_header Connection '';
}
结论与建议
核心要点
- SSE是AI对话流式输出的最佳选择:基于HTTP、浏览器原生支持、自动重连
- Spring Boot 3的SseEmitter封装了SSE协议细节,使用简单
- POST方式SSE需要前端使用fetch + ReadableStream,但功能更强大
实践建议
- 超时设置:SSE超时时间应大于大模型最长响应时间,建议5分钟以上
- 心跳机制:长时间无数据推送时发送心跳包,防止中间代理关闭连接
- 错误处理:前端需要处理网络断开、服务端异常等场景,提供重试机制
- Nginx配置:必须关闭proxy_buffering,否则SSE消息会被缓冲
- 资源释放:SSE连接关闭时,务必取消大模型API的流式订阅,释放资源