SSE 服务端推送
一、SSE 概述
Server-Sent Events(SSE)是 HTML5 规范的一部分,允许服务器主动向客户端推送事件数据。与 WebSocket 双工通信不同,SSE 是单向通信——数据仅从服务器流向客户端。客户端通过 HTTP 连接发起请求,服务器以 text/event-stream 格式持续返回数据。
1.1 SSE 核心特性
| 特性 | 说明 |
|---|---|
| 传输方向 | 服务器 → 客户端(单工) |
| 协议 | HTTP/1.1 及 HTTP/2 |
| 数据格式 | UTF-8 文本,text/event-stream |
| 浏览器 API | EventSource |
| 自动重连 | 标准内置支持 |
| 同源限制 | 遵循同源策略 |
1.2 消息格式规范
SSE 数据流使用 UTF-8 编码的纯文本格式,每个事件由一个或多个字段行组成:
event: orderUpdate
id: 1724169600001
retry: 3000
data: {"amount": 15800, "orderNo": "ORD20250321001"}
event: userOnline
id: 1724169600002
data: {"onlineCount": 1286}每个事件以空行结束。字段说明:
event:事件类型,客户端通过addEventListener(type, callback)监听data:事件数据,支持多行拼接id:事件 ID,断线重连时通过Last-Event-ID头通知服务器retry:重连间隔(毫秒)
二、SseEmitter 源码原理
SseEmitter 是 Spring MVC 4.2 引入的 SSE 推送类,位于 org.springframework.web.servlet.mvc.method.annotation 包下,继承自 ResponseBodyEmitter,基于 Servlet 3.0 异步请求机制封装。
2.1 类层次与构造方法
ResponseBodyEmitter
└── SseEmitter| 构造方法 | 说明 |
|---|---|
SseEmitter() | 默认超时,继承 AsyncRequestTimeoutException 默认值 |
SseEmitter(Long timeout) | 自定义超时时间(毫秒) |
2.2 SseEventBuilder
内部构建器,用于构造符合 SSE 协议的事件:
SseEmitter emitter = new SseEmitter();
SseEmitter.SseEventBuilder builder = SseEmitter.event()
.id("1001")
.name("message")
.reconnectTime(3000L)
.data("{\"content\": \"hello\"}");
emitter.send(builder);2.3 核心源码分析
ResponseBodyEmitter 基类
// 异步请求回调接口
public interface Callback {
void onCompletion(Runnable callback);
void onTimeout(Runnable callback);
void onError(Consumer<Throwable> callback);
}
// 发送消息到响应输出流
public void send(Object object, MediaType mediaType) throws IOException { ... }
// 异步请求超时时间,0 表示无超时
private Long timeout;SseEmitter 扩展逻辑
public class SseEmitter extends ResponseBodyEmitter {
private static final MediaType SSE_CONTENT_TYPE =
new MediaType("text", "event-stream", StandardCharsets.UTF_8);
public SseEmitter() { super(); }
public SseEmitter(Long timeout) { super(timeout); }
// 将事件格式化为 SSE 协议文本
@Override
protected String formatEvent(Object data, String name, String id,
long expireTime) {
StringBuilder sb = new StringBuilder();
if (id != null) {
sb.append("id:").append(id).append('\n');
}
if (name != null) {
sb.append("event:").append(name).append('\n');
}
sb.append("data:").append(data).append('\n').append('\n');
return sb.toString();
}
}2.4 工作流程
客户端 服务器
│ │
│ GET /api/sse/push │
│ ──────────────────────────> │
│ │ SseEmitter emitter = new SseEmitter()
│ 200 OK │ return emitter; // 不关闭连接
│ Content-Type: │
│ text/event-stream │
│ <────────────────────────── │
│ │ new Thread(() -> {
│ │ emitter.send(data);
│ data: {...} │ // 持续推送...
│ <────────────────────────── │ }).start();
│ │
│ [客户端断线] │ emitter.onCompletion();服务器返回 SseEmitter 后,Servlet 容器线程立即释放,请求转入异步上下文。业务线程通过 emitter.send() 推送数据,通过 emitter.complete() 结束推送。
三、WebMvcConfigurer 异步请求配置
SseEmitter 正常工作需要配置 Spring MVC 的异步请求支持。
3.1 异步请求配置
@Configuration
@EnableWebMvc
public class WebMvcConfig implements WebMvcConfigurer {
@Override
public void configureAsyncSupport(AsyncSupportConfigurer configurer) {
configurer.setDefaultTimeout(30_000L);
configurer.setTaskExecutor(mvcAsyncExecutor());
}
@Bean
public ThreadPoolTaskExecutor mvcAsyncExecutor() {
ThreadPoolTaskExecutor executor = new ThreadPoolTaskExecutor();
executor.setCorePoolSize(5);
executor.setMaxPoolSize(20);
executor.setQueueCapacity(100);
executor.setThreadNamePrefix("mvc-async-");
executor.setWaitForTasksToCompleteOnShutdown(true);
executor.setAwaitTerminationSeconds(60);
executor.initialize();
return executor;
}
}3.2 配置项说明
| 配置项 | 作用 |
|---|---|
setDefaultTimeout(long) | 默认异步请求超时(毫秒) |
setTaskExecutor(AsyncTaskExecutor) | 异步请求线程池,防止 Servlet 容器线程耗尽 |
3.3 Spring Boot 配置
spring:
mvc:
async:
request-timeout: 30000注意:Spring Boot 中
@EnableWebMvc会接管自动配置,通常只需保留WebMvcConfigurer的 Bean 定义。
四、超时与异常处理
4.1 三层超时机制
| 层级 | 配置位置 | 默认值 |
|---|---|---|
| 1. 实例级 | new SseEmitter(timeout) | 继承层级 2 |
| 2. AsyncSupportConfigurer | setDefaultTimeout(ms) | 取决于 Servlet 容器 |
| 3. Servlet 容器 | asyncSupported timeout | Tomcat 默认 30 秒 |
优先级:层级 1 > 层级 2 > 层级 3。
4.2 超时与异常回调
@Component
public class SseService {
private final Map<String, SseEmitter> emitterMap = new ConcurrentHashMap<>();
public SseEmitter createEmitter(String clientId) {
SseEmitter emitter = new SseEmitter(5 * 60 * 1000L);
emitter.onTimeout(() -> {
System.err.println("[" + clientId + "] 超时");
emitterMap.remove(clientId);
});
emitter.onError(throwable -> {
System.err.println("[" + clientId + "] 异常: " + throwable.getMessage());
emitterMap.remove(clientId);
});
emitter.onCompletion(() -> {
System.out.println("[" + clientId + "] 关闭");
emitterMap.remove(clientId);
});
emitterMap.put(clientId, emitter);
return emitter;
}
}4.3 AsyncRequestTimeoutException 统一处理
@ControllerAdvice
public class AsyncExceptionHandler {
@ExceptionHandler(AsyncRequestTimeoutException.class)
public ResponseEntity<String> handleAsyncTimeout(AsyncRequestTimeoutException e) {
return ResponseEntity.status(HttpStatus.SERVICE_UNAVAILABLE)
.body("{\"code\": 503, \"message\": \"异步请求超时\"}");
}
}4.4 常见异常
| 异常 | 产生原因 | 处理方式 |
|---|---|---|
IllegalStateException | 连接已关闭后调用 send() | 发送前检查客户端状态 |
AsyncRequestTimeoutException | 超过异步请求超时时间 | 增大超时或心跳保活 |
IOException | 客户端断连导致写入失败 | 移除断连客户端的 emitter |
五、重连机制(Last-Event-ID、Retry)
SSE 协议原生支持断线重连,这是其相对轮询的核心优势。
5.1 Last-Event-ID 原理
连接中断后浏览器自动向同一 URL 发起重连,携带最后一次收到的事件 ID:
GET /api/sse/push HTTP/1.1
Accept: text/event-stream
Last-Event-ID: 1724169600042服务器据此只推送增量数据,避免重复推送。
5.2 服务端实现
@RestController
@RequestMapping("/api/sse")
public class SseController {
@GetMapping(value = "/push", produces = MediaType.TEXT_EVENT_STREAM_VALUE)
public SseEmitter push(HttpServletRequest request) {
String lastEventId = request.getHeader("Last-Event-ID");
String clientId = request.getParameter("clientId");
SseEmitter emitter = sseService.createEmitter(clientId);
if (StringUtils.hasText(lastEventId)) {
sseService.sendMissedEvents(clientId, lastEventId, emitter);
}
return emitter;
}
}5.3 Retry 与心跳保活
// 设置重连间隔
SseEmitter.SseEventBuilder event = SseEmitter.event()
.reconnectTime(5000L) // 5 秒重连间隔
.name("heartbeat")
.data("keep-alive");
emitter.send(event);
// 心跳保活——防止代理服务器空闲超时关闭连接
public void startHeartbeat(SseEmitter emitter) {
ScheduledExecutorService scheduler = Executors.newSingleThreadScheduledExecutor();
scheduler.scheduleAtFixedRate(() -> {
try {
emitter.send(SseEmitter.event().name("heartbeat").data("ping"));
} catch (IOException e) {
scheduler.shutdown();
}
}, 0, 10, TimeUnit.SECONDS);
}5.4 完整重连流程
客户端 服务器
│ GET /api/sse/push │
│ ───────────────────────> │
│ id: 100 │
│ <─────────────────────── │
│ id: 101 │
│ <─────────────────────── │
│ [网络中断] │
│ [等待 retry 时长] │
│ GET /api/sse/push │
│ Last-Event-ID: 101 │ ← 浏览器自动携带
│ ───────────────────────> │
│ id: 102 │ 从 ID 102 开始推送
│ <─────────────────────── │六、SSE vs WebSocket 对比
6.1 核心差异
| 对比维度 | SSE | WebSocket |
|---|---|---|
| 通信方向 | 服务器 → 客户端(单工) | 双向(全双工) |
| 协议 | HTTP/HTTPS | ws:// / wss://(独立协议) |
| 数据格式 | 文本(UTF-8) | 文本 + 二进制 |
| 浏览器 API | EventSource | WebSocket |
| 自动重连 | 内建支持 | 需自行实现 |
| 连接数限制 | HTTP/1.1 同域 6~8 个 | 无此限制 |
| 跨防火墙 | HTTP 端口,友好 | 需配置代理支持 WS |
| 实现复杂度 | 低 | 中高 |
6.2 协议差异
SSE 基于 HTTP:以 Transfer-Encoding: chunked 持续响应,仍是标准 HTTP 请求-响应模型。
WebSocket 升级协议:通过 HTTP Upgrade 切换到独立 WS 协议:
GET /ws/chat HTTP/1.1
Upgrade: websocket
Connection: Upgrade
Sec-WebSocket-Key: dGhlIHNhbXBsZSBub25jZQ==
HTTP/1.1 101 Switching Protocols
Upgrade: websocket
Connection: Upgrade
Sec-WebSocket-Accept: s3pPLMBiTxaQ9kYGzzhZRbK+xOo=升级后使用 WS 协议帧(frame)进行双向数据交换。
6.3 适用场景
| 场景 | 推荐方案 | 理由 |
|---|---|---|
| 数据大屏 / 监控面板 | SSE | 单向推送,实现简单 |
| AI 流式输出 | SSE | 天然适合流式文本 |
| 通知推送 | SSE | 浏览器内置重连 |
| 在线聊天 | WebSocket | 需要双向通信 |
| 在线游戏 | WebSocket | 低延迟双向通信 |
6.4 优缺点
| 优点 | 缺点 | |
|---|---|---|
| SSE | 基于 HTTP + 原生自动重连 + 实现简单 | 单向 + 不支持二进制 + HTTP/1.1 连接数限制 |
| WebSocket | 全双工低延迟 + 无连接数限制 + 支持二进制 | 需额外握手 + 自行实现重连 + 代理配置复杂 |
七、实战:运营后台数据大屏 SSE 实时推送
实现一个运营数据大屏,实时推送:订单成交额、用户在线数、今日订单数、系统健康状态(CPU/内存)。
7.1 项目依赖
<parent>
<groupId>org.springframework.boot</groupId>
<artifactId>spring-boot-starter-parent</artifactId>
<version>3.2.5</version>
</parent>
<dependencies>
<dependency>
<groupId>org.springframework.boot</groupId>
<artifactId>spring-boot-starter-web</artifactId>
</dependency>
</dependencies>7.2 异步配置
@Configuration
public class AsyncConfig implements WebMvcConfigurer {
@Override
public void configureAsyncSupport(AsyncSupportConfigurer configurer) {
configurer.setDefaultTimeout(0L);
configurer.setTaskExecutor(sseTaskExecutor());
}
@Bean
public ThreadPoolTaskExecutor sseTaskExecutor() {
ThreadPoolTaskExecutor executor = new ThreadPoolTaskExecutor();
executor.setCorePoolSize(10);
executor.setMaxPoolSize(50);
executor.setQueueCapacity(200);
executor.setThreadNamePrefix("sse-push-");
executor.setWaitForTasksToCompleteOnShutdown(true);
executor.setAwaitTerminationSeconds(30);
executor.initialize();
return executor;
}
}7.3 数据模型
public class DashboardData {
private BigDecimal orderAmount; // 订单成交额
private Integer onlineUserCount; // 用户在线数
private Integer todayOrderCount; // 今日订单数
private Double cpuUsage; // CPU 使用率
private Double memoryUsage; // 内存使用率
private Long timestamp; // 数据时间戳
// getters/setters 省略(可使用 Lombok @Data)
}7.4 SSE 推送服务
@Service
public class DashboardSseService {
private final Map<String, SseEmitter> clients = new ConcurrentHashMap<>();
private final Random random = new Random();
public SseEmitter createConnection(String clientId, String lastEventId) {
removeClient(clientId);
SseEmitter emitter = new SseEmitter(3600_000L);
emitter.onTimeout(() -> clients.remove(clientId));
emitter.onError(ex -> clients.remove(clientId));
emitter.onCompletion(() -> clients.remove(clientId));
clients.put(clientId, emitter);
if (StringUtils.hasText(lastEventId)) {
sendMissedData(clientId, lastEventId, emitter);
}
return emitter;
}
public void broadcastDashboardData(DashboardData data) {
if (clients.isEmpty()) return;
List<String> disconnected = new ArrayList<>();
clients.forEach((clientId, emitter) -> {
try {
SseEmitter.SseEventBuilder event = SseEmitter.event()
.id(String.valueOf(System.currentTimeMillis()))
.name("dashboard")
.data(JsonUtils.toJson(data), MediaType.APPLICATION_JSON);
emitter.send(event);
} catch (IOException e) {
disconnected.add(clientId);
}
});
disconnected.forEach(clients::remove);
}
private void sendMissedData(String clientId, String lastEventId,
SseEmitter emitter) {
long lastId = Long.parseLong(lastEventId);
long currentId = System.currentTimeMillis();
if (currentId - lastId > 300_000) {
try {
DashboardData summary = new DashboardData();
summary.setOrderAmount(BigDecimal.valueOf(random.nextInt(50000) + 100000));
summary.setOnlineUserCount(random.nextInt(3000) + 500);
summary.setTodayOrderCount(random.nextInt(2000) + 500);
summary.setCpuUsage(45.0 + random.nextDouble() * 40);
summary.setMemoryUsage(50.0 + random.nextDouble() * 30);
summary.setTimestamp(System.currentTimeMillis());
emitter.send(SseEmitter.event()
.id(String.valueOf(currentId))
.name("dashboard-replay")
.data(JsonUtils.toJson(summary), MediaType.APPLICATION_JSON));
} catch (IOException ignored) {}
}
}
public void sendHeartbeat() {
clients.forEach((clientId, emitter) -> {
try {
emitter.send(SseEmitter.event().name("heartbeat").data("ping"));
} catch (IOException ignored) {}
});
}
public void removeClient(String clientId) {
SseEmitter existing = clients.remove(clientId);
if (existing != null) existing.complete();
}
public int getClientCount() { return clients.size(); }
}7.5 Controller
@RestController
@RequestMapping("/api/dashboard")
public class DashboardController {
private final DashboardSseService sseService;
@GetMapping(value = "/stream", produces = MediaType.TEXT_EVENT_STREAM_VALUE)
public SseEmitter stream(HttpServletRequest request) {
String clientId = request.getParameter("clientId");
if (!StringUtils.hasText(clientId)) {
clientId = UUID.randomUUID().toString();
}
String lastEventId = request.getHeader("Last-Event-ID");
return sseService.createConnection(clientId, lastEventId);
}
}7.6 数据推送调度
@Component
public class DashboardDataPushScheduler {
private final DashboardSseService sseService;
private final Random random = new Random();
@Scheduled(fixedRate = 2000)
public void pushDashboardData() {
DashboardData data = new DashboardData();
data.setOrderAmount(BigDecimal.valueOf(random.nextInt(100000) + 50000));
data.setOnlineUserCount(random.nextInt(5000) + 1000);
data.setTodayOrderCount(random.nextInt(3000) + 800);
data.setCpuUsage(35.0 + random.nextDouble() * 50);
data.setMemoryUsage(45.0 + random.nextDouble() * 40);
data.setTimestamp(System.currentTimeMillis());
sseService.broadcastDashboardData(data);
}
@Scheduled(fixedRate = 10000)
public void pushHeartbeat() {
sseService.sendHeartbeat();
}
}7.7 前端实现
<!DOCTYPE html>
<html lang="zh-CN">
<head>
<meta charset="UTF-8">
<title>运营数据大屏</title>
<style>
.dashboard { display: grid; grid-template-columns: repeat(4, 1fr); gap: 20px; padding: 20px; }
.card { background: #1a1a2e; color: #fff; border-radius: 12px; padding: 24px; text-align: center; }
.card .value { font-size: 2.5rem; font-weight: bold; color: #00d4ff; margin: 12px 0; }
.card .label { font-size: 0.9rem; color: #8892b0; }
.status-bar { position: fixed; bottom: 0; left: 0; right: 0; padding: 8px 20px; background: #0f0f23; color: #8892b0; }
</style>
</head>
<body>
<div class="dashboard" id="dashboard"></div>
<div class="status-bar">
<span id="connectionStatus">连接状态: 未连接</span>
<span id="lastUpdate">最后更新: --</span>
</div>
<script>
const cards = [
{ id: 'orderAmount', label: '今日成交额', fmt: d => '¥' + Number(d).toLocaleString() },
{ id: 'onlineUserCount', label: '用户在线', fmt: d => Number(d).toLocaleString() },
{ id: 'todayOrderCount', label: '今日订单数', fmt: d => Number(d).toLocaleString() },
{ id: 'systemHealth', label: 'CPU / 内存', fmt: (d, data) => data.cpuUsage.toFixed(1) + '% / ' + data.memoryUsage.toFixed(1) + '%' }
];
document.getElementById('dashboard').innerHTML = cards.map(c =>
'<div class="card"><div class="label">' + c.label + '</div><div class="value" id="' + c.id + '">--</div></div>'
).join('');
const es = new EventSource('/api/dashboard/stream?clientId=dashboard-' + Date.now());
es.addEventListener('dashboard', function(event) {
const d = JSON.parse(event.data);
cards.forEach(c => {
document.getElementById(c.id).textContent = c.fmt(d[c.id] ?? d, d);
});
document.getElementById('lastUpdate').textContent = '最后更新: ' + new Date(d.timestamp).toLocaleTimeString();
});
es.onopen = function() {
const el = document.getElementById('connectionStatus');
el.textContent = '连接状态: 已连接'; el.style.color = '#00d4ff';
};
es.onerror = function() {
const el = document.getElementById('connectionStatus');
el.textContent = '连接状态: 已断开,正在重连...'; el.style.color = '#ff6b6b';
};
window.addEventListener('beforeunload', function() { es.close(); });
</script>
</body>
</html>7.8 启动配置与运行
spring:
task:
scheduling:
pool:
size: 5
server:
port: 8080在启动类上启用定时调度:
@SpringBootApplication
@EnableScheduling
public class DashboardApplication {
public static void main(String[] args) {
SpringApplication.run(DashboardApplication.class, args);
}
}7.9 流程图
浏览器 EventSource DashboardController DashboardSseService DataPushScheduler
│ │ │ │
│ GET /stream?clientId=xxx │ │ │
│ ──────────────────────────────────> │ │ │
│ 200 text/event-stream │ │ │
│ <────────────────────────────────── │ │ │
│ │ │ │
│ event: dashboard (每2秒) │ │
│ <══════════════════════════════════════════════════════════╡ @Scheduled(2000ms) │
│ data: {orderAmount,...} │ │ │
│ │ │ │
│ event: heartbeat (每10秒) │ │
│ <══════════════════════════════════════════════════════════╡ @Scheduled(10000ms) │
│ data: ping │ │ │
│ │ │ │
│ [断线] 自动重连携带 Last-Event-ID │ │ │
│ ──────────────────────────────────> │ │ │
│ event: dashboard-replay │ 补发摘要数据 │ │
│ <════════════════════════════════════╡ │ │7.10 最佳实践
- 线程安全:
ConcurrentHashMap管理客户端,send()在单线程中串行调用 - 资源释放:必须注册
onTimeout/onError/onCompletion回调清理断连客户端,防止连接泄漏 - 心跳与容错:定期发送心跳防止代理服务器关闭连接;
broadcast中捕获 IOException 避免单客户端影响全局 - 连接数控制:设置上限(如 1000),超限拒绝新连接;HTTP/2 下不考虑同域限制
- CORS:跨域场景需允许 SSE 端点跨域访问
八、总结
SSE 作为轻量级服务端推送方案,在 Spring MVC 中通过 SseEmitter 类得到优雅支持。相比 WebSocket,SSE 的优势在于极低的实现复杂度、原生 HTTP 协议兼容性以及内置的重连机制。对于单向数据流场景——如运营数据大屏、实时监控、通知推送、AI 流式输出——SSE 是比 WebSocket 更合适的选择。Spring MVC 的异步请求机制配合 SseEmitter,使 SSE 方案可以无缝融入现有 Spring Boot 项目,仅需少量配置即可实现生产级实时数据推送。