目录
一、引言
二、请求率与数据公式推导
三、WebSocket
1、引入WebSocket依赖
2、实现WebSocketConfiguration配置类
3、WebSocket两种消息
4、两个关键回调方法
5、广播群发
6、前端开启WebSocket连接
四、轮询
五、实验
六、总结
一、引言
假设有这样一个场景:在炒股时,股价不是固定不变的,是动态变化的。我们希望只要股价发生变动,我们能第一时间收到消息,并作出决策。那么我们就需要时刻关注股价变化。
但是,传统的网页都是浏览器主动发起请求,请求后端数据。用户如果不手动点击刷新,股价就不会变,因此造成不好的体验。
那么我们可以考虑使用轮询,前端每隔interval时间段就向后端查询一次,看看股价有没有更新。如果有更新,那么直接带着新数据返回,如果没有更新,那么这次请求就是无效的。轮询的工作由客户端做,但是如果用户切换到别的页面那么定时器可能不会正常工作。
同时我们还可以考虑使用WebSocket,这是一种服务器主动向客户端推送的消息推送机制。后端检测到新数据,那么就会向已经建立了WebScoket连接的客户端推送消息。WebSocket是一种网络通信协议,它允许在客户端(如浏览器)和服务器之间建立持久的、双向的、全双工的通信连接。传统实时访问就是基于上面的轮询来做的,如果用户希望使用WebSocket协议,那么就会和后端进行握手自动从http协议升级到WebSocket协议。
WebSocket有如下核心特点:
- 全双工通信:客户端和服务器可以同时互相发送数据,互不阻塞
- 持久连接:一次握手后,连接一直保持,直到一方主动关闭
- 低开销:握手阶段用HTTP,之后的数据帧头部很小,比HTTP每次携带完整头部高效得多
- 基于TCP:底层是TCP连接,保证可靠传输
- 协议表示:URL使用ws://(明文)或wss://(加密)传输
其工作原理可以如下简述:
- 握手:客户端发送一个带有Upgrade:websocket头得HTTP请求。服务器如果支持,返回101Switching Protocols,协议从HTTP升级为WebSocket
- 数据传输:握手成功后,双方通过数据帧通信,不在使用HTTP请求-响应模式
- 关闭:任意一方发送关闭帧,连接断开
二、请求率与数据公式推导
假设每一条新的数据每隔 t 时间产生,本次实验总的观察时间为T,数据从生产到使用的时效为L,轮询间隔为interval
那么,如果使用轮询,则interval 需要满足约束:
就需要轮询次,但是成功的只有
次,其他的都是无效请求
如果使用WebScoket则数据产生后就触发推送,不会产生无效请求,总推送数次(只算业务请求,不算建立连接的一次握手)
则可以计算出轮询的请求率为(req/s),WebSocket的请求率为
,WebSocket相对轮询的请求率降幅为
%
我们没有考虑网络传输过程中的时延,只是理论推导,实际结果还需要设计实验。
在股价实时场景中,t= 3s , 延迟L = 2s,interval 通常设置为1s.我们观察T = 60s,统计两种方式的请求数量,以及有效请求和无效请求数量,计算出降幅。
利用上述公式,轮询需要发出60次请求,WebSocket只需要推送20次,降幅为67.7%。
三、WebSocket
1、引入WebSocket依赖
使用WebSocket需要引入依赖
<dependency> <groupId>org.springframework.boot</groupId> <artifactId>spring-boot-starter-websocket</artifactId> </dependency>2、实现WebSocketConfiguration配置类
编写配置类,实现WebSocketConfiguration
@Configuration @EnableWebSocket public class WebSocketConfig implements WebSocketConfigurer { private final DataWebSocketHandler handler; public WebSocketConfig(DataWebSocketHandler handler) { this.handler = handler; } @Override public void registerWebSocketHandlers(WebSocketHandlerRegistry registry) { // 前端连接 ws://localhost:8080/ws registry.addHandler(handler, "/ws"); } }需要重写注册处理器方法,这种方式有点像拦截器.这里的registry注册表需要把URL路径和WebScoket处理器(handler)关联在一起,登记到框架中。
WebScoketHandlerRegistry接口有唯一的方法,addhadler,需要参数handler和url路径。 因此核心推送逻辑在WebScoket处理器中实现,url为webscoket连接提供一个唯一的入口地址,让客户端知道连接到哪里,也让服务端知道该把这条连接交给哪个handler处理。
这里我们使用构造方法注入了DataWebSocketHandler处理器。
3、WebSocket两种消息
WebSocket有两种消息:文本消息(Text)和二进制消息(Binary)
Spring的顶层接口WebSocketHandler中的handler方法里消息类型为泛型,需要手动处理。
void handleMessage(WebSocketSession session, WebSocketMessage<?> message) throws Exception;于是子类AbstractWebSocketHandler做了区分:
public void handleMessage(WebSocketSession session, WebSocketMessage<?> message) throws Exception { if (message instanceof TextMessage textMessage) { this.handleTextMessage(session, textMessage); } else if (message instanceof BinaryMessage binaryMessage) { this.handleBinaryMessage(session, binaryMessage); } else { if (!(message instanceof PongMessage)) { throw new IllegalStateException("Unexpected WebSocket message type: " + message); } PongMessage pongMessage = (PongMessage)message; this.handlePongMessage(session, pongMessage); } }会先判断消息的类型,然后交给不同的handler去处理。而TextWebSocketHandler实现了AbstractWebSocketHandler,因此可以针对文本消息做很好的处理。那我们使用文本消息自然就需要实现TextWebSocketHandler
4、两个关键回调方法
需要实现里面几个关键方法:
我们把用户连接存储起来:concurrentHashMap是线程安全的
/** 保存所有在线连接,数据产生时广播给所有人 */ private final Set<WebSocketSession> sessions = ConcurrentHashMap.newKeySet();@Override public void afterConnectionEstablished(WebSocketSession session) { sessions.add(session); System.out.println("[ws] connected: " + session.getId() + " total=" + sessions.size()); } @Override public void afterConnectionClosed(WebSocketSession session, CloseStatus status) { sessions.remove(session); System.out.println("[ws] disconnected: " + session.getId() + " total=" + sessions.size()); }每一个连接都携带一个Session,id是他们的唯一标识,后面可以遍历session集合来群发消息。这两个方法会在WebSocket连接建立后、WebSocket连接释放后自动调用。作用就是把session及时添加和移除。
完整的Handler如下:
@Component public class DataWebSocketHandler extends TextWebSocketHandler { private final DataStore dataSource; /** 保存所有在线连接,数据产生时广播给所有人 */ private final Set<WebSocketSession> sessions = ConcurrentHashMap.newKeySet(); public DataWebSocketHandler(DataStore dataSource) { this.dataSource = dataSource; } @Override public void afterConnectionEstablished(WebSocketSession session) { sessions.add(session); System.out.println("[ws] connected: " + session.getId() + " total=" + sessions.size()); } @Override public void afterConnectionClosed(WebSocketSession session, CloseStatus status) { sessions.remove(session); System.out.println("[ws] disconnected: " + session.getId() + " total=" + sessions.size()); } /** 由 DataGenerator 在产生数据后调用,主动推给所有客户端 */ public void broadcast(DataStore.DataPoint p) { String payload = String.format( "{\"seq\":%d,\"value\":\"%s\",\"produceTs\":%d,\"pushTs\":%d}", p.seq(), p.value(), p.ts(), System.currentTimeMillis() ); for (WebSocketSession s : sessions) { try { if (s.isOpen()) { s.sendMessage(new TextMessage(payload)); } } catch (Exception e) { e.printStackTrace(); } } System.out.println("[ws] pushed: seq=" + p.seq() + " to " + sessions.size() + " clients"); } }通过WebSocketSession中的isOpen和sendMessage来判断连接是否可用,发送消息。消息内容是自定义的。
5、广播群发
当有数据产生时,我们希望调用broadcat()方法。
@Component public class DataGenerator { private final DataStore dataSource; private final DataWebSocketHandler wsHandler; // 注入 private final ScheduledExecutorService scheduler = Executors.newScheduledThreadPool(1); private static final long T_DATA_MS = 3_000; public DataGenerator(DataStore dataSource, DataWebSocketHandler wsHandler) { this.dataSource = dataSource; this.wsHandler = wsHandler; } @PostConstruct public void start() { scheduler.scheduleAtFixedRate(() -> { DataStore.DataPoint p = dataSource.produce(); wsHandler.broadcast(p); // 产生即推送 }, T_DATA_MS, T_DATA_MS, TimeUnit.MILLISECONDS); } }数据产生的时间是3s,每3s产生一个数据,然后广播推送。数据如何产生(DataStore)我们不必关心,数据如何格式我们也不用关心。
目前我们就把WebSocket配置好了,看看前端怎么做。
6、前端开启WebSocket连接
// ============ WebSocket ============ function connectWs() { ws = new WebSocket('ws://' + location.host + '/ws'); ws.onopen = () => console.log('[ws] open'); ws.onmessage = (ev) => { stats.ws.total++; const data = JSON.parse(ev.data); const now = Date.now(); if (data.seq !== stats.ws.lastSeq) { stats.ws.valid++; stats.ws.lastSeq = data.seq; recordLatency(stats.ws, data.produceTs); const li = document.createElement('li'); li.textContent = `seq=${data.seq} value=${data.value} lat=${now - data.produceTs}ms`; $('wsList').prepend(li); } else { stats.ws.invalid++; // 理论上 WS 不会重复,这里兜底 } }; ws.onclose = () => console.log('[ws] closed'); ws.onerror = (e) => console.error('[ws] error', e); }还记得我们之前说过WebSocket协议是ws或者wss开头的,前端new WebSocket(),配置好连接的地址,当得到服务器允许后协议自动从Http升级到WebSocket.
通过这个新协议我们可以接收和发送消息,这和我们之前学习的HTTP完全不同:HTTP走Controller这个我们非常熟悉,但是WebSocket走Handler,但是后面大家都可以走Service和Mapper.
那么前端发送后端接收参数这块儿写法就不太一样,HTTP借助controller封装,WebSocket就需要借助Interceptor拦截器。不过归根结底,还是Spring框架对HTTP进行了大量封装而已。
后面我们会具体学习使用STOMP模式来简化WebSocket的使用,本节不深入。
四、轮询
轮询的主要工作是前端每隔interval时间就向后端发送请求,负责此功能的代码在前端:
// 一次轮询 async function pollOnce() { total++; const res = await fetch('/poll'); const data = await res.json(); if (data.seq !== lastSeq) { // 有新数据 valid++; lastSeq = data.seq; const lat = Date.now() - data.produceTs; const li = document.createElement('li'); li.textContent = `seq=${data.seq} value=${data.value} lat=${lat}ms`; $('list').prepend(li); } else { // 没新数据 invalid++; } $('total').textContent = total; $('valid').textContent = valid; $('invalid').textContent = invalid; } // 开始轮询 $('startBtn').onclick = () => { timer = setInterval(pollOnce, 1000); // 每 1 秒问一次 $('startBtn').disabled = true; $('stopBtn').disabled = false; }; // 停止 $('stopBtn').onclick = () => { clearInterval(timer); $('startBtn').disabled = false; $('stopBtn').disabled = true; };使用html提供的关键函数timer = setInterval(function,interval)就可以实现轮询调用,使用clearInterval(timer)清除定时器就可以实现结束轮询。
所以,轮询是比较简单的。
五、实验
我们实验参数如下:
在股价实时场景中,t= 3s , 延迟L = 2s,interval 通常设置为1s.我们观察T = 60s,统计两种方式的请求数量,以及有效请求和无效请求数量,计算出降幅。
利用上述公式,轮询需要发出60次请求,WebSocket只需要推送20次,降幅为67.7%。
得到如下结果
可以计算总请求率降幅为67.22%,符合预期。
六、总结
在追求数据高时效性的场景下,建议使用WebSocket而非轮询。使用WebSocket会降低请求率,从而减轻服务器压力。但是WebSocket并不是免费的午餐,是不是需要心跳保活机制判断连接健康状态?最大连接数是不是需要限制?是否占用线程,未来连接过多线程够用吗?如果消息某次推送失败了,是不是需要重发?未来机器横向扩展,session怎么跨机器使用?鉴权到期,安全措施,性能...将会面临很多问题。
针对不同场景,两种方式都有自己的优势,有时候轮询的代价可能远小于WebSocket。