SSE技术详解:从HTTP实时推送到Spring Boot实战应用
1. 从轮询到流式:为什么我们需要SSE?
做后端开发的朋友,对“实时数据推送”这个需求肯定不陌生。早些年,我们最朴素的想法就是让前端隔三差五来问一次:“后端,你有新数据吗?”这就是所谓的轮询(Polling)。简单粗暴,但问题一大堆:无效请求多,服务器压力大,数据延迟高,用户体验差。后来,WebSocket横空出世,它允许服务器和客户端建立一个全双工的、长久的连接,双方可以随时互发消息,真正实现了实时双向通信。一时间,WebSocket成了实时应用的代名词,聊天室、在线游戏、协同编辑,哪里都有它的身影。
但WebSocket真的是所有实时场景的“银弹”吗?未必。我遇到过很多这样的需求:后台有一个耗时的数据处理任务,前端需要实时看到任务进度;或者,服务器需要向客户端持续推送新闻快讯、股票价格、监控日志。在这些场景里,数据流主要是单向的,从服务器流向客户端。如果为此专门搭建一套WebSocket服务,引入复杂的握手协议、心跳维护、双向消息处理逻辑,无异于“杀鸡用牛刀”,增加了不必要的复杂性和开发成本。
这时候,SSE(Server-Sent Events)就该登场了。它就像给HTTP协议装上了一根“单向水管”。客户端发起一个普通的HTTP请求,服务器就可以通过这个连接,源源不断地、主动地向客户端发送数据流,而且这个连接会一直保持。对于上面提到的那些“服务器主动推送,客户端被动接收”的场景,SSE方案简洁、高效、天然契合。
最近在折腾一些AI应用和实时数据仪表盘,SSE更是成了我的首选。比如,用Spring AI处理大语言模型的流式响应,或者向管理后台实时推送系统告警,用SSE来实现“流式输出”,代码写起来那叫一个清爽。所以,今天我就结合自己踩过的坑和积累的经验,来好好聊聊SSE协议,从原理、实战到避坑,给你讲明白。
2. SSE协议核心原理与工作机制拆解
要理解SSE,我们不能只停留在“会用”的层面,得先看看它的“内脏”是怎么工作的。这能帮你未来在排查一些诡异问题时,心里有张清晰的地图。
2.1 协议基础:建立在HTTP之上的事件流
SSE本质上不是一个独立的协议,而是对HTTP协议的一种使用约定。它完全基于标准的HTTP/1.1或HTTP/2,这意味着它不需要像WebSocket那样升级协议(Upgrade: websocket),几乎能被所有现代浏览器和HTTP客户端库原生支持。
它的核心在于两个HTTP头:
Content-Type: text/event-stream:这是SSE的“身份证”。当服务器返回这个响应头时,就是在告诉客户端:“接下来我发送的不是一个普通的HTML或JSON文档,而是一个遵循SSE格式的事件流。”Cache-Control: no-cache:确保中间代理或浏览器不会缓存这个流式响应。Connection: keep-alive:这是保持长连接的关键,让TCP连接在多次数据传输后依然存活。
客户端的工作很简单,就是使用EventSourceAPI(浏览器原生)或类似的库,向一个特定的URL发起一个GET请求。当它收到text/event-stream的响应后,便会挂起这个请求,进入等待状态,而不是像普通请求那样立即关闭连接。
2.2 数据格式:看似简单,暗藏玄机
服务器端发送的数据,必须遵循一个简单的文本格式。每一段消息由若干行组成,行与行之间用换行符(\n)分隔。关键的行类型有四种:
data::这是消息的主体行。一行data:后面跟着实际的数据内容。如果消息内容很长,可以分成多行data:,它们最终会被连接成一个字符串。data: 这是一条消息的第一部分\n data: 这是第二部分\n客户端收到后,
EventSource的onmessage事件会接收到一个完整的数据:这是一条消息的第一部分\n这是第二部分。event::用于指定事件类型。这允许你在一个连接上发送多种不同类型的消息。客户端可以为不同的事件类型绑定不同的监听器。event: statusUpdate\n data: {"progress": 50}\n \n event: logMessage\n data: 任务执行到第二步\n \n客户端可以通过
addEventListener('statusUpdate', ...)和addEventListener('logMessage', ...)来分别处理。id::用于设置消息的ID。这个ID有两个重要作用:一是客户端在连接断开后重连时,可以通过Last-Event-ID请求头将这个ID发送给服务器,请求从断点之后的数据开始发送,实现断点续传;二是作为一种客户端收到消息的确认机制。retry::用于建议客户端在连接断开后,等待多少毫秒再进行重连。单位是毫秒。例如retry: 10000表示建议10秒后重试。
一个完整的消息块以一个空行(即连续两个换行符\n\n)结束。EventSourceAPI在收到这个空行后,才会触发对应的事件,将累积的data行数据传递给应用层。
注意:这个格式是强制的。如果你在
data:行前不小心多了一个空格,或者忘了在消息结束时发送空行,客户端的EventSource可能无法正确解析,导致消息无法触发。我在早期调试时,经常用curl命令直接请求SSE端点,观察原始的、未经解析的流数据,这是排查格式问题最有效的方法。
2.3 与WebSocket的核心差异:单向 vs. 双向
这是面试常考题,也是技术选型的决策点。我们不能简单说谁好谁坏,而要看场景。
| 特性 | Server-Sent Events (SSE) | WebSocket |
|---|---|---|
| 通信方向 | 单向(服务器 -> 客户端) | 全双工(服务器 <-> 客户端) |
| 协议基础 | HTTP(长连接) | 独立的ws://或wss://协议,需要握手升级 |
| 浏览器支持 | 除IE外的主流浏览器原生支持 | 广泛支持 |
| 数据格式 | 纯文本,UTF-8编码。格式固定(data:,event:等) | 二进制帧或文本帧,格式完全自定义 |
| 自动重连 | 原生支持。EventSource自动处理连接断开和重试 | 需要手动实现心跳和重连逻辑 |
| 断点续传 | 原生支持。通过Last-Event-ID头实现 | 需要应用层协议自行设计 |
| 复杂度 | 低。基于HTTP,无需额外库,服务器端逻辑简单 | 中高。需要处理握手、帧解析、心跳等 |
| 适用场景 | 实时通知、新闻推送、股票行情、任务进度、日志流、AI流式响应 | 聊天室、在线游戏、实时协作、双向数据同步 |
选择SSE的关键时刻:当你确定数据流主要是服务器向客户端的单向推送,并且你希望实现成本最低、利用HTTP现有生态(如身份认证、缓存、代理兼容性)、且需要自动重连和断点续传时,SSE是更优雅的选择。Spring AI的流式输出选择SSE,正是看中了其与HTTP服务无缝集成、前端接入简单的特点。
3. 服务端实现:以Spring Boot为例的深度实践
理论说再多,不如一行代码。我们以最常用的Java Spring Boot为例,看看如何构建一个健壮、高效的SSE服务端。这里我会分享两种主流方式,并深入生产环境中必须考虑的细节。
3.1 方式一:使用SseEmitter(推荐用于Spring MVC)
SseEmitter是Spring框架专门为SSE提供的一个抽象,它封装了底层的HTTP响应流,让我们的代码更专注于业务逻辑。
基础实现:
@RestController @RequestMapping("/api/sse") public class SseController { // 用于保存所有客户端的连接,实现广播等功能 private final Map<String, SseEmitter> emitters = new ConcurrentHashMap<>(); @GetMapping(path = "/subscribe", produces = MediaType.TEXT_EVENT_STREAM_VALUE) public SseEmitter subscribe(@RequestParam String clientId) { // 1. 创建SseEmitter,设置超时时间(0表示永不超时,生产环境建议设置) SseEmitter emitter = new SseEmitter(30 * 60 * 1000L); // 30分钟超时 // 2. 保存连接 emitters.put(clientId, emitter); // 3. 设置回调:连接完成或出错时清理资源 emitter.onCompletion(() -> { log.info("Client {} SSE connection completed.", clientId); emitters.remove(clientId); }); emitter.onTimeout(() -> { log.warn("Client {} SSE connection timed out.", clientId); emitter.complete(); emitters.remove(clientId); }); emitter.onError((ex) -> { log.error("Error on SSE connection for client {}: ", clientId, ex); emitters.remove(clientId); }); // 4. 发送初始连接成功消息(可选,但推荐) try { emitter.send(SseEmitter.event() .name("connect") // 事件类型 .data("Connected as: " + clientId) .id(UUID.randomUUID().toString()) // 初始ID .reconnectTime(5000L) // 建议重连时间 ); } catch (IOException e) { log.error("Failed to send initial message to client {}", clientId, e); emitter.completeWithError(e); } return emitter; } // 一个模拟向特定客户端推送进度的方法 @PostMapping("/progress/{clientId}") public void sendProgress(@PathVariable String clientId, @RequestParam int progress) { SseEmitter emitter = emitters.get(clientId); if (emitter != null) { try { emitter.send(SseEmitter.event() .name("progress") .data("{\"value\": " + progress + "}") .id(String.valueOf(System.currentTimeMillis())) // 用时间戳做ID ); } catch (IOException e) { log.warn("Client {} may be disconnected, failed to send progress.", clientId, e); // 发送失败通常意味着客户端已断开,清理资源 emitter.completeWithError(e); emitters.remove(clientId); } } } }关键点解析与避坑:
超时设置:
SseEmitter构造函数可以传入超时时间(毫秒)。生产环境切忌设置为0(无限长)。因为TCP连接可能因为网络波动而僵死,但应用层不知情。设置一个合理的超时(如30分钟),可以让僵死的连接被及时清理,释放服务器资源(如线程、文件描述符)。超时后,会触发onTimeout回调。资源清理:
onCompletion、onTimeout、onError这三个回调是必须设置的。它们的核心任务就是从你的连接池(如上面的emittersMap)中移除对应的SseEmitter。如果只创建不清理,会导致内存泄漏,连接数稍微一多,服务器就可能被拖垮。异常处理:
emitter.send()方法会抛出IOException。这个异常在绝大多数情况下,都意味着客户端已经主动断开连接(比如用户关闭了浏览器标签)。此时,你应该捕获这个异常,调用emitter.complete()或emitter.completeWithError(e)来正式结束这个发射器,并执行资源清理。不要尝试重发,因为连接已经不存在了。数据格式与编码:
data()方法默认发送文本。如果你想发送JSON,直接传入JSON字符串即可,前端EventSource会以字符串形式接收,你需要手动JSON.parse()。确保字符串是UTF-8编码。发送二进制数据不是SSE的标准用法,通常需要Base64编码后以文本形式发送。
3.2 方式二:响应式编程(WebFlux)
如果你的项目使用的是Spring WebFlux,那么实现SSE会更加直观和函数式,因为它本身就是为处理流数据而设计的。
@RestController @RequestMapping("/api/sse-flux") public class SseFluxController { @GetMapping(path = "/stream", produces = MediaType.TEXT_EVENT_STREAM_VALUE) public Flux<ServerSentEvent<Object>> streamEvents() { // 模拟一个每秒钟发送一个数字的无限流 return Flux.interval(Duration.ofSeconds(1)) .map(sequence -> ServerSentEvent.builder() .id(String.valueOf(sequence)) // ID .event("periodic-event") // 事件类型 .data("Current count is: " + sequence) // 数据 .build() ); } // 更复杂的例子:集成Spring AI的流式响应 // 假设有一个AiService,其streamResponse方法返回Flux<String> @GetMapping(path = "/ai-chat", produces = MediaType.TEXT_EVENT_STREAM_VALUE) public Flux<ServerSentEvent<String>> chatStream(@RequestParam String question) { return aiService.streamResponse(question) .map(chunk -> ServerSentEvent.builder(chunk).build()) .onErrorResume(e -> { log.error("AI stream error", e); return Flux.just(ServerSentEvent.builder("[ERROR] Stream interrupted").build()); }); } }WebFlux的方式非常简洁,它利用Flux这个响应式流,天然地支持背压(Backpressure),能更好地处理大量并发连接。返回类型Flux<ServerSentEvent>会被Spring自动转换为符合SSE格式的HTTP流。
3.3 生产环境进阶考量
连接管理:上面的例子用了简单的
ConcurrentHashMap。在生产中,你可能需要更强大的结构,比如支持按主题订阅/发布的机制,或者使用Redis等外部存储来支持集群部署下的连接管理(因为HTTP连接是绑定到具体服务器实例的)。心跳机制:虽然SSE连接是长连接,但一些代理服务器或负载均衡器(如Nginx)可能会因为长时间没有数据传输而切断连接。为了防止这种情况,一个常见的做法是定期从服务器发送“心跳”消息。这可以是一个特定事件类型的空数据消息。
// 在创建emitter后,启动一个心跳任务 ScheduledExecutorService scheduler = Executors.newScheduledThreadPool(1); scheduler.scheduleAtFixedRate(() -> { try { emitter.send(SseEmitter.event().name("heartbeat").data("ping").id("hb")); } catch (IOException ignored) { // 发送失败,连接可能已断,任务会被取消 scheduler.shutdown(); } }, 0, 25, TimeUnit.SECONDS); // 每25秒发送一次注意:心跳间隔要小于你的负载均衡器或操作系统的TCP Keep-Alive超时时间。通常设置为20-30秒比较安全。
背压处理:对于
SseEmitter,如果客户端网络很慢,而服务器发送数据很快,可能会导致服务器端数据积压,最终内存溢出。Spring的SseEmitter本身对背压的支持有限。一个缓解策略是,在发送前检查emitter的状态,或者使用有界队列来暂存待发送的消息。而在WebFlux中,背压是原生支持的,Flux会根据下游(客户端)的消费能力来调整数据推送速率。
4. 客户端实现:从前端到移动端的全栈指南
服务端准备好了,客户端如何优雅地接入呢?这里覆盖从浏览器到移动端(如React Native/Flutter)的关键实践。
4.1 浏览器原生EventSourceAPI
这是最简单直接的方式,现代浏览器(除IE)都支持。
// 创建EventSource连接 const eventSource = new EventSource('http://your-api.com/api/sse/subscribe?clientId=frontend_001'); // 监听通用消息(未指定event类型,或event为'message') eventSource.onmessage = (event) => { console.log('Generic message:', event.data); const data = JSON.parse(event.data); // 如果数据是JSON updateUI(data); }; // 监听特定事件类型的消息 eventSource.addEventListener('progress', (event) => { console.log('Progress update:', event.data); updateProgressBar(JSON.parse(event.data).value); }); eventSource.addEventListener('heartbeat', (event) => { console.log('Heartbeat received:', event.data); // 可以在这里重置一个连接健康状态的计时器 }); // 监听连接打开事件 eventSource.onopen = (event) => { console.log('SSE Connection opened.'); }; // 监听错误事件 eventSource.onerror = (event) => { console.error('SSE Error:', event); // 根据eventSource.readyState判断状态 // 0: CONNECTING, 1: OPEN, 2: CLOSED if (eventSource.readyState === EventSource.CLOSED) { console.log('Connection was closed.'); // 可以在这里尝试手动重连 // setTimeout(() => connectSSE(), 5000); } }; // 主动关闭连接 function closeConnection() { eventSource.close(); console.log('SSE connection closed by client.'); }客户端避坑要点:
- 跨域问题:
EventSource默认遵守同源策略。如果你的前端和后端不在同一个域名下,服务器端必须设置正确的CORS头,例如Access-Control-Allow-Origin: *或你的前端域名。对于携带身份认证(如Cookie)的请求,可能还需要设置Access-Control-Allow-Credentials: true,并且前端在创建EventSource时可能需要额外配置(但标准API不支持,通常需要服务器支持预检请求或使用其他方式认证,如URL token)。 - 错误处理:
onerror事件会被多次触发,例如网络错误、服务器返回非200状态码、或者流格式错误。readyState属性是判断当前连接状态的关键。不要一发生错误就立即无脑重连,这可能会对故障中的服务器造成雪崩。合理的策略是采用“指数退避”重试,例如第一次等1秒,第二次等2秒,第三次等4秒,以此类推。 - 数据解析:
event.data永远是字符串。如果服务器发送的是JSON,你需要手动调用JSON.parse()。记得用try...catch包裹,防止格式错误导致整个脚本崩溃。
4.2 移动端与网络不稳定的挑战
文章开头热词里提到了一个非常经典的问题:“移动端如何让SSE请求在网络短时断开重连后可以继续获取数据”。移动网络(4G/5G/WIFI切换、进出电梯、隧道)的不稳定性是SSE面临的一大挑战。
核心诉求:断点续传。这正是SSE协议中id:字段和Last-Event-ID头的用武之地。
服务端配合实现:
- 服务器在发送每条重要消息时,都附带一个递增的、唯一的
id。 - 客户端连接时,检查本地存储(如
localStorage或AsyncStorage)中是否保存了上一次收到的最新消息ID。 - 在创建新的
EventSource连接时,如果浏览器原生API支持(部分浏览器可能需要polyfill),可以在URL中携带?lastEventId=xxx参数。更通用的做法是,服务器端在连接建立时,读取客户端通过Cookie或首次握手请求体传递的lastEventId。 - 服务器收到
lastEventId后,从该ID之后的消息开始推送。
一个增强型客户端的思路(伪代码):由于原生EventSource在网络断开后自动重连时,会自动发送Last-Event-ID头,但移动端短时断开(如3秒)可能被迅速重连,而服务器端可能还没清理旧连接,导致数据混乱。因此,我们需要更精细的控制。
// 使用一个更可控的SSE库,或者用Fetch API模拟SSE class RobustSSEClient { constructor(url) { this.url = url; this.lastEventId = localStorage.getItem('last_sse_id') || null; this.reconnectDelay = 1000; // 初始重连延迟 this.maxReconnectDelay = 30000; // 最大重连延迟 this.connect(); } async connect() { try { const headers = {}; if (this.lastEventId) { headers['Last-Event-ID'] = this.lastEventId; } const response = await fetch(this.url, { headers }); if (!response.ok || !response.body) throw new Error('Connect failed'); const reader = response.body.getReader(); const decoder = new TextDecoder('utf-8'); let buffer = ''; while (true) { const { done, value } = await reader.read(); if (done) break; buffer += decoder.decode(value, { stream: true }); const lines = buffer.split('\n'); buffer = lines.pop(); // 最后一行可能是不完整的,放回buffer for (const line of lines) { this.parseSSELine(line); } } } catch (error) { console.warn('SSE连接断开,准备重连...', error); this.scheduleReconnect(); } } parseSSELine(line) { // 简化的SSE行解析逻辑 if (line.startsWith('id:')) { this.lastEventId = line.substring(3).trim(); localStorage.setItem('last_sse_id', this.lastEventId); } else if (line.startsWith('data:')) { const data = line.substring(5).trim(); if (data) { this.onMessage(data); } } // 忽略event:, retry: 或空行解析... } onMessage(data) { // 处理业务消息 console.log('Received:', data); } scheduleReconnect() { // 指数退避策略 const delay = Math.min(this.reconnectDelay, this.maxReconnectDelay); console.log(`将在 ${delay}ms 后重连`); setTimeout(() => { this.reconnectDelay *= 2; // 延迟加倍 this.connect(); }, delay); } // 连接成功后,重置重连延迟 resetReconnectDelay() { this.reconnectDelay = 1000; } }这种用Fetch API读取流并手动解析的方式,给了我们最大的控制权,可以自定义重连逻辑、断点续传策略和错误处理。对于React Native或Flutter,原理类似,可以使用fetch或dio(Flutter)库来读取流响应,并实现相应的解析器。
5. 实战问题排查与性能优化
理论、代码都齐了,但在真实的生产环境跑起来,总会遇到一些“坑”。下面是我总结的几个典型问题和优化点。
5.1 常见问题与解决方案速查表
| 问题现象 | 可能原因 | 排查步骤与解决方案 |
|---|---|---|
| 前端收不到任何消息 | 1. 连接未成功建立。 2. 服务器响应格式错误。 3. 跨域问题。 | 1. 浏览器开发者工具 -> Network,查看SSE请求状态码是否为200,响应头Content-Type是否为text/event-stream。2. 用 curl -N <url>直接请求接口,查看原始输出格式是否正确(每段消息以空行结束)。3. 检查Console是否有CORS错误,确保服务器设置了正确的CORS头。 |
| 消息延迟很高或时断时续 | 1. 服务器端处理阻塞。 2. 网络问题或代理缓冲。 3. 客户端事件循环阻塞。 | 1. 检查服务器端发送消息的代码是否在同步阻塞操作(如复杂计算、同步IO)。考虑异步非阻塞发送。 2. 检查Nginx等代理配置,确保 proxy_buffering off;对于SSE流,缓冲会导致数据积压然后一次性吐出。3. 检查前端JavaScript是否在执行耗时操作,阻塞了 onmessage回调。 |
| 连接频繁断开重连 | 1. 服务器或代理超时设置过短。 2. 防火墙或负载均衡器切断空闲连接。 3. 移动网络抖动。 | 1. 调整服务器端SseEmitter超时时间(如30分钟)。检查Nginx的proxy_read_timeout(建议设置很大,如1小时)。2. 引入服务器端心跳机制,定期发送注释消息( :开头的行)或空事件,保持TCP连接活跃。3. 客户端实现指数退避重连,避免重连风暴。 |
| 内存泄漏(服务器端) | SseEmitter实例未被正确清理。 | 1.务必实现onCompletion,onTimeout,onError回调,并从连接管理Map中移除失效的emitter。2. 定期扫描Map,清理长时间未活动(如未收到心跳确认)的连接。 |
| 移动端退后台后连接断开 | 操作系统为省电可能暂停网络请求。 | 1. 这是正常行为。应用回到前台时,需要检测连接状态并重新建立SSE连接,并携带Last-Event-ID尝试续传。2. 对于关键实时性要求不高的数据,可以考虑退而使用短轮询或WebSocket(部分场景下保活能力更强)。 |
5.2 性能优化与最佳实践
连接数限制:一个浏览器标签页对同一个域名有并发连接数限制(通常是6个)。SSE长连接会占用其中一个。如果你的页面同时有多个SSE连接或其他HTTP长连接需求,需要注意这个限制。可以考虑将多个逻辑流合并到一个SSE连接中,通过不同的
event类型来区分。数据量优化:SSE传输的是文本。对于频繁更新且数据量大的场景(如每秒推送大量日志),要考虑数据压缩。虽然HTTP层可以开启gzip,但对于流式响应,压缩效果和方式需要测试。更有效的方法是精简数据格式,只发送变化的部分(增量更新)。
服务端资源管理:每个SSE连接在服务器端都会持有一个线程(Tomcat等传统Servlet容器)或一个反应式订阅(WebFlux)。在连接数很高时(比如上万),对服务器资源是巨大考验。使用WebFlux等非阻塞框架能显著提升并发能力。此外,要做好监控,关注活跃连接数、内存使用情况。
安全与认证:SSE基于HTTP,因此可以使用所有HTTP标准的认证方式,如Cookie、JWT Token(通常放在URL参数或自定义Header,但原生
EventSource不支持设置Header,需使用Fetch模拟或后端从Cookie读取)。务必注意,如果Token放在URL中,可能会被日志记录,存在泄漏风险。与Spring AI等框架集成:这是当前的热门应用。Spring AI的流式Chat模型接口通常会返回一个
Flux<ChatResponse>,你只需要像前面WebFlux例子中那样,将其映射为ServerSentEvent流即可。关键在于处理好流中的每个“块”(chunk),并确保在发生错误时能优雅地关闭流并发送错误信息给前端,而不是让连接直接挂断。
SSE是一个在特定场景下极其优雅和高效的解决方案。它复用了HTTP的广阔生态,降低了开发和运维的复杂度。下次当你面临服务器向客户端单向推送数据的场景时,不妨先问问自己:“用SSE是不是更简单?” 很多时候,答案都是肯定的。