ShotAI SSE 聊天流实现详解
📚 概述
本文深入讲解 ShotAI 项目中如何使用 SSE(Server-Sent Events)技术实现实时流式对话,打造类似 ChatGPT 的打字机效果。
🤔 为什么选择 SSE?
传统轮询 vs SSE vs WebSocket
| 特性 |
传统轮询 |
SSE |
WebSocket |
| 实时性 |
❌ 差 |
✅ 好 |
✅ 好 |
| 服务器推送 |
❌ 需要轮询 |
✅ 原生支持 |
✅ 原生支持 |
| 双向通信 |
❌ 否 |
❌ 否 |
✅ 是 |
| 实现复杂度 |
⭐ 简单 |
⭐⭐ 中等 |
⭐⭐⭐ 复杂 |
| 浏览器支持 |
✅ 全部 |
✅ 现代浏览器 |
✅ 现代浏览器 |
| 自动重连 |
❌ 需手动实现 |
✅ 浏览器自动 |
❌ 需手动实现 |
选择 SSE 的理由:
- ✅ 单向推送满足需求(服务器 → 客户端)
- ✅ 实现简单,基于 HTTP 协议
- ✅ 浏览器自动重连,无需额外处理
- ✅ 完美适配流式 AI 响应场景
🏗️ 架构设计
整体流程图
┌──────────┐ ┌──────────┐ │ 前端 │ │ 后端 │ └────┬─────┘ └────┬─────┘ │ │ │ 1. 建立 SSE 连接 │ ├──────────────────────────────>│ │ GET /api/sse/connect?userId │ │ │ │ 2. 返回 SseEmitter │ │<──────────────────────────────┤ │ (连接保持打开) │ │ │ │ 3. 发送聊天消息 │ ├──────────────────────────────>│ │ POST /api/chat/send │ │ │ │ 4. 立即返回 200 │ │<──────────────────────────────┤ │ │ │ │ 5. 异步处理 │ │ 调用 AI │ │ │ 6. 推送文本片段 (event: add) │ │<──────────────────────────────┤ │ data: "你" │ │ │ │ 7. 继续推送 │ │<──────────────────────────────┤ │ data: "好" │ │ │ │ 8. 推送完成 (event: finish) │ │<──────────────────────────────┤ │ data: "done" │ │ │ │ 9. 用户点击停止 │ ├──────────────────────────────>│ │ POST /api/sse/stop │ │ │ │ 10. 停止生成 │ │<──────────────────────────────┤ │ data: "stopped" │ │ │
|
💻 后端实现
1. SSEMsgType 枚举
public enum SSEMsgType { ADD("add"), FINISH("finish"), ERROR("error"); private final String type; SSEMsgType(String type) { this.type = type; } public String getType() { return type; } }
|
设计说明:
ADD:每次推送一小段文本,前端逐字显示
FINISH:标记生成完成,前端停止接收
ERROR:通知前端发生错误
2. SSEServer 工具类
2.1 核心数据结构
@Slf4j public class SSEServer {
private static final Map<String, SseEmitter> SSE_CACHE = new ConcurrentHashMap<>(); }
|
为什么使用 ConcurrentHashMap?
- 线程安全,支持高并发
- 无需额外加锁,性能优秀
- 适合多用户同时连接的场景
2.2 创建 SSE 连接
public static SseEmitter createSse(String userId) { if (SSE_CACHE.containsKey(userId)) { return SSE_CACHE.get(userId); } SseEmitter sseEmitter = new SseEmitter(1800_000L); sseEmitter.onCompletion(() -> { log.info("【SSE完成】用户: {}", userId); }); sseEmitter.onTimeout(() -> { log.warn("【SSE超时】用户: {}", userId); SSE_CACHE.remove(userId); }); sseEmitter.onError(throwable -> { log.error("【SSE错误】用户: {}, 错误: {}", userId, throwable.getMessage()); SSE_CACHE.remove(userId); }); SSE_CACHE.put(userId, sseEmitter); log.info("【SSE连接】用户: {}, 当前连接数: {}", userId, SSE_CACHE.size()); try { sseEmitter.send(SseEmitter.event() .name("connected") .data("SSE连接已建立")); } catch (IOException e) { log.error("【SSE初始化失败】用户: {}, 错误: {}", userId, e.getMessage()); SSE_CACHE.remove(userId); throw new RuntimeException("SSE连接初始化失败", e); } return sseEmitter; }
|
关键设计点:
- 超时设置:30 分钟,平衡资源占用和用户体验
- 初始化消息:立即发送一条消息,触发前端
onopen 事件
- 回调处理:完善的生命周期管理
- 幂等性:重复调用返回同一个连接
2.3 发送消息
public static void sendMsg(String userId, String content, SSEMsgType msgType) { SseEmitter sseEmitter = SSE_CACHE.get(userId); if (sseEmitter == null) { log.warn("【SSE不存在】用户: {}", userId); return; } try { sseEmitter.send(SseEmitter.event() .data(content) .name(msgType.getType())); } catch (IOException e) { log.error("【SSE发送失败】用户: {}, 错误: {}", userId, e.getMessage()); removeSse(userId); } }
|
消息格式:
event: add data: 你好
event: add data: ,我是
event: finish data: done
|
2.4 连接管理
public static boolean exists(String userId) { return SSE_CACHE.containsKey(userId); }
public static void removeSse(String userId) { SseEmitter sseEmitter = SSE_CACHE.remove(userId); if (sseEmitter != null) { try { sseEmitter.complete(); } catch (Exception e) { log.error("【用户: {}】关闭 SSE 连接失败: {}", userId, e.getMessage()); } } }
public static int getConnectionCount() { return SSE_CACHE.size(); }
|
3. SseController 控制器
@Slf4j @RestController @RequestMapping("/api/sse") public class SseController { @Autowired private ChatService chatService;
@GetMapping("/connect") public SseEmitter connect(@RequestParam String userId) { log.info("建立 SSE 连接,用户: {}", userId); return chatService.createSseConnection(userId); }
@PostMapping("/stop") public ChatResponse stopGeneration(@RequestParam String userId) { log.info("收到停止生成请求,用户: {}", userId); try { chatService.stopGeneration(userId); return ChatResponse.success("已停止生成"); } catch (Exception e) { log.error("停止生成失败", e); return ChatResponse.error("停止失败: " + e.getMessage()); } } }
|
4. 与 AI 流的集成
private void processChat(String userId, String message) { try { Prompt prompt = new Prompt(message); Flux<String> stream = chatClient.prompt(prompt).stream().content(); Disposable subscription = stream .subscribe( content -> { if (SSEServer.exists(userId)) { SSEServer.sendMsg(userId, content, SSEMsgType.ADD); } }, error -> { log.error("【对话流失败】用户: {}", userId, error); SSEServer.sendMsg(userId, "抱歉,出现了错误", SSEMsgType.ERROR); userSubscriptions.remove(userId); }, () -> { log.info("【对话完成】用户: {}", userId); SSEServer.sendMsg(userId, "done", SSEMsgType.FINISH); SSEServer.close(userId); userSubscriptions.remove(userId); } ); userSubscriptions.put(userId, subscription); } catch (Exception e) { log.error("【对话异常】用户: {}", userId, e); SSEServer.sendMsg(userId, "系统错误", SSEMsgType.ERROR); } }
|
数据流转:
DeepSeek API → Flux<String> → subscribe → SSEServer.sendMsg → 前端
|
🌐 前端实现
1. 建立 SSE 连接
const connectSSE = (userId) => { const eventSource = new EventSource(`/api/sse/connect?userId=${userId}`); eventSource.addEventListener('connected', (event) => { console.log('SSE 连接已建立:', event.data); }); eventSource.addEventListener('add', (event) => { currentMessage.value += event.data; }); eventSource.addEventListener('finish', (event) => { console.log('生成完成:', event.data); isGenerating.value = false; }); eventSource.addEventListener('error', (event) => { console.error('SSE 错误:', event); eventSource.close(); }); return eventSource; };
|
2. 发送消息
const sendMessage = async (message) => { if (!eventSource) { eventSource = connectSSE(userId); } const result = await api.sendMessage(userId, message, 'normal'); if (result.error) { console.error('发送失败:', result.message); return; } isGenerating.value = true; currentMessage.value = ''; };
|
3. 停止生成
const stopGeneration = async () => { const result = await api.stopGeneration(userId); if (!result.error) { isGenerating.value = false; console.log('已停止生成'); } };
|
🎯 核心技术点
1. 连接复用
问题: 每次对话都创建新连接,浪费资源
解决方案:
if (SSE_CACHE.containsKey(userId)) { return SSE_CACHE.get(userId); }
|
2. 并发安全
问题: 多个请求同时到达,可能产生竞态条件
解决方案:
private static final Map<String, SseEmitter> SSE_CACHE = new ConcurrentHashMap<>();
Disposable oldSubscription = userSubscriptions.get(userId); if (oldSubscription != null && !oldSubscription.isDisposed()) { oldSubscription.dispose(); }
|
3. 资源清理
问题: 连接断开后,资源未释放
解决方案:
sseEmitter.onTimeout(() -> { SSE_CACHE.remove(userId); });
sseEmitter.onError(throwable -> { SSE_CACHE.remove(userId); });
|
4. 错误恢复
问题: 发送失败后,连接状态不一致
解决方案:
try { sseEmitter.send(event); } catch (IOException e) { log.error("发送失败,移除连接", e); removeSse(userId); }
|
📊 性能优化
1. 批量发送
private StringBuilder buffer = new StringBuilder();
public void sendWithBuffer(String userId, String content) { buffer.append(content); if (buffer.length() >= 10) { SSEServer.sendMsg(userId, buffer.toString(), SSEMsgType.ADD); buffer.setLength(0); } }
|
2. 连接池管理
private static final int MAX_CONNECTIONS = 1000;
public static SseEmitter createSse(String userId) { if (SSE_CACHE.size() >= MAX_CONNECTIONS) { throw new RuntimeException("连接数已达上限"); } }
|
3. 心跳保活
@Scheduled(fixedRate = 30000) public void sendHeartbeat() { SSE_CACHE.forEach((userId, emitter) -> { try { emitter.send(SseEmitter.event() .name("heartbeat") .data("ping")); } catch (IOException e) { SSEServer.removeSse(userId); } }); }
|
⚠️ 常见问题
1. 连接立即断开
原因: 没有发送初始化消息
解决:
sseEmitter.send(SseEmitter.event() .name("connected") .data("SSE连接已建立"));
|
2. 中文乱码
原因: 字符编码问题
解决:
sseEmitter.send(SseEmitter.event() .data(content, MediaType.TEXT_PLAIN));
|
3. 浏览器自动重连
原因: SSE 默认会自动重连
解决:
🔍 调试技巧
1. 查看 SSE 事件
Chrome DevTools → Network → 选择 SSE 请求 → EventStream
2. 日志追踪
log.info("【SSE】用户: {}, 事件: {}, 数据: {}", userId, msgType, content);
|
3. 监控连接数
@GetMapping("/api/sse/stats") public Map<String, Object> getStats() { return Map.of( "connectionCount", SSEServer.getConnectionCount(), "timestamp", System.currentTimeMillis() ); }
|
🎯 总结
通过 SSE 技术,我们实现了:
- ✅ 实时推送:毫秒级延迟,打字机效果流畅
- ✅ 简单可靠:基于 HTTP,无需复杂配置
- ✅ 自动重连:浏览器原生支持,无需手动处理
- ✅ 资源高效:连接复用,减少服务器压力
SSE 是实现 AI 流式对话的最佳选择,代码简洁、性能出色、用户体验优秀。
📚 参考资料
下一篇:RAG 的实现
分享ShotAI-SSE流实现 
微信扫码打开,再通过右上角菜单分享。
交流与讨论
评论由 Valine / LeanCloud 提供,点击后连接第三方服务。