Java实现SSE案例:从零搭建服务端实时推送的完整指南
目录导读
- 什么是SSE?为何在Java中实现它?
- SSE vs WebSocket:选型对比与适用场景
- 环境准备与技术栈(JDK 17 + Spring Boot 3)
- 核心实现:基于Spring MVC的SSE端点
- 进阶玩法:异步线程池 + 心跳机制 + 断线重连
- 前端接入:EventSource API详解
- 实战案例:股票行情推送系统(附完整代码)
- 常见问题问答(FAQ)
- 性能调优与生产级注意事项
- 总结与最佳实践
什么是SSE?为何在Java中实现它?
SSE(Server-Sent Events,服务端发送事件)是一种基于HTTP的轻量级实时通信协议,它允许服务器主动向客户端推送数据,客户端只需通过EventSource接口即可订阅,与WebSocket不同,SSE是单向的(仅服务端→客户端),但它在Java中实现极其简单——你不需要额外的库,Spring框架原生支持。

为何选择Java实现?
Java生态成熟、Spring Boot的自动配置让SSE端点开发变得异常简洁,且天然支持线程池异步化,非常适合股票、日志流、AI聊天等实时场景。
SSE vs WebSocket:选型对比
| 维度 | SSE | WebSocket |
|---|---|---|
| 通信方向 | 单向(服务端→客户端) | 双向 |
| 协议 | HTTP/1.1+(天然穿透代理) | 独立协议(需要握手升级) |
| 自动重连 | 内置(浏览器自动重连) | 需手动实现 |
| 二进制支持 | 仅文本(Base64编码可绕过) | 原生支持 |
| 实现复杂度 | 极低(Spring注解) | 较高(需要处理连接状态) |
| 适用场景 | 通知推送、实时图表、日志流 | 在线游戏、聊天、协同编辑 |
如果你的业务只需要服务端主动推送,SSE是更轻、更可靠的选择。
环境准备与技术栈
- JDK 17+(推荐使用记录类型简化代码)
- Spring Boot 3.x(内置Spring MVC 6)
- Maven/Gradle(构建工具)
- Lombok(可选,减少样板代码)
<dependency>
<groupId>org.springframework.boot</groupId>
<artifactId>spring-boot-starter-web</artifactId>
</dependency>
核心实现:基于Spring MVC的SSE端点
基础版:直接返回SseEmitter
@RestController
@RequestMapping("/api/sse")
public class SseController {
@GetMapping(value = "/stream", produces = MediaType.TEXT_EVENT_STREAM_VALUE)
public SseEmitter stream() {
SseEmitter emitter = new SseEmitter(0L); // 0表示永不超时
// 模拟数据推送
new Thread(() -> {
try {
for (int i = 0; i < 10; i++) {
emitter.send(SseEmitter.event()
.name("message")
.data(Map.of("id", i, "time", LocalTime.now().toString())));
Thread.sleep(1000);
}
emitter.complete();
} catch (Exception e) {
emitter.completeWithError(e);
}
}).start();
return emitter;
}
}
关键点:
produces = MediaType.TEXT_EVENT_STREAM_VALUE指定响应类型为SSE流。SseEmitter.event().name("xxx")自定义事件名,前端监听对应事件。- 必须手动调用
complete()或completeWithError()来结束连接。
改进版:使用ExecutorService线程池
private final ExecutorService executor = Executors.newCachedThreadPool();
@GetMapping("/stream2")
public SseEmitter stream2() {
SseEmitter emitter = new SseEmitter(60_000L); // 60秒超时
executor.execute(() -> {
try {
// 业务逻辑
emitter.send("初始化数据...");
emitter.complete();
} catch (IOException e) {
emitter.completeWithError(e);
}
});
emitter.onTimeout(() -> System.out.println("连接超时"));
emitter.onCompletion(() -> System.out.println("连接关闭"));
return emitter;
}
进阶玩法:心跳机制与断线重连
在生产环境中,代理服务器(如Nginx)默认会关闭空闲连接,必须定期发送注释(如heartbeat)保持活跃。
@Scheduled(fixedRate = 15000)
public void heartbeat() {
// 遍历所有活跃的SseEmitter,发送注释
activeEmitters.forEach((id, emitter) -> {
try {
emitter.send(SseEmitter.event().comment("heartbeat"));
} catch (Exception e) {
activeEmitters.remove(id);
}
});
}
前端自动重连:EventSource内建重连机制,服务端只需正常发送数据,客户端断开后会自动尝试重新连接(默认延迟3秒)。
前端接入:EventSource API详解
const eventSource = new EventSource('http://localhost:8080/api/sse/stream');
// 监听默认消息(无event字段)
eventSource.onmessage = (event) => {
console.log('收到数据:', JSON.parse(event.data));
};
// 监听自定义事件
eventSource.addEventListener('message', (event) => {
// 如果服务端设置了事件名
});
// 连接错误处理
eventSource.onerror = (event) => {
if (event.target.readyState === EventSource.CLOSED) {
console.log('连接已关闭');
} else {
console.log('正在重连...');
}
};
实战案例:股票行情推送系统(完整代码)
@RestController
@RequestMapping("/stock")
public class StockController {
private final Map<String, SseEmitter> clients = new ConcurrentHashMap<>();
@GetMapping("/subscribe/{code}")
public SseEmitter subscribe(@PathVariable String code) {
SseEmitter emitter = new SseEmitter(0L);
clients.put(code + "-" + UUID.randomUUID(), emitter);
emitter.onCompletion(() -> clients.remove(code));
emitter.onTimeout(() -> clients.remove(code));
return emitter;
}
@PostMapping("/push/{code}")
public String pushPrice(@PathVariable String code, @RequestBody PriceData data) {
for (Map.Entry<String, SseEmitter> entry : clients.entrySet()) {
if (entry.getKey().startsWith(code)) {
try {
entry.getValue().send(SseEmitter.event()
.name("price")
.data(data));
} catch (IOException e) {
clients.remove(entry.getKey());
}
}
}
return "推送成功";
}
}
前端调用:
const es = new EventSource('/stock/subscribe/600519');
es.addEventListener('price', (e) => {
document.getElementById('price').textContent = JSON.parse(e.data).price;
});
常见问题问答(FAQ)
Q1: SSE连接为什么会频繁断开?
A: 很可能是代理服务器空闲超时所致,解决方案:① 服务端开启心跳线程定时发送注释;② 配置Nginx的proxy_read_timeout。
Q2: 多个客户端共享同一个Emitter会怎样?
A: 每个客户端必须有独立的SseEmitter实例,否则数据会串流,务必使用Map存储每个连接的Emitter。
Q3: SSE支持自定义HTTP状态码吗?
A: 支持,但一旦响应头已发送则无法修改,建议在业务逻辑开始前设置响应状态。
Q4: 如何处理客户端异常断开?
A: 在onCompletion和onTimeout回调中清理资源,同时捕获IOException。
性能调优与生产级注意事项
- 连接限制:默认Tomcat支持200个并发线程,若推送频率高需调大
server.tomcat.threads.max。 - 异步化:推送业务务必使用
@Async或独立线程池,避免阻塞Tomcat工作线程。 - 数据压缩:对大型JSON启用GZIP压缩(
server.compression.enabled=true)。 - 监控:通过
Metrics记录活跃连接数、推送失败率。 - 安全:生产环境使用HTTPS并添加鉴权(如JWT放在URL参数,SSE不支持自定义Header)。
总结与最佳实践
SSE是Java世界中最被低估的实时通信方案,通过Spring Boot的SseEmitter,你可以在30分钟内完成一个可靠的实时推送系统,最佳实践归纳:
- ✅ 优先选择SSE而非WebSocket,除非需要双向通信
- ✅ 务必实现心跳机制应对网络代理
- ✅ 使用线程池隔离业务逻辑,保证响应速度
- ✅ 前端配合
EventSource自动重连,提升用户体验
拿起你的IDE,动手实现第一个Java SSE案例吧!如果你有任何问题,欢迎在评论区留言讨论。