Spring Boot SSE 流式输出实现智能体实时响应

Server-Sent Events (SSE) 是一种基于 HTTP 的轻量级技术,允许服务器向客户端推送实时更新。Spring Boot 提供了简洁的 API 实现 SSE,适用于智能体对话、实时监控等场景。

服务端实现

在 Spring Boot 中创建 SSE 控制器需要 @RestControllerSseEmitter 对象。以下示例展示如何建立连接并发送流式数据:

@RestController
@RequestMapping("/sse")
public class SseController {

    private final Map<String, SseEmitter> emitters = new ConcurrentHashMap<>();

    @GetMapping("/connect/{clientId}")
    public SseEmitter connect(@PathVariable String clientId) {
        SseEmitter emitter = new SseEmitter(60_000L); // 超时时间60秒
        emitters.put(clientId, emitter);

        emitter.onCompletion(() -> emitters.remove(clientId));
        emitter.onTimeout(() -> emitters.remove(clientId));
        
        return emitter;
    }

    @Async
    public void sendData(String clientId, String data) {
        SseEmitter emitter = emitters.get(clientId);
        if (emitter != null) {
            try {
                emitter.send(SseEmitter.event()
                    .data(data)
                    .name("message"));
            } catch (IOException e) {
                emitter.completeWithError(e);
            }
        }
    }
}
客户端实现

前端通过 EventSource API 监听服务器推送:

const eventSource = new EventSource('/sse/connect/user123');

eventSource.addEventListener('message', (event) => {
    const data = JSON.parse(event.data);
    console.log('Received update:', data);
});

eventSource.onerror = (err) => {
    console.error('SSE error:', err);
};
智能体集成模式

对于 AI 智能体场景,可将流式响应拆分为多个数据块发送:

  1. 分块处理:将大语言模型生成的响应按语义段落拆分
  2. 标记状态:每个数据块包含类型标记(如 partial|complete
  3. 错误恢复:实现重连机制和序列号验证
public void streamAgentResponse(String sessionId, String query) {
    List<String> chunks = aiService.generateStreamingResponse(query);
    for (int i = 0; i < chunks.size(); i++) {
        String payload = String.format("{'index':%d,'content':'%s','status':'%s'}", 
            i, chunks.get(i), i == chunks.size()-1 ? "complete" : "partial");
        sseService.sendData(sessionId, payload);
    }
}
性能优化建议
  • 设置合理的 keepalive 间隔防止连接断开
  • 使用线程池管理异步发送任务
  • 添加心跳包维持连接活跃
  • 采用 Protobuf 或 MessagePack 替代 JSON 减少带宽
异常处理机制

实现健壮性需要处理以下场景:

  • 客户端意外断开时自动清理资源
  • 网络抖动时支持事件ID追踪和断点续传
  • 服务端背压控制防止消息堆积
emitter.onError((ex) -> {
    log.warn("Client {} error: {}", clientId, ex.getMessage());
    cleanupResources(clientId);
});

这种模式已成功应用于智能客服、实时数据分析等场景,相比 WebSocket 具有更简单的协议栈和自动重连优势。通过合理设计数据格式和状态机,可以构建高可用的实时交互系统。

Logo

北京人形旗下天工造物具身智能开源社区,聚焦具身天工与慧思开物两大平台

更多推荐