发布于2026-07-17 阅读(0)
扫一扫,手机访问
说起来,现在的互联网应用开发,几乎都在跟实时数据较劲。从股票交易系统里实时跳动的股价,到社交平台上瞬间送达的私信,再到物联网场景下传感器数据的连续回传,哪一样离得开后端的实时响应能力?传统的轮询机制虽然简单粗暴,但效率低、资源消耗大,就像隔几分钟去问一遍“有没有新消息”,既笨拙又浪费。而WebSocket虽然功能全面,但在某些场景下又显得过于沉重,实现成本不低。在这种背景下,Server-Sent Events(SSE)作为一个轻量级、基于HTTP协议的单向实时通信方案,逐渐走进了开发者的视野。
SSE的核心思路是:服务器主动往客户端“扔”数据,不用客户端一遍遍地主动问“有新消息吗?”。这不仅提升了数据传输效率,也减轻了服务器的负担。更关键的是,它基于标准的HTTP协议,实现起来相当简单,不需要额外的协议支持。而Spring Boot,作为目前Ja va微服务领域最主流的框架之一,在2.7.8版本中已经提供了对SSE的原生支持,能让我们轻松地搭建起实时数据推送服务。

这篇文章会从零开始,一步步带你搭建一个基于Spring Boot 2.7.8的SSE服务端应用。我们会从环境搭建、引入依赖这些基础工作开始,然后深入探讨SSE的核心概念,比如事件流的格式、数据推送机制,以及如何处理客户端的连接和重连。通过具体的代码示例,你会看到如何在Spring Boot里配置和使用SSE,实现从服务器到客户端的实时数据推送。此外,我们还会聊聊一些常见的坑和挑战,比如如何保证数据的实时性和准确性,高并发场景下怎么优化性能,以及如何确保服务的稳定可靠。无论你是刚接触Spring Boot的新手,还是有一定经验的开发者,希望这篇文章能帮你掌握SSE的关键要点,并最终运用到自己的项目中去。
在正式动手实现之前,先给刚接触的朋友打个底,简单介绍一下SSE的机制。我们主要讲三件事:SSE是什么、它的工作原理是什么,以及它适合用在哪些场景。
SSE(Server-Sent Events)是HTML 5规范的一部分,官方有详细的文档。这个规范本身比较简单,主要包括两部分:一是服务器端与浏览器端之间的通讯协议,二是浏览器端供Ja vaScript使用的EventSource对象。通讯协议基于纯文本,服务器的响应内容类型被设定为“text/event-stream”。整个响应文本可以看作一个事件流,由不同的事件组成。每个事件包含类型和数据两部分,还可以有一个可选的标识符。不同事件之间通过一个仅包含回车和换行的空行(“\r\n”)来分隔。每个事件的数据可以写成多行。
客户端发起请求:客户端通过EventSource API向服务器发起一个HTTP GET请求,请求头里带上Accept: text/event-stream,表明自己希望接收事件流。
服务器响应:服务器收到请求后,保持连接不断开,并设置响应头Content-Type: text/event-stream和Cache-Control: no-cache,确保数据流不会被缓存。
数据推送:服务器通过这个保持开放的连接,以事件流的形式向客户端发送数据。每个事件由几个字段组成,比如data(消息内容)、event(事件类型)、id(消息编号)和retry(重连间隔)。
客户端接收:客户端通过监听事件流来获取数据,并在收到事件后进行处理。
自动重连:如果连接中途断开了,客户端会根据retry字段的值自动尝试重新连接。
说到SSE的使用场景,就不得不提它的老对手——WebSocket。WebSocket是双向全双工的通道,可以同时收发消息,而且在HTTPS安全方面的处理也比较严格。但在某些场景下,我们其实并不需要这么复杂的双向通信,只需要被动地接收服务器推送的信息就行。整理一下,SSE适合的场景大概有这些:
实时通知:比如新闻更新、消息提醒、股票价格变动,服务器可以实时向客户端推送最新信息。
流式数据:比如日志流、传感器数据,客户端可以持续接收服务器送来的数据流。
单向通信场景:当只需要服务器向客户端推送数据,客户端不需要回传数据时,SSE是个简单高效的选择。
讲完这些基础知识,下面我们就正式进入正题,看看在Spring Boot中如何实现SSE服务。
Spring Boot 2.7.8版本对SSE提供了原生支持,实现起来相当方便。这一节我们会从Ma ven依赖引入、SSE服务类实现、SSE控制器类和基于Thymeleaf的简单页面实现这四个方面来展开介绍。
先给出一个比较简单的Pom依赖示例,关键代码如下:
4.0.0 org.yelang baidu-sse-client 0.0.1-SNAPSHOT org.springframework.boot spring-boot-starter-parent 2.7.18 1.8 UTF-8 UTF-8 org.springframework.boot spring-boot-starter-webflux org.springframework.boot spring-boot-starter-web com.fasterxml.jackson.core jackson-databind org.springframework.boot spring-boot-starter-thymeleaf org.springframework.boot spring-boot-starter-test test io.projectreactor reactor-test test org.springframework.boot spring-boot-ma ven-plugin org.apache.ma ven.plugins ma ven-compiler-plugin 1.8 1.8
引入依赖之后,接下来就是重头戏——SSE服务类的实现,这是整个SSE的核心。它不仅要负责连接的创建和销毁,还要处理消息的发送,包括群发和单点发送。下面分别来看这些功能的具体实现,核心代码如下:
package org.yelang.service;
import org.springframework.stereotype.Service;
import org.springframework.web.servlet.mvc.method.annotation.SseEmitter;
import ja va.io.IOException;
import ja va.util.Map;
import ja va.util.concurrent.ConcurrentHashMap;
import ja va.util.concurrent.atomic.AtomicInteger;
@Service
public class SseService {
// 保存所有连接的 emitter
private final Map emitters = new ConcurrentHashMap<>();
private final AtomicInteger counter = new AtomicInteger(0);
/**
* -创建新的 SSE 连接
*/
public SseEmitter createEmitter(String clientId) {
// 设置超时时间(0表示永不超时)
SseEmitter emitter = new SseEmitter(0L);
emitters.put(clientId, emitter);
// 设置完成和超时回调
emitter.onCompletion(() -> {
emitters.remove(clientId);
System.out.println("SSE连接完成: " + clientId);
});
emitter.onTimeout(() -> {
emitters.remove(clientId);
System.out.println("SSE连接超时: " + clientId);
});
emitter.onError((e) -> {
emitters.remove(clientId);
System.out.println("SSE连接错误: " + clientId + ", 错误: " + e.getMessage());
});
return emitter;
}
/**
* -发送消息给所有客户端
*/
public void sendToAll(String message) {
emitters.forEach((clientId, emitter) -> {
try {
SseEmitter.SseEventBuilder event = SseEmitter.event().data(message)
.id(String.valueOf(counter.incrementAndGet())).name("message").reconnectTime(5000L);
emitter.send(event);
} catch (IOException e) {
emitter.completeWithError(e);
emitters.remove(clientId);
}
});
}
/**
* -发送消息给特定客户端
*/
public void sendToClient(String clientId, String message) {
SseEmitter emitter = emitters.get(clientId);
if (emitter != null) {
try {
SseEmitter.SseEventBuilder event = SseEmitter.event().data(message)
.id(String.valueOf(counter.incrementAndGet())).name("message");
emitter.send(event);
} catch (IOException e) {
emitter.completeWithError(e);
emitters.remove(clientId);
}
}
}
/**
* -获取当前连接数
*/
public int getConnectionCount() {
return emitters.size();
}
}
这里为了演示方便,连接时间设置成了永不超时。同时,为了方便统一管理各个连接,我们用了一个HashMap来保存。群发和私发的区别在于,群发是向所有客户端广播消息,而私发则是只有通信双方才知道的悄悄话。
和常规的MVC应用一样,后台也需要一个SSE控制器,才能为前端页面提供连接和消息发布的能力。这里我们提供了以下几个方法:
| 序号 | 方法名 | 参数说明 |
| 1 | public String index(Model model) | 跳转SSE管理首页 |
| 2 | public SseEmitter streamSse(@RequestParam(value = "clientId", required = false) String clientId) | 建立 SSE 连接 |
| 3 | public String broadcastMessage(@RequestParam String message) | 广播消息 |
| 4 | public String sendToClient(@RequestParam String clientId, @RequestParam String message) | 点对点发送消息 |
| 5 | public int getConnectionCount() | 获取SSE连接数 |
下面给出实例代码:
package org.yelang.controller;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.http.MediaType;
import org.springframework.stereotype.Controller;
import org.springframework.ui.Model;
import org.springframework.web.bind.annotation.*;
import org.springframework.web.servlet.mvc.method.annotation.SseEmitter;
import org.yelang.service.SseService;
import ja va.util.UUID;
@Controller
@RequestMapping("/sseman")
public class SseController {
@Autowired
private SseService sseService;
/**
* -首页
*/
@GetMapping("/index")
public String index(Model model) {
model.addAttribute("connectionCount", sseService.getConnectionCount());
return "sse/index";
}
/**
* -建立 SSE 连接
*/
@GetMapping(value = "/sse", produces = MediaType.TEXT_EVENT_STREAM_VALUE)
public SseEmitter streamSse(@RequestParam(value = "clientId", required = false) String clientId) {
if (clientId == null || clientId.trim().isEmpty()) {
clientId = UUID.randomUUID().toString();
}
return sseService.createEmitter(clientId);
}
/**
* -发送消息给所有客户端
*/
@PostMapping("/broadcast")
@ResponseBody
public String broadcastMessage(@RequestParam String message) {
sseService.sendToAll(message);
return "消息已广播";
}
/**
* - 发送消息给特定客户端
*/
@PostMapping("/send-to-client")
@ResponseBody
public String sendToClient(@RequestParam String clientId, @RequestParam String message) {
sseService.sendToClient(clientId, message);
return "消息已发送给客户端: " + clientId;
}
/**
* -获取连接数
*/
@GetMapping("/connection-count")
@ResponseBody
public int getConnectionCount() {
return sseService.getConnectionCount();
}
}
下面以Thymeleaf为例,重点讲解前端页面如何集成SSE。大家可以根据实际业务需要来调整。首先,定义一下页面样式:
页面区域分为连接区、发送区和消息展示区,我们提供了两个面板来分别处理数据收集和事件发送。
连接状态
当前连接数:
0
我的客户端ID:
未连接
发送消息
接收消息
为了方便页面标识,我们创建一个生成clientId的方法,使用随机数生成,参考代码如下:
function generateClientId() {
return 'client_' + Math.random().toString(36).substr(2, 9);
}
在HTML中创建SSE连接并与后台连接的代码如下:
let eventSource = null;
let clientId = null;
var ctx = "/bdsse/sseman";
function connectSSE() {
if (eventSource) {
addMessage('警告', '已经连接到SSE服务器');
return;
}
// 生成客户端ID
clientId = generateClientId();
document.getElementById('clientId').textContent = clientId;
// 创建 EventSource 连接
eventSource = new EventSource(ctx + '/sse?clientId=' + clientId);
// 处理消息事件
eventSource.onmessage = function(event) {
addMessage('服务器消息', event.data);
};
// 处理自定义事件
eventSource.addEventListener('message', function(event) {
addMessage('自定义消息', event.data);
});
// 处理连接打开
eventSource.onopen = function(event) {
addMessage('系统', 'SSE连接已建立');
refreshConnectionCount();
};
// 处理错误
eventSource.onerror = function(event) {
if (eventSource.readyState === EventSource.CLOSED) {
addMessage('系统', 'SSE连接已关闭');
} else {
addMessage('错误', 'SSE连接错误: ' + event);
}
};
addMessage('系统', '正在连接SSE服务器...');
}
通过这种方式,前端就可以请求后台接口,成功连接上SSE服务端。如果想断开连接,可以调用以下方法:
function disconnectSSE() {
console.log("断开连接");
if (eventSource) {
eventSource.close();
eventSource = null;
addMessage('系统', 'SSE连接已断开');
refreshConnectionCount();
} else {
addMessage('警告', '没有活动的SSE连接');
}
}
更多具体的应用和代码,我们会在下一节中详细讲解。
这一节,我们结合具体的页面和SSE的处理方法,来演示一个实际的消息发送与接收的完整流程。
在控制台启动服务后,在浏览器中输入访问地址,就能看到下面的界面:

