JOURNAL / 随笔

ShotAI-SSE流实现

ShotAI-SSE流实现 封面
本文目录
  1. ShotAI SSE 聊天流实现详解
    1. 📚 概述
    2. 🤔 为什么选择 SSE?
      1. 传统轮询 vs SSE vs WebSocket
    3. 🏗️ 架构设计
      1. 整体流程图
    4. 💻 后端实现
      1. 1. SSEMsgType 枚举
      2. 2. SSEServer 工具类
      3. 3. SseController 控制器
      4. 4. 与 AI 流的集成
    5. 🌐 前端实现
      1. 1. 建立 SSE 连接
      2. 2. 发送消息
      3. 3. 停止生成
    6. 🎯 核心技术点
      1. 1. 连接复用
      2. 2. 并发安全
      3. 3. 资源清理
      4. 4. 错误恢复
    7. 📊 性能优化
      1. 1. 批量发送
      2. 2. 连接池管理
      3. 3. 心跳保活
    8. ⚠️ 常见问题
      1. 1. 连接立即断开
      2. 2. 中文乱码
      3. 3. 浏览器自动重连
    9. 🔍 调试技巧
      1. 1. 查看 SSE 事件
      2. 2. 日志追踪
      3. 3. 监控连接数
    10. 🎯 总结
    11. 📚 参考资料

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 枚举

/**
* SSE 消息类型枚举
*/
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 {

/**
* 存储所有用户的 SSE 连接
* Key: userId (用户唯一标识)
* Value: SseEmitter (SSE 发射器)
*/
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);
}

// 超时时间设置为 30 分钟
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());

// 立即发送一条初始化消息,确保触发前端的 onopen 事件
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;
}

关键设计点:

  1. 超时设置:30 分钟,平衡资源占用和用户体验
  2. 初始化消息:立即发送一条消息,触发前端 onopen 事件
  3. 回调处理:完善的生命周期管理
  4. 幂等性:重复调用返回同一个连接

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;

/**
* 建立 SSE 连接
*
* @param userId 用户ID
* @return SSE 发射器
*/
@GetMapping("/connect")
public SseEmitter connect(@RequestParam String userId) {
log.info("建立 SSE 连接,用户: {}", userId);
return chatService.createSseConnection(userId);
}

/**
* 停止生成
*
* @param userId 用户ID
* @return 响应结果
*/
@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 {
// 1. 创建提示词
Prompt prompt = new Prompt(message);

// 2. 获取流式响应
Flux<String> stream = chatClient.prompt(prompt).stream().content();

// 3. 订阅流并通过 SSE 推送
Disposable subscription = stream
.subscribe(
// onNext: 每收到一个文本片段,立即推送
content -> {
if (SSEServer.exists(userId)) {
SSEServer.sendMsg(userId, content, SSEMsgType.ADD);
}
},
// onError: 错误处理
error -> {
log.error("【对话流失败】用户: {}", userId, error);
SSEServer.sendMsg(userId, "抱歉,出现了错误", SSEMsgType.ERROR);
userSubscriptions.remove(userId);
},
// onComplete: 完成回调
() -> {
log.info("【对话完成】用户: {}", userId);
SSEServer.sendMsg(userId, "done", SSEMsgType.FINISH);
SSEServer.close(userId);
userSubscriptions.remove(userId);
}
);

// 4. 保存订阅引用
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 连接

// 创建 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) => {
// 1. 先建立 SSE 连接(如果还没有)
if (!eventSource) {
eventSource = connectSSE(userId);
}

// 2. 发送消息到后端
const result = await api.sendMessage(userId, message, 'normal');

if (result.error) {
console.error('发送失败:', result.message);
return;
}

// 3. 等待 SSE 推送结果
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. 并发安全

问题: 多个请求同时到达,可能产生竞态条件

解决方案:

// 使用 ConcurrentHashMap
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) { // 累积 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) // 每 30 秒
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)); // 指定 MediaType

3. 浏览器自动重连

原因: SSE 默认会自动重连

解决:

// 前端主动关闭连接
eventSource.close();

🔍 调试技巧

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 提供,点击后连接第三方服务。

搜索文章