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服务。

Logo

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

更多推荐