ARTICLE DETAIL

资讯详情

深耕网站建设与运营推广的一线实战洞察。

Java WebSocket 实战:服务端客户端、心跳机制与断线重连

Java WebSocket 实战:服务端客户端、心跳机制与断线重连 做后端这几年WebSocket 是我反复接触、也反复踩坑的一个点。最近整理项目时发现网上关于 Java WebSocket 的使用示例大多只给一段片段要么服务端注解贴一段要么前端 onMessage 贴一段真正把服务端、客户端、心跳机制、断线重连和问题排查串成一条线的东西很少。这篇文章就把我实战中验证过的 Java WebSocket 使用示例汇总一下覆盖标准 API、Spring Boot 集成、JS 和 Java 客户端写法以及生产环境里必须处理的心跳与重连问题。适合正在做实时推送、站内信、客服系统、数据大屏的后端同学也适合面试前想系统梳理一遍的 Java 开发。1. 为什么是 WebSocket它解决了什么以及什么时候别用它1.1 HTTP 轮询的两大痛点在 WebSocket 出现之前要实现服务端主动推送最常见的方案是轮询。短轮询就是浏览器每隔几秒发一次 HTTP 请求服务端不管有没有新数据都响应这个请求本身的开销非常大长轮询稍微聪明一点服务端在拿到请求后不立即返回而是把连接挂着有数据或者超时才响应但每个连接依然是一次完整的 HTTP 生命周期Header、Cookie、鉴权每次都要重复走一遍。我做过一个在线用户状态模块最开始用短轮询每 3 秒拉一次单机 2000 人同时在线时网关层每秒 QPS 直接飙到 700 左右其中 90% 的请求都是不带新增数据的空响应。这不是带宽问题而是无意义的请求放大了连接开销、CPU 开销和日志量。轮询还有一个体验问题延迟不可控。3 秒轮询意味着消息最快也可能要等 3 秒才到达做聊天和实时大屏根本不可接受。1.2 WebSocket 的一次握手一次升级WebSocket 是一个在单个 TCP 连接上进行全双工通信的协议。它复用了 HTTP 的握手机制客户端发起一个带有Upgrade: websocket头的 HTTP 请求服务端如果支持就返回101 Switching Protocols这样双方的通信就升级成了 WebSocket 连接。后续的数据传递通过帧frame完成不再需要 HTTP 头开销大幅下降。帧类型主要有文本帧、二进制帧、Ping/Pong 控制帧和 Close 帧。文本帧用于传输字符串消息二进制帧可以传文件流或序列化对象Ping/Pong 专用于连接保活这也对应到后面要讲的心跳机制。理解握手和帧类型后许多问题就好办了。比如你看到连接频繁掉线第一反应不该是调大代理超时而是先确认代理层有没有正确转发 Upgrade 头。另外要泼一盆冷水WebSocket 不是万能的。它适合服务端主动推送、双向交互、长连接场景简单的一次性查询、低频的 API 调用用 HTTP 反而更合适。长连接本身会占用文件描述符和内存如果只是做个定时刷新接口强行上 WebSocket 只会给自己添麻烦。2. 服务端实现注解式与 Spring Boot 两种写法2.1 标准 API 注解式开发Java 标准 WebSocket API 的基本用法是ServerEndpoint加回调方法。现在 Spring Boot 3 用的是jakarta.websocket包Spring Boot 2 用javax.websocket写法基本一致。import jakarta.websocket.*; import jakarta.websocket.server.ServerEndpoint; import java.io.IOException; import java.util.concurrent.CopyOnWriteArraySet; ServerEndpoint(/ws/chat) public class ChatWebSocket { private static final CopyOnWriteArraySetSession SESSIONS new CopyOnWriteArraySet(); private Session session; private String userId; OnOpen public void onOpen(Session session, PathParam(userId) String userId) { this.session session; this.userId userId; SESSIONS.add(session); System.out.println(连接建立: userId , 当前连接数: SESSIONS.size()); } OnMessage public void onMessage(String message, Session session) { // 解析消息、做业务处理、转发给所有在线用户 for (Session s : SESSIONS) { if (s.isOpen()) { s.getBasicRemote().sendText(message); } } } OnClose public void onClose(Session session, CloseReason reason) { SESSIONS.remove(session); System.out.println(连接关闭: reason.getReasonPhrase()); } OnError public void onError(Session session, Throwable error) { error.printStackTrace(); } }这里有几个容易忽略的细节。SESSIONS 集合我用了CopyOnWriteArraySet因为 WebSocket 消息回调可能来自多个线程普通 HashSet 在并发遍历时很容易报ConcurrentModificationException。每个客户端的 Session 是否 open发送前必须判断不然你以为对方还在实际上已经掉线一sendText就抛 IOException。在 Spring Boot 项目里如果直接用ServerEndpoint还需要额外注册一个ServerEndpointExporterBean否则 Endpoint 不会被扫描进容器。Configuration public class WebSocketConfig { Bean public ServerEndpointExporter serverEndpointExporter() { return new ServerEndpointExporter(); } }2.2 为什么我更推荐 Spring 的 WebSocketHandler如果你希望把 WebSocket 跟 Spring 容器真正打通我更推荐实现WebSocketHandler而不是用ServerEndpoint注解。原因很实在注解类默认不是 Spring 管理的 Bean它的生命周期由容器管理在里面直接Autowired注入 Service 往往是 null。虽然可以用静态工具类曲线救国但代码维护起来很别扭。Configuration EnableWebSocket public class WebSocketConfig implements WebSocketConfigurer { Override public void registerWebSocketHandlers(WebSocketHandlerRegistry registry) { registry.addHandler(chatHandler(), /ws/chat) .addInterceptors(new AuthHandshakeInterceptor()) .setAllowedOrigins(*); } Bean public WebSocketHandler chatHandler() { return new ChatWebSocketHandler(); } }Handler 类的写法如下。注意这里用的是org.springframework.web.socket.WebSocketSession它跟jakarta.websocket.Session不是同一个类千万别混用。Component public class ChatWebSocketHandler extends TextWebSocketHandler { private static final MapString, WebSocketSession SESSIONS new ConcurrentHashMap(); Override public void afterConnectionEstablished(WebSocketSession session) { String userId (String) session.getAttributes().get(userId); SESSIONS.put(userId, session); // 这里可以注入 Service 发送上线通知 } Override protected void handleTextMessage(WebSocketSession session, TextMessage message) throws Exception { String payload message.getPayload(); // 解析 payload按 type 分发 session.sendMessage(new TextMessage({\code\:0})); } Override public void afterConnectionClosed(WebSocketSession session, CloseStatus status) { String userId (String) session.getAttributes().get(userId); SESSIONS.remove(userId); } }处理消息时还是要注意handleTextMessage会被多个线程并发调用如果要在里面更新共享数据要么用ConcurrentHashMap要么加锁别图省事直接用一个普通 HashMap。2.3 鉴权与跨域处理WebSocket 握手本质是 HTTP 请求所以鉴权可以放在HandshakeInterceptor里。在beforeHandshake方法中拿到请求参数或 Header校验 token不合法直接返回 false 拒绝握手。合法就把用户信息放进 attributes后续afterConnectionEstablished、handleTextMessage都能从session.getAttributes()取到。Component public class AuthHandshakeInterceptor implements HandshakeInterceptor { Override public boolean beforeHandshake(ServerHttpRequest request, ServerHttpResponse response, WebSocketHandler wsHandler, MapString, Object attributes) { String token request.getHeaders().getFirst(Authorization); if (token null || !token.equals(expected-token)) { response.setStatusCode(HttpStatus.UNAUTHORIZED); return false; } attributes.put(userId, user-1024); return true; } }跨域方面setAllowedOrigins(*)是最省事的写法但生产环境建议换成具体域名列表否则任何页面都能连你的 WebSocket这不只是安全问题还会增加无效连接压力。3. 客户端用法浏览器 JS 端与 Java 客户端3.1 前端 WebSocket 客户端标准写法浏览器端的 API 很简洁核心就是四个事件回调加一个send方法。我见过不少前端同事把readyState判断漏了导致连接断开后 send 静默失败所以下面示例会带上状态判断。const ws new WebSocket(ws://localhost:8080/ws/chat?userId1001); ws.onopen function () { console.log(连接建立); ws.send(JSON.stringify({ type: login, data: { token: xxx } })); }; ws.onmessage function (event) { const msg JSON.parse(event.data); // 根据 msg.type 分发处理 if (msg.type chat) { renderMessage(msg.data); } }; ws.onclose function (event) { console.warn(连接关闭, event.code, event.reason); }; ws.onerror function (error) { console.error(连接异常, error); }; function safeSend(data) { if (ws.readyState WebSocket.OPEN) { ws.send(JSON.stringify(data)); } }readyState一共有四种状态CONNECTING、OPEN、CLOSING、CLOSED。很多人只判断OPEN这没错但要注意在CONNECTING状态下调用 send 也会报错。所以最稳的方式是等 onopen 之后设置一个标志位再发送业务消息。3.2 Java 客户端JDK 自带 API 与 Java-WebSocket 库服务端写完调试时总得有客户端工具。除了在线调试工具Java 项目里自己写客户端测试也很方便。JDK 11 开始提供了java.net.http.WebSocket不需要额外依赖import java.net.URI; import java.net.http.HttpClient; import java.net.http.WebSocket; import java.util.concurrent.CompletionStage; public class JdkClientExample { public static void main(String[] args) throws Exception { HttpClient client HttpClient.newHttpClient(); WebSocket ws client.newWebSocketBuilder() .buildAsync(URI.create(ws://localhost:8080/ws/chat), new WebSocket.Listener() { Override public void onOpen(WebSocket webSocket) { webSocket.sendText({\type\:\login\}, true); } Override public CompletionStage? onText(WebSocket webSocket, CharSequence data, boolean last) { System.out.println(收到: data); return null; } Override public CompletionStage? onClose(WebSocket webSocket, int statusCode, String reason) { System.out.println(关闭: reason); return null; } }).join(); Thread.sleep(5000); } }如果想用轻量级第三方库Java-WebSocket 也比较稳定Maven 坐标是org.java-websocket:Java-WebSocket。它的回调风格更像浏览器端写起来直观import org.java_websocket.client.WebSocketClient; import org.java_websocket.handshake.ServerHandshake; import java.net.URI; public class ThirdPartyClient extends WebSocketClient { public ThirdPartyClient(URI uri) { super(uri); } Override public void onOpen(ServerHandshake handshake) { System.out.println(连接建立); } Override public void onMessage(String message) { System.out.println(收到: message); } public static void main(String[] args) throws Exception { ThirdPartyClient client new ThirdPartyClient(new URI(ws://localhost:8080/ws/chat)); client.connect(); client.send({\type\:\ping\}); Thread.sleep(5000); } }Java 客户端通常用于两种场景一是服务端之间的消息互通比如一个后端需要连接另一个系统的 WebSocket 接收推送二是自动化测试真实模拟客户端行为。写测试时我建议把消息收发的关键日志都打出来调试长连接问题时效太高了。3.3 消息格式前后端必须约定清楚长连接项目里消息格式混乱是最痛的问题之一。我建议从一开始就统一成 JSON并且带一个type字段做路由data字段装业务数据必要的话再加一个traceId用于排查链路。{ type: chat, traceId: abc-123, data: { from: user-1001, to: user-1002, content: 你好 } }服务端收到字符串后第一步是解析 JSON 获取 type然后按 type 分发到不同的业务方法。不要图省事把所有逻辑都写在 onMessage 里否则后面加一个消息类型就得动核心类维护成本直线上升。4. 心跳机制与断线重连生产环境必踩的坑4.1 WebSocket 为什么会静默断开服务端和客户端之间如果一段时间没有数据帧传输中间的网络设备比如 Nginx、防火墙、负载均衡器就可能认为连接空闲强制回收。实际上TCP 层面的有效数据包越少越容易被判定为僵尸连接。一旦被回收两端都不会立刻感知之后任何一方再发数据才会触发错误但那一刻业务已经被打断了。心跳机制就是解决这个问题的标准手段。简单说就是定期发送一个小数据包证明我还活着。WebSocket 协议层面有 Ping/Pong 帧服务端可以发 Ping客户端回 Pong浏览器和多数框架会自动回复应用层也可以自己约定比如发{type:ping}收到{type:pong}就算存活。4.2 双重心跳实现方案我推荐在实际项目里做双重心跳服务端定时扫描发心跳客户端再做一个应用级 ping 并统计未响应次数。这样即使某一层被代理挡了另一层也能兜底。服务端定时扫描思路用Scheduled每隔 30 秒遍历所有 Session通过session.isOpen()判断存活再发一个应用层 ping 消息。注意遍历时不能直接修改集合结构否则可能抛异常所以 Session 存储得用支持并发遍历的CopyOnWriteArraySet或ConcurrentHashMap。Scheduled(fixedDelay 30000) public void heartBeat() { for (WebSocketSession session : SESSIONS.values()) { if (session.isOpen()) { try { session.sendMessage(new TextMessage({\type\:\ping\})); } catch (IOException e) { // 发送失败关闭连接移出集合 try { session.close(CloseStatus.SESSION_NOT_RELIABLE); } catch (IOException ex) { // ignore } SESSIONS.values().remove(session); } } } }客户端收到 ping 后回一个 pongws.onmessage function (event) { const msg JSON.parse(event.data); if (msg.type ping) { ws.send(JSON.stringify({ type: pong })); return; } // 其他业务消息处理 };心跳间隔的取值有讲究。30 秒是一个双保险的折中选择——太频繁会增加无谓流量太疏则起不到保活效果。如果代理层的proxy_read_timeout默认是 60 秒心跳至少要比 60 秒短否则连接会在两次心跳之间被掐断。生产环境我会先确认链路里所有网络设备的超时时间再定心跳间隔。4.3 断线重连与指数退避心跳能尽量保活但网络抖动、服务重启都会造成连接断开。断线后必须重连但重连不能太激进否则服务端会被打爆。指数退避是常见策略第一次 1 秒第二次 2 秒第三次 4 秒最大不超过 30 秒并且加一点随机抖动防止大量客户端同时重连形成雪崩。let retryCount 0; const maxRetry 8; function connect() { const ws new WebSocket(ws://localhost:8080/ws/chat); ws.onclose function () { if (retryCount maxRetry) { const delay Math.min(1000 * Math.pow(2, retryCount), 30000) Math.random() * 1000; retryCount; setTimeout(connect, delay); } }; ws.onopen function () { retryCount 0; }; } connect();重连成功后还有一个容易忽略的问题消息补偿。断线期间用户可能漏掉一批消息正确做法是重连后让服务端把增量消息补回来。比如客户端连接时带上最后一条消息的 id服务端从该 id 之后开始补发。这个机制是否做取决于业务容忍度但对于支付通知、订单状态这类关键消息我强烈建议做。5. 高频问题与排查技巧实录5.1 Nginx 代理导致连接频繁断开用 Nginx 做反向代理时若没有正确配置WebSocket 连接通常撑不过 60 秒。这是因为 Nginx 默认的proxy_read_timeout是 60 秒而且默认 HTTP 版本是 1.0不支持 Upgrade。关键是 location 里必须显式声明 Upgrade 头并保持连接location /ws/ { proxy_pass http://backend; proxy_http_version 1.1; proxy_set_header Upgrade $http_upgrade; proxy_set_header Connection upgrade; proxy_read_timeout 3600s; proxy_send_timeout 3600s; }这里我把超时调成了 1 小时但注意这只是配合心跳的一种边界条件。如果应用层有 30 秒心跳这个超时时间其实不需要那么极端60 秒到 5 分钟足够。核心逻辑是心跳间隔必须小于所有代理层的超时时间否则再大的超时也会被网络设备掐断。5.2 同一 Session 并发写导致异常我在项目里遇到过Multiple messages sent to this WebSocketSession的报错。原因是一个 Session 的sendMessage方法不是线程安全的多个线程同时往同一个连接写数据就会抛出 IllegalStateException。解决办法很简单对每个 Session 的发送操作加锁或者用一个专用线程池统一发消息。public synchronized void sendToUser(String userId, String message) { WebSocketSession session SESSIONS.get(userId); if (session ! null session.isOpen()) { try { session.sendMessage(new TextMessage(message)); } catch (IOException e) { SESSIONS.remove(userId); } } }5.3 Session 管理不当导致内存泄漏长连接场景里Session 集合如果没有及时清理掉线用户的 Session 会一直留在内存里最终拖垮服务。我见过比较常见的两个坑一是afterConnectionClosed里忘了SESSIONS.remove二是后端主动关闭连接时没有从集合中移除。建议统一封装一个closeSession方法无论异常还是正常关闭都走同一条清理逻辑并且定期打印连接数做观察。5.4 集群环境下消息不到达单机 WebSocket 不存在这个问题但一旦部署了多个实例用户 A 连在实例 1用户 B 连在实例 2A 给 B 发消息时实例 1 无法直接找到 B 的 Session。这不是 WebSocket 本身的问题而是分布式状态下 Session 不共享。常规方案是用 Redis 发布订阅或消息中间件做广播实例 1 把消息发到 Redis channel实例 2 订阅到后找到本机上的 B 的 Session 再推送。关键点是建立一个用户标识 - 节点标识的全局路由表不然你不知道该发给谁。5.5 高频问题速查表现象可能原因处理方式连接 60 秒左右被断开代理超时或未配置心跳配置 Upgrade 头缩短心跳间隔调高代理超时发送消息抛 IllegalStateException同一 Session 并发写入对发送方法加 synchronized 或用发送线程池连接正常但收不到推送集群模式下目标 Session 不在当前节点引入 Redis 发布订阅或 MQ 广播内存持续增长Session 未在关闭时清理封装统一的连接清理逻辑并定期观察连接数ServerEndpoint中注入的 Service 为 null注解类不受 Spring 容器管理换用 WebSocketHandler 方式或通过静态 ApplicationContext 获取浏览器收到的消息乱码编码不一致统一使用 UTF-8服务端设置 charset我在实际排查中还有一个体会遇到连接异常不要急着改代码先用日志把服务端和客户端的收发时间线拉出来对一遍很多时候问题就出在时机上。比如客户端以为连接还在服务端其实已经关闭了或者服务端发了 ping客户端收到了但回包被代理丢弃。把日志按 sessionId 过滤后这些一眼就能看出来。心跳和重连这两个机制往往是长连接系统稳定性的分水岭。功能写完只是开始真正上线后能否扛住断线风暴、网络抖动靠的就是这两块细节。我把它们放在文章最后写也是因为它们最容易被新手忽略却最值得花时间打磨。
返回列表