Spring Boot实现智能体实时SSE流式响应,Vue三元表达式。
·
Spring Boot SSE 流式输出实现智能体实时响应
SSE(Server-Sent Events)是一种基于HTTP的服务器推送技术,允许服务器主动向客户端发送事件流数据。Spring Boot提供了对SSE的原生支持,结合智能体的实时响应需求,可以构建高效的数据流传输方案。
依赖配置
在Spring Boot项目中添加Web依赖即可支持SSE功能,无需额外库。Maven配置如下:
<dependency>
<groupId>org.springframework.boot</groupId>
<artifactId>spring-boot-starter-web</artifactId>
</dependency>
服务端实现
创建SSE事件流控制器,使用SseEmitter对象管理连接:
@RestController
@RequestMapping("/sse")
public class SseController {
private final Map<String, SseEmitter> emitters = new ConcurrentHashMap<>();
@GetMapping("/connect/{id}")
public SseEmitter connect(@PathVariable String id) {
SseEmitter emitter = new SseEmitter(30_000L); // 30秒超时
emitters.put(id, emitter);
emitter.onCompletion(() -> emitters.remove(id));
emitter.onTimeout(() -> emitters.remove(id));
return emitter;
}
@Async
public void pushEvent(String id, String data) {
SseEmitter emitter = emitters.get(id);
if (emitter != null) {
try {
emitter.send(SseEmitter.event()
.data(data)
.name("message"));
} catch (IOException e) {
emitter.complete();
emitters.remove(id);
}
}
}
}
客户端实现
浏览器端通过EventSource API接收事件流:
const eventSource = new EventSource('/sse/connect/user123');
eventSource.onmessage = (e) => {
console.log('Received:', e.data);
};
eventSource.addEventListener('message', (e) => {
document.getElementById('output').innerHTML += e.data + '<br>';
});
智能体集成模式
将AI模型响应拆分为多个数据块进行流式传输:
public void streamAIResponse(String prompt, String clientId) {
List<String> chunks = aiService.generateStream(prompt);
for (String chunk : chunks) {
sseController.pushEvent(clientId, chunk);
Thread.sleep(100); // 控制发送间隔
}
sseController.pushEvent(clientId, "[DONE]");
}
性能优化策略
设置合理的心跳机制保持连接:
@Scheduled(fixedRate = 15000)
public void sendHeartbeat() {
emitters.forEach((id, emitter) -> {
try {
emitter.send(SseEmitter.event()
.comment("heartbeat"));
} catch (IOException ex) {
emitter.complete();
emitters.remove(id);
}
});
}
采用背压控制防止消息堆积:
private final Semaphore semaphore = new Semaphore(100);
public void pushEventWithBackpressure(String id, String data) {
try {
semaphore.acquire();
SseEmitter emitter = emitters.get(id);
if (emitter != null) {
emitter.send(data, () -> semaphore.release());
}
} catch (InterruptedException e) {
Thread.currentThread().interrupt();
}
}
异常处理机制
增强连接稳定性处理:
emitter.onError((ex) -> {
log.error("SSE error for {}: {}", id, ex.getMessage());
emitter.complete();
emitters.remove(id);
// 触发重连逻辑
});
安全增强措施
添加CSRF防护和权限验证:
@GetMapping("/connect/{id}")
public SseEmitter connect(
@PathVariable String id,
@RequestHeader("X-Auth-Token") String token) {
if (!authService.validateToken(token)) {
throw new SecurityException("Invalid token");
}
// ...原有连接逻辑
}
监控指标收集
通过Micrometer暴露SSE指标:
@Bean
public MeterRegistryCustomizer<MeterRegistry> sseMetrics() {
return registry -> Gauge.builder("sse.active_connections",
emitters::size)
.register(registry);
}
集群环境适配
Redis发布订阅实现跨节点消息广播:
@Bean
public RedisMessageListenerContainer redisContainer(
RedisConnectionFactory factory) {
RedisMessageListenerContainer container = new RedisMessageListenerContainer();
container.setConnectionFactory(factory);
container.addMessageListener((message, pattern) -> {
String clientId = new String(message.getBody());
String data = redisTemplate.opsForValue().get(clientId);
pushEvent(clientId, data);
}, new ChannelTopic("sse:notify"));
return container;
}
这种实现方式特别适合需要持续输出AI生成内容、实时数据监控、金融行情推送等场景。通过合理的超时设置、心跳机制和错误处理,可以构建高可用的SSE服务。
更多推荐
所有评论(0)