流式编程
约 1228 个字 76 行代码 预计阅读时间 5 分钟
SSE协议介绍
HTTP协议本身设计为无状态的请求-响应模式,严格来说,是无法做到服务器主动推送消息到客户端,但通过Server-Sent Events(服务器发送事件,简称SSE)技术可实现流式传输,允许服务器主动向浏览器推送数据流
也就是说,服务器向客户端声明,接下来要发送的是流消息(streaming),这时客户端不会关闭连接,会一直等待服务器发送过来新的数据流
SSE(Server-Sent Events)是一种基于HTTP的轻量级实时通信协议,浏览器通过内置的EventSource API接收并处理这些实时事件
核心特点:
| 特点 | 说明 |
| 基于HTTP协议 | 复用标准HTTP/HTTPS协议,无需额外端口或协议,兼容性好且易于部署 |
| 单向通信机制 | SSE仅支持服务器向客户端的单向数据推送,客户端通过普通HTTP请求建立连接后,服务器可持续发送数据流,但客户端无法通过同一连接向服务器发送数据 |
| 自动重连机制 | 支持断线重连,连接中断时,浏览器会自动尝试重新连接(支持retry字段指定重连间隔) |
| 定义消息类型 | 客户端发起请求后,服务器保持连接开放,响应头设置Content-Type: text/event-stream,标识为事件流格式,持续推送事件流 |
SSE协议数据格式
服务端向浏览器发送SSE数据,需要设置必要的HTTP头信息:
| Text Only |
|---|
| Content-Type: text/event-stream;charset=utf-8
Connection: keep-alive
|
需要注意的是,编码必须是utf-8
每一次发送的消息,由若干个message组成,每个message之间由\n\n分隔,每个message内部由若干行组成,每一行都是如下格式:
field可以取值为:
data(必需):数据内容 event(非必需):表示自定义的事件类型,默认是message事件 id(非必需):数据标识符,相当于每一条数据的编号 retry(非必需):指定浏览器重新发起连接的时间间隔
除此之外,还可以有冒号:开头的行,表示注释
例如下面的示例:
| Text Only |
|---|
| 第一条message:
event: foo\n
data: a foo event\n\n
第二条message:
data: an unnamed event\n\n
第三条message:
event: end\n
data: a bar event\n\n
|
Java后端发送SSE协议数据
原生格式
根据SSE协议的格式,在Java中可以自行构建SSE数据,例如下面的案例:
| Java |
|---|
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17 | @RestController
@RequestMapping("/sse")
public class SseController {
@RequestMapping("/data")
public void data(HttpServletResponse response) throws IOException, InterruptedException {
// 必须设置ContentType为text/event-stream;charset=utf-8
response.setContentType("text/event-stream;charset=utf-8");
PrintWriter writer = response.getWriter();
// 写10次内容:正计时10次
for (int i = 0; i < 10; i++) {
String time = "data: " + new Date() + "\n\n"; // \n\n表示一条消息结束
writer.write(time);
writer.flush();
Thread.sleep(1000L);
}
}
}
|
对应的,在前端编写用浏览器内置对象EventSource接收SSE数据的代码:
| HTML |
|---|
| <body>
<div id="content"></div>
</body>
<script>
let eventSource = new EventSource("/sse/data"); // 此处参数为请求接口地址
eventSource.onmessage = (event) => {
document.getElementById("content").innerText = event.data;
};
</script>
|
默认情况下,SSE协议是会进行自动重连的,也就是说,如果没有手动关闭,在一次请求之后会再一次发送请求。如果需要指定重连时间,可以通过retry来指定具体的时间(单位为毫秒)。如果需要关闭连接,可以使用自定义SSE事件,例如发送下面的数据:
| Text Only |
|---|
| 使用自定义事件标记开始:
event: start\n
data: 时间\n\n
使用自定义事件标记结束:
event: end\n
data: end\n\n
|
对应的后端代码如下:
| Java |
|---|
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20 | @RequestMapping("/event")
public void event(HttpServletResponse response) throws IOException, InterruptedException {
// 必须设置ContentType为text/event-stream;charset=utf-8
response.setContentType("text/event-stream;charset=utf-8");
PrintWriter writer = response.getWriter();
// 自定义事件
// 写10次内容:正计时10次
for (int i = 0; i < 10; i++) {
String e = "event: start\n";
e += "data: " + new Date() + "\n\n"; // \n\n表示一条消息结束
writer.write(e);
writer.flush();
Thread.sleep(1000L);
}
// 10次之后不要重连
String end = "event: end\ndata: end\n\n";
writer.write(end);
writer.flush();
}
|
前端代码中就需要针对自定义事件进行处理,默认的是message事件,所以有默认的onmessage方法,对于自定义事件,需要使用addEventListener:
| HTML |
|---|
1
2
3
4
5
6
7
8
9
10
11
12
13 | <body>
<div id="content"></div>
</body>
<script>
// 自定义事件
let eventSource = new EventSource("/sse/event");
eventSource.addEventListener("start", (event)=>{
document.getElementById("content").innerText = event.data;
});
eventSource.addEventListener("end", ()=>{
eventSource.close(); // 关闭连接
});
</script>
|
Spring与SSE
Spring 4.2开始就已经支持SSE,从Spring 5开始我们可以使用WebFlux更优雅的实现SSE协议。Flux是WebFlux的核心API
可以把Flux想象成一条传送带:
- 异步传送:数据像快递包裹一样逐个到达,不需等全部到齐
- 灵活加工:支持中途修改数据(如过滤/转换)
- 弹性控制:接收方可以调速(背压机制)
Flux的流程分为三个步骤:
- 创建Flux:创建一个Flux数据流,并有数据源
- 处理数据:使用操作符对数据进行处理
- 订阅数据:订阅Flux来消费数据,触发数据的流动
创建Flux可以使用到just()方法创建一个包含指定元素的Flux:
| Java |
|---|
| Flux<String> fruitFlux = Flux.just("Apple", "Banana", "Cherry");
|
也有其他的创建方式,如fromIterable、range、interval:
| Java |
|---|
| // 从集合(如List、Set)创建
Flux.fromIterable(Arrays.asList(1, 2, 3))
// 生成1~5的Flux
Flux.range(1, 5)
// 每隔1秒生成一次数据
Flux.interval(Duration.ofSeconds(1))
|
处理数据有以下几种常见的操作符:
| 操作符 | 作用 | 示例代码 |
map() | 元素一对一转换 | fruitFlux.map(String::toUpperCase) |
filter() | 条件过滤 | .filter(s -> s.startsWith("A")) |
take() | 限制元素数量 | .take(2) |
merge() | 合并多个Flux(不保证顺序) | Flux.merge(Flux.just("A"), Flux.just("B")) |
concat() | 顺序拼接多个Flux(保证顺序) | Flux.concat(Flux.just("A"), Flux.just("B")) |
delayElements() | 延迟元素发射 | .delayElements(Duration.ofSeconds(1)) |
订阅数据以实现数据的流动,使用subscribe方法实现,例如:
| Java |
|---|
| Flux<String> fruitFlux = Flux.just("Apple", "Banana", "Cherry");
Flux<String> newFlux = fruitFlux.map(String::toUpperCase) // 将每个字符串转为大写
.filter(s -> s.startsWith("A")); // 筛选出以A开头的元素
newFlux.subscribe(System.out::println);
|
对于前面每隔1秒返回一次时间的后端代码可以修改为如下:
| Java |
|---|
| @RequestMapping(value = "/flux", produces = MediaType.TEXT_EVENT_STREAM_VALUE)
public Flux<String> flux(HttpServletResponse response) throws IOException, InterruptedException {
// 每隔1秒生成当前时间数据
return Flux.interval(Duration.ofSeconds(1)).map(s -> new Date().toString());
}
|
因为默认是message事件,所以前端直接调用onmessage方法即可