可以看到,消息栏已经显示连接SSE服务器成功。接下来,我们就可以进行消息广播和点对点发送了。
群发消息很好理解,就是通知所有连接的客户端,同时在客户端显示发送的消息。群发消息的处理代码如下:
function broadcastMessage() {
const message = document.getElementById('broadcastMessage').value;
if (!message) {
alert('请输入消息');
return;
}
fetch(ctx + '/broadcast', {
method: 'POST',
headers: {
'Content-Type': 'application/x-www-form-urlencoded',
},
body: 'message=' + encodeURIComponent(message)
})
.then(response => response.text())
.then(data => {
addMessage('操作', data);
document.getElementById('broadcastMessage').value = '';
})
.catch(error => {
addMessage('错误', '发送失败: ' + error);
});
}
为了方便演示,我们打开两个标签页,让两个客户端都连上这个SSE服务,界面如下:

在任意一个客户端中输入需要群发的消息,比如“hello world,大家好”,然后点击“广播给所有客户端”按钮,再来看看每个客户端收到的内容:

使用SSE除了可以群发消息,也可以向指定客户端发送消息,也就是点对点私发。实现代码如下:
function sendPrivateMessage() {
const targetClientId = document.getElementById('targetClientId').value;
const message = document.getElementById('privateMessage').value;
if (!targetClientId || !message) {
alert('请输入客户端ID和消息');
return;
}
fetch(ctx + '/send-to-client', {
method: 'POST',
headers: {
'Content-Type': 'application/x-www-form-urlencoded',
},
body: 'clientId=' + encodeURIComponent(targetClientId) +
'&message=' + encodeURIComponent(message)
})
.then(response => response.text())
.then(data => {
addMessage('操作', data);
document.getElementById('privateMessage').value = '';
})
.catch(error => {
addMessage('错误', '发送失败: ' + error);
});
}
记下目标客户端的clientID之后,就可以向这个客户端发送消息了。点击发送后,在接收端和发送端可以看到以下内容:

再来检查一下其他第三方客户端,看看它们能否收到消息:

从图上可以清楚地看到,目标客户端成功接收到了消息,而非目标客户端则没有收到,这就实现了点对点的私发消息。
以上就是本文的全部内容。我们从零开始,逐步构建了一个基于Spring Boot 2.7.8的SSE服务端应用。从环境搭建、引入依赖,到深入探讨SSE的核心概念,包括事件流格式、数据推送机制,以及如何处理客户端的连接和重连。通过具体的代码示例,你应该已经掌握了如何在Spring Boot中配置和使用SSE,实现从服务器到客户端的实时数据推送。行文仓促,定有不足之处,欢迎各位朋友在评论区批评指正,不胜感激。
售后无忧
立即购买>office旗舰店
售后无忧
立即购买>office旗舰店
售后无忧
立即购买>office旗舰店
售后无忧
立即购买>office旗舰店
正版软件
正版软件
正版软件
正版软件
正版软件
1
2
3
7
8