1. 项目概述从零构建一个健壮的WebSocket消息推送服务最近在做一个需要实时消息推送的后台管理系统比如订单状态变更、系统告警、客服聊天这些场景前端得立刻知道。用HTTP轮询太笨重长轮询也麻烦最后选了WebSocket。但真上手才发现光把连接建起来只是第一步离“能用”和“好用”还差得远。比如用户怎么认证连接断了怎么知道怎么给特定的一群人发消息而不是广播这些才是工程实践里的硬骨头。这个项目就是用SpringBoot搭一个完整的WebSocket服务端重点解决四个核心问题连接建立时的身份校验、维持连接可用的心跳机制、服务端主动推送消息以及按业务逻辑对用户进行分组管理。网上很多教程只讲怎么握手成功但生产环境里没心跳的连接说断就断没校验的服务谁都能连分组广播实现不好性能就崩。我会结合我趟过的坑把每个环节的原理、代码和配置细节掰开揉碎了讲目标是让你看完就能搭出一个稳定、安全、可维护的WebSocket推送服务。2. 核心组件选型与环境搭建在SpringBoot里集成WebSocket主流有两种方式一是直接用Spring提供的spring-boot-starter-websocket它底层封装了标准的Java WebSocket APIJSR-356并提供了更Spring风格的编程模型二是通过STOMP子协议它更像一个消息代理适合复杂的消息路由场景。对于我们的需求——点对点、分组推送、心跳保活——直接使用spring-boot-starter-websocket更轻量、更可控也更容易理解底层机制。首先在pom.xml里引入依赖。这里注意我们通常还会引入spring-boot-starter-security来做连接时的安全校验但为了聚焦WebSocket核心我们先手动实现一个简单的token校验。dependency groupIdorg.springframework.boot/groupId artifactIdspring-boot-starter-websocket/artifactId /dependency !-- 用于JSON消息序列化 -- dependency groupIdcom.fasterxml.jackson.core/groupId artifactIdjackson-databind/artifactId /dependency接着需要一个核心配置类来开启WebSocket支持。这里我直接给出一个增强版的配置它做了三件事注册我们的WebSocket处理器、配置允许跨域前端独立部署时必需、以及设置消息缓冲区大小。import org.springframework.context.annotation.Configuration; import org.springframework.web.socket.config.annotation.EnableWebSocket; import org.springframework.web.socket.config.annotation.WebSocketConfigurer; import org.springframework.web.socket.config.annotation.WebSocketHandlerRegistry; Configuration EnableWebSocket public class WebSocketConfig implements WebSocketConfigurer { Override public void registerWebSocketHandlers(WebSocketHandlerRegistry registry) { registry.addHandler(new MyWebSocketHandler(), /ws) .addInterceptors(new AuthHandshakeInterceptor()) // 添加握手拦截器 .setAllowedOrigins(*); // 生产环境应指定具体域名 } }这个配置类里MyWebSocketHandler是我们处理消息的核心类AuthHandshakeInterceptor则负责在握手阶段拦截请求进行身份校验。把校验放在握手阶段比连接建立后再处理要安全得多无效的连接请求在握手时就会被拒绝不会占用服务器资源。注意setAllowedOrigins(*)在开发时图个方便上线前必须改为具体的前端域名如https://yourdomain.com这是防止跨站WebSocket劫持Cross-Site WebSocket Hijacking最基本的一步。从相关热词里看到bp靶场cross-site websocket hijacking指的就是这种攻击攻击者可以在用户浏览器中构造恶意网站利用用户已有的身份比如Cookie向你的WebSocket服务发起连接并窃听或篡改消息。严格限制Origin是首要防线。3. 握手拦截器实现连接身份校验WebSocket协议本身不处理身份认证我们需要在HTTP升级为WebSocket协议的那个握手请求Handshake Request里做文章。Spring WebSocket提供了HandshakeInterceptor接口允许我们在握手前和握手后插入逻辑。我实现了一个简单的基于Token的校验拦截器。思路是前端在建立WebSocket连接时不能像普通HTTP请求那样在Header里带Authorization但可以将Token作为一个查询参数Query Parameter附在连接URL上例如ws://localhost:8080/ws?tokeneyJhbGciOiJ...。我们在拦截器里解析这个token并进行验证。import org.springframework.http.server.ServerHttpRequest; import org.springframework.http.server.ServerHttpResponse; import org.springframework.web.socket.WebSocketHandler; import org.springframework.web.socket.server.HandshakeInterceptor; import java.util.Map; public class AuthHandshakeInterceptor implements HandshakeInterceptor { Override public boolean beforeHandshake(ServerHttpRequest request, ServerHttpResponse response, WebSocketHandler wsHandler, MapString, Object attributes) throws Exception { // 从请求URI中获取token参数 String query request.getURI().getQuery(); if (query null || !query.contains(token)) { // 可以返回false拒绝握手也可以返回401状态码 response.setStatusCode(HttpStatus.UNAUTHORIZED); return false; } // 简单解析token实际项目应使用JWT等规范方案 String token query.substring(query.indexOf(token) 6); // 这里模拟一个简单的校验逻辑 if (!isValidToken(token)) { response.setStatusCode(HttpStatus.FORBIDDEN); return false; } // 校验通过可以将用户信息放入attributes后续在Handler中可取用 String userId extractUserIdFromToken(token); attributes.put(userId, userId); return true; } Override public void afterHandshake(ServerHttpRequest request, ServerHttpResponse response, WebSocketHandler wsHandler, Exception exception) { // 握手成功后调用可用于记录日志等 } private boolean isValidToken(String token) { // 实现你的token验证逻辑例如调用认证服务或校验JWT签名 return token ! null token.startsWith(valid_); } private String extractUserIdFromToken(String token) { // 从token中解析出用户ID这里简单演示 return token.replace(valid_, ); } }这里有个关键点attributes参数。它是一个Map数据在握手阶段被放入后会在WebSocket Session建立后传递给我们自定义的WebSocketHandler。这样我们就成功地将用户身份信息从HTTP握手请求传递到了WebSocket会话中为后续的用户分组和定向推送打下了基础。实操心得把Token放在URL参数里是一种常见做法但要注意其可能被浏览器历史记录、日志服务器记录存在泄漏风险。对于安全性要求极高的场景可以考虑在握手前先通过一个普通的HTTP接口认证服务端返回一个一次性的、短效的connectTicket前端再用这个ticket来建立WebSocket连接。不过大多数内部管理系统使用HTTPS Token的方式已经足够。4. 消息处理器连接、消息与异常的生命周期管理WebSocketHandler是处理所有WebSocket事件的核心。Spring提供了一个方便的适配器类TextWebSocketHandler我们继承它并重写关键方法。这里我设计了一个不仅能处理消息还能管理用户会话和心跳的增强处理器。import org.springframework.web.socket.handler.TextWebSocketHandler; import org.springframework.web.socket.WebSocketSession; import org.springframework.web.socket.TextMessage; import java.util.concurrent.ConcurrentHashMap; public class MyWebSocketHandler extends TextWebSocketHandler { // 存储用户ID与WebSocketSession的映射 private static final ConcurrentHashMapString, WebSocketSession userSessionMap new ConcurrentHashMap(); // 存储Session与最后活跃时间戳用于心跳检测 private static final ConcurrentHashMapString, Long sessionLastActiveTime new ConcurrentHashMap(); Override public void afterConnectionEstablished(WebSocketSession session) throws Exception { // 连接建立成功 String userId (String) session.getAttributes().get(userId); if (userId ! null) { userSessionMap.put(userId, session); sessionLastActiveTime.put(session.getId(), System.currentTimeMillis()); System.out.println(用户 userId 连接成功Session ID: session.getId()); // 可以在这里向该用户发送一条欢迎消息或连接成功通知 session.sendMessage(new TextMessage({\type\:\system\, \msg\:\WebSocket连接已建立\})); } else { // 理论上不会走到这里因为拦截器已校验 session.close(CloseStatus.NOT_ACCEPTABLE.withReason(未识别的用户)); } } Override protected void handleTextMessage(WebSocketSession session, TextMessage message) throws Exception { // 处理客户端发来的文本消息 String payload message.getPayload(); String userId (String) session.getAttributes().get(userId); System.out.println(收到来自用户 userId 的消息: payload); // 更新该会话的最后活跃时间 sessionLastActiveTime.put(session.getId(), System.currentTimeMillis()); // 这里可以解析消息内容根据不同的type执行不同逻辑 // 例如{type: ping} 表示心跳{type: chat, to: user2, content: hello} 表示私聊 // 具体业务逻辑解析... } Override public void handleTransportError(WebSocketSession session, Throwable exception) throws Exception { // 传输过程出错比如解码失败 System.err.println(WebSocket传输错误Session ID: session.getId() , 错误: exception.getMessage()); } Override public void afterConnectionClosed(WebSocketSession session, CloseStatus status) throws Exception { // 连接关闭 String userId (String) session.getAttributes().get(userId); if (userId ! null) { userSessionMap.remove(userId); sessionLastActiveTime.remove(session.getId()); System.out.println(用户 userId 连接关闭状态码: status.getCode() , 原因: status.getReason()); } } }这个处理器框架已经具备了会话管理的基础。userSessionMap让我们能通过用户ID找到对应的WebSocketSession这是实现点对点推送的关键。sessionLastActiveTime则为后续实现心跳超时断开提供了数据支持。5. 心跳机制PING-PONG保活与连接健康度检测WebSocket连接可能因为网络波动、代理超时、客户端崩溃等原因无声无息地断开即“死连接”。服务器和客户端都无法立即感知。心跳机制Heartbeat就是双方定期发送一个小数据包PING/PONG帧来确认对方是否还在线。WebSocket协议层面其实定义了PING和PONG控制帧但Java WebSocket APIJSR-356并没有直接向应用层暴露发送PING帧的接口。通常我们在应用层自己实现一个基于文本或二进制消息的“伪心跳”。我的实现方案是双重的服务端主动探测用一个定时任务定期检查所有会话的最后活跃时间如果超过阈值比如30秒则主动向客户端发送一个PING消息并等待PONG回应。如果一定时间内没收到PONG则认为连接已死主动关闭它。客户端主动上报要求前端每隔一段时间比如25秒主动向服务端发送一个{type:ping}的消息。服务端收到后更新该会话的“最后活跃时间”。这种方式更简单但依赖于客户端的配合。这里重点讲服务端主动探测的实现。我们需要一个Spring的定时任务组件。首先在处理器里增加发送PING和接收PONG的逻辑public class MyWebSocketHandler extends TextWebSocketHandler { // ... 之前的成员变量和方法 ... Override protected void handleTextMessage(WebSocketSession session, TextMessage message) throws Exception { String payload message.getPayload(); // 更新活跃时间 sessionLastActiveTime.put(session.getId(), System.currentTimeMillis()); // 解析消息 ObjectMapper mapper new ObjectMapper(); try { JsonNode node mapper.readTree(payload); String type node.get(type).asText(); if (pong.equals(type)) { // 收到客户端对服务端PING的回应 System.out.println(收到来自会话 session.getId() 的PONG回应); return; // 心跳回应不处理其他业务 } else if (ping.equals(type)) { // 收到客户端主动发来的PING立即回复PONG session.sendMessage(new TextMessage({\type\:\pong\})); return; } // ... 其他业务消息处理 ... } catch (Exception e) { // 消息格式错误处理 } } }然后创建一个定时任务类定期扫描并清理死连接import org.springframework.scheduling.annotation.EnableScheduling; import org.springframework.scheduling.annotation.Scheduled; import org.springframework.stereotype.Component; import org.springframework.web.socket.WebSocketSession; import java.io.IOException; import java.util.Iterator; import java.util.Map; Component EnableScheduling public class WebSocketHeartbeatTask { // 心跳超时时间单位毫秒 private static final long HEARTBEAT_TIMEOUT 35000; // 35秒 // 发送PING的间隔应小于超时时间 private static final long PING_INTERVAL 30000; // 30秒 Scheduled(fixedRate 15000) // 每15秒执行一次检查 public void checkAlive() { long currentTime System.currentTimeMillis(); IteratorMap.EntryString, Long iterator MyWebSocketHandler.sessionLastActiveTime.entrySet().iterator(); while (iterator.hasNext()) { Map.EntryString, Long entry iterator.next(); String sessionId entry.getKey(); Long lastActiveTime entry.getValue(); WebSocketSession session findSessionById(sessionId); // 需要根据sessionId找到session对象 if (session null || !session.isOpen()) { // 会话已不存在或已关闭清理记录 iterator.remove(); MyWebSocketHandler.userSessionMap.values().removeIf(s - s.getId().equals(sessionId)); continue; } long inactiveDuration currentTime - lastActiveTime; if (inactiveDuration HEARTBEAT_TIMEOUT) { // 超过超时时间直接关闭连接 try { session.close(CloseStatus.SESSION_NOT_RELIABLE); System.out.println(因心跳超时关闭会话: sessionId); } catch (IOException e) { e.printStackTrace(); } iterator.remove(); MyWebSocketHandler.userSessionMap.values().removeIf(s - s.getId().equals(sessionId)); } else if (inactiveDuration PING_INTERVAL) { // 超过PING间隔但未超时发送一个PING探测 try { session.sendMessage(new TextMessage({\type\:\ping\})); System.out.println(向会话 sessionId 发送PING探测); } catch (IOException e) { // 发送失败可能连接已失效下次检查会处理 e.printStackTrace(); } } } } // 一个辅助方法根据sessionId从userSessionMap中查找session // 注意userSessionMap存储的是userId-session这里需要遍历查找实际可优化数据结构 private WebSocketSession findSessionById(String targetSessionId) { for (WebSocketSession session : MyWebSocketHandler.userSessionMap.values()) { if (session.getId().equals(targetSessionId)) { return session; } } return null; } }踩坑实录心跳超时时间HEARTBEAT_TIMEOUT和PING间隔PING_INTERVAL的设定需要谨慎。间隔太短会产生大量无用心跳包增加服务器和网络负担间隔太长则无法及时发现死连接。一般建议PING间隔为25-30秒超时时间比间隔多5-10秒给网络延迟和客户端处理留出余量。另外像Nginx这样的反向代理默认会对WebSocket连接有一个60秒的超时proxy_read_timeout你的服务端心跳超时必须小于这个值否则连接会被代理服务器先掐断。6. 用户分组与定向消息推送广播消息很简单遍历userSessionMap的所有session发送即可。但实际业务中更多是需要按条件推送比如推送给某个在线的特定用户点对点、推送给某个部门的所有人、推送给具有某个角色标签的所有用户。这就是用户分组。我的实现思路是在用户连接建立时不仅记录userId-session的映射还根据业务规则将用户加入到不同的“逻辑组”中。这个组可以用一个ConcurrentHashMapString, SetString来维护key是组名如dept:finance,role:adminvalue是该组内所有用户的ID集合。首先在Handler中增加分组管理容器public class MyWebSocketHandler extends TextWebSocketHandler { // ... 已有成员变量 ... // 分组存储组名 - 用户ID集合 private static final ConcurrentHashMapString, SetString userGroupMap new ConcurrentHashMap(); Override public void afterConnectionEstablished(WebSocketSession session) throws Exception { String userId (String) session.getAttributes().get(userId); if (userId ! null) { userSessionMap.put(userId, session); sessionLastActiveTime.put(session.getId(), System.currentTimeMillis()); // --- 关键根据业务逻辑将用户加入分组 --- // 假设我们从数据库或上下文中获取用户的部门、角色等信息 ListString userGroups getUserGroupsFromDatabase(userId); // 伪方法 for (String group : userGroups) { userGroupMap.computeIfAbsent(group, k - ConcurrentHashMap.newKeySet()).add(userId); } System.out.println(用户 userId 加入分组: userGroups); // --- 分组结束 --- } } Override public void afterConnectionClosed(WebSocketSession session, CloseStatus status) throws Exception { String userId (String) session.getAttributes().get(userId); if (userId ! null) { userSessionMap.remove(userId); sessionLastActiveTime.remove(session.getId()); // --- 关键用户断开时从所有分组中移除 --- for (SetString groupUsers : userGroupMap.values()) { groupUsers.remove(userId); } // 注意这里不移除空的组如果需要可以定期清理 // --- 分组结束 --- } } // ... 其他方法 ... }然后提供几个核心的推送方法public class MyWebSocketHandler extends TextWebSocketHandler { // ... 其他代码 ... /** * 向单个用户发送消息 */ public static void sendMessageToUser(String userId, String message) throws IOException { WebSocketSession session userSessionMap.get(userId); if (session ! null session.isOpen()) { synchronized (session) { // 发送消息需要同步防止多线程并发写冲突 session.sendMessage(new TextMessage(message)); } } else { // 用户不在线可以存入消息队列或数据库待其上线后推送 System.out.println(用户 userId 不在线消息无法实时送达); } } /** * 向特定分组的所有在线用户广播消息 */ public static void sendMessageToGroup(String groupName, String message) { SetString userIds userGroupMap.get(groupName); if (userIds ! null !userIds.isEmpty()) { for (String userId : userIds) { try { sendMessageToUser(userId, message); } catch (IOException e) { System.err.println(向分组用户 userId 发送消息失败: e.getMessage()); } } } } /** * 全局广播给所有在线用户 */ public static void broadcastMessage(String message) { for (WebSocketSession session : userSessionMap.values()) { if (session.isOpen()) { try { synchronized (session) { session.sendMessage(new TextMessage(message)); } } catch (IOException e) { System.err.println(广播消息失败Session ID: session.getId()); } } } } }这样在任何一个Spring管理的Bean中比如Service层你都可以通过MyWebSocketHandler.sendMessageToUser(userId, msg)或MyWebSocketHandler.sendMessageToGroup(dept:sales, msg)来触发消息推送了。性能优化提示当分组内用户数量巨大比如上万人时遍历发送可能会阻塞业务线程。可以考虑将消息推送任务提交给一个专门的线程池来异步执行。另外userGroupMap的数据结构可以进一步优化例如使用Guava的Multimap或者维护一个userId-groups的反向索引方便用户下线时快速从所有组中移除。7. 消息格式设计与业务处理前面我们一直用简单的JSON字符串{type:ping}作为消息。在实际项目中需要设计一个统一的消息格式方便前端和后端解析。一个常见的格式如下{ type: chat/notification/system/..., sender: user123, receiver: user456 / group:dept1 / all, timestamp: 1640995200000, payload: { // 实际的消息内容结构根据type不同而变化 title: 新订单, content: 您有一笔新的订单待处理订单号202412310001, url: /order/detail/202412310001 } }在Handler的handleTextMessage方法中我们需要根据type字段来路由到不同的业务处理器Override protected void handleTextMessage(WebSocketSession session, TextMessage message) throws Exception { String payload message.getPayload(); sessionLastActiveTime.put(session.getId(), System.currentTimeMillis()); // 更新心跳时间 ObjectMapper mapper new ObjectMapper(); try { JsonNode rootNode mapper.readTree(payload); String type rootNode.path(type).asText(null); // 安全获取避免NPE String sender (String) session.getAttributes().get(userId); if (ping.equals(type)) { // 心跳 session.sendMessage(new TextMessage({\type\:\pong\})); return; } else if (chat.equals(type)) { // 私聊 String receiver rootNode.path(receiver).asText(); JsonNode msgPayload rootNode.path(payload); handlePrivateChat(sender, receiver, msgPayload); } else if (join_group.equals(type)) { // 动态加入分组 String groupName rootNode.path(group).asText(); joinGroup(sender, groupName); } else if (leave_group.equals(type)) { // 动态离开分组 String groupName rootNode.path(group).asText(); leaveGroup(sender, groupName); } // ... 其他业务类型 ... } catch (JsonProcessingException e) { session.sendMessage(new TextMessage({\type\:\error\, \msg\:\消息格式错误\})); } catch (Exception e) { session.sendMessage(new TextMessage({\type\:\error\, \msg\:\服务器处理异常\})); } } private void handlePrivateChat(String sender, String receiver, JsonNode payload) throws IOException { String content payload.path(content).asText(); String formattedMsg String.format({\type\:\chat\, \from\:\%s\, \content\:\%s\}, sender, content); sendMessageToUser(receiver, formattedMsg); // 可选也发一份给发送者自己作为发送成功回执 sendMessageToUser(sender, {\type\:\chat_status\, \status\:\sent\, \to\:\ receiver \}); } private void joinGroup(String userId, String groupName) { userGroupMap.computeIfAbsent(groupName, k - ConcurrentHashMap.newKeySet()).add(userId); // 通知该用户已加入组 try { sendMessageToUser(userId, {\type\:\system\, \msg\:\你已成功加入分组: groupName \}); } catch (IOException e) { e.printStackTrace(); } }这种设计使得消息处理逻辑清晰易于扩展新的消息类型。8. 生产环境部署与进阶考量把代码跑起来只是开始要上线还得过好几关。连接数限制与资源管理每个WebSocket连接都会占用一个线程取决于容器实现如Tomcat的NIO模式和内存。默认的Tomcat配置可能只能处理几千个并发连接。你需要调整application.properties# 增大最大连接数 server.tomcat.max-connections10000 # 增大工作线程数 server.tomcat.threads.max200 # 调整WebSocket相关的缓冲区大小 server.tomcat.max-swallow-size2MB同时必须在代码中做好连接管理及时清理无效会话我们的心跳机制就在做这件事防止内存泄漏。集群部署与会话共享上面的代码把所有会话信息userSessionMap,userGroupMap都存在单个应用实例的内存里。一旦部署多台实例用户可能连到A实例但推送消息的请求发到了B实例B实例上根本没有这个用户的session导致推送失败。解决方案是引入外部存储来共享会话信息例如Redis。连接建立时将userId、instanceId实例标识如IP:PORT和group信息存入Redis并设置过期时间略大于心跳超时时间。发送消息时先根据userId或group从Redis查出所有在线的用户及其所在的instanceId。如果目标用户就在当前实例直接发送如果在其他实例则需要通过消息队列如RabbitMQ、Kafka或者HTTP调用需要实例间有通信能力将消息转发到对应实例去发送。这是一个架构上的重大变化通常会引入Spring Cloud、WebSocket Stomp Broker Relay或者自研一个轻量的消息转发层。前端连接与重连策略前端不能只连一次就完事。需要监听WebSocket的onclose和onerror事件实现自动重连并采用指数退避策略比如第一次断线等1秒重连第二次等2秒第三次等4秒...。重连时同样要带上认证Token。监控与日志务必记录连接建立、关闭、消息收发、异常断开等关键事件并监控活跃连接数、消息吞吐量等指标这对排查线上问题至关重要。我自己在项目上线后就遇到过因为Nginx的proxy_read_timeout设置比服务端心跳间隔短导致连接频繁被代理断开的问题。后来统一了超时配置并在服务端增加了更细致的连接状态日志才稳定下来。WebSocket看似简单真想在生产环境扛住流量每一个细节都得抠。