基于Spring Messaging WebSocket,如何通过STOMP检测客户端网络故障
检测WebSocket客户端网络故障的方案(Spring Messaging 5.3.21)
1. 基于STOMP心跳机制的自动检测
Spring WebSocket原生支持STOMP心跳协议,通过配置心跳参数,服务端会自动检测客户端心跳超时(即网络故障导致的连接中断),并触发断开事件。
配置心跳参数
在WebSocket配置类中配置心跳间隔与超时:
@Configuration @EnableWebSocketMessageBroker public class WebSocketConfig implements WebSocketMessageBrokerConfigurer { @Override public void configureMessageBroker(MessageBrokerRegistry config) { // 客户端每10秒发送一次心跳,服务端每30秒发送一次,超时阈值设为15秒 config.setHeartbeatValue(new long[]{10000, 30000}) .setTaskScheduler(taskScheduler()); } @Bean public TaskScheduler taskScheduler() { ThreadPoolTaskScheduler scheduler = new ThreadPoolTaskScheduler(); scheduler.setPoolSize(1); scheduler.setThreadNamePrefix("websocket-heartbeat-"); scheduler.initialize(); return scheduler; } }
监听断开事件
实现ApplicationListener监听SessionDisconnectEvent,当客户端心跳超时或网络故障断开时,会触发该事件:
@Component public class WebSocketDisconnectListener implements ApplicationListener<SessionDisconnectEvent> { @Override public void onApplicationEvent(SessionDisconnectEvent event) { StompHeaderAccessor headerAccessor = StompHeaderAccessor.wrap(event.getMessage()); String sessionId = headerAccessor.getSessionId(); // 处理客户端断开逻辑:日志记录、资源清理等 System.out.println("客户端因网络故障断开连接,Session ID: " + sessionId); // 可通过CloseStatus获取断开原因 CloseStatus closeStatus = event.getCloseStatus(); if (closeStatus != null) { System.out.println("断开原因: " + closeStatus.getReason()); } } }
2. 自定义拦截器跟踪连接状态
在你已有的ChannelInterceptor基础上,维护客户端心跳时间戳,定时检查超时情况:
@Component public class WebSocketAuthInterceptor implements ChannelInterceptor { private final Map<String, Long> sessionHeartbeatMap = new ConcurrentHashMap<>(); private final ScheduledExecutorService scheduler = Executors.newSingleThreadScheduledExecutor(); @PostConstruct public void init() { // 每15秒检查一次心跳超时 scheduler.scheduleAtFixedRate(this::checkHeartbeatTimeouts, 15, 15, TimeUnit.SECONDS); } @Override public Message<?> preSend(Message<?> message, MessageChannel channel) { try { StompHeaderAccessor headerAccessor = StompHeaderAccessor.wrap(message); String sessionId = headerAccessor.getSessionId(); if (StompCommand.CONNECT.equals(headerAccessor.getCommand())) { // 原有验证逻辑 // 记录新连接的初始心跳时间 sessionHeartbeatMap.put(sessionId, System.currentTimeMillis()); } else if (StompCommand.HEARTBEAT.equals(headerAccessor.getCommand())) { // 更新心跳时间戳 sessionHeartbeatMap.put(sessionId, System.currentTimeMillis()); } else if (StompCommand.DISCONNECT.equals(headerAccessor.getCommand())) { // 客户端主动断开,移除Session记录 sessionHeartbeatMap.remove(sessionId); } } catch (Exception e) { // 异常处理逻辑 } return message; } private void checkHeartbeatTimeouts() { long currentTime = System.currentTimeMillis(); // 设置超时时间为20秒(略大于客户端心跳间隔) long timeout = 20000; sessionHeartbeatMap.entrySet().removeIf(entry -> { if (currentTime - entry.getValue() > timeout) { String sessionId = entry.getKey(); // 处理网络故障导致的超时断开 System.out.println("客户端网络故障,Session ID: " + sessionId + " 心跳超时"); return true; } return false; }); } @PreDestroy public void destroy() { scheduler.shutdown(); } }
3. WebSocketSession生命周期回调
通过HandshakeInterceptor获取WebSocketSession,添加关闭监听:
@Component public class WebSocketHandshakeInterceptor implements HandshakeInterceptor { @Override public boolean beforeHandshake(ServerHttpRequest request, ServerHttpResponse response, WebSocketHandler wsHandler, Map<String, Object> attributes) throws Exception { // 握手前验证逻辑 return true; } @Override public void afterHandshake(ServerHttpRequest request, ServerHttpResponse response, WebSocketHandler wsHandler, Exception exception) { if (request instanceof ServletServerHttpRequest) { ServletServerHttpRequest servletRequest = (ServletServerHttpRequest) request; javax.websocket.Session session = (javax.websocket.Session) servletRequest.getServletRequest().getAttribute("javax.websocket.Session"); if (session != null) { session.addCloseListener(closeReason -> { // 连接关闭时触发,包含网络故障场景 System.out.println("客户端断开连接,原因: " + closeReason.getReasonPhrase()); }); } } } }
在WebSocket配置类中注册该拦截器:
@Configuration @EnableWebSocketMessageBroker public class WebSocketConfig implements WebSocketMessageBrokerConfigurer { @Autowired private WebSocketHandshakeInterceptor handshakeInterceptor; @Override public void registerStompEndpoints(StompEndpointRegistry registry) { registry.addEndpoint("/ws") .addInterceptors(handshakeInterceptor) .withSockJS(); } // 其他配置... }
内容的提问来源于stack exchange,提问作者Ammar Zafar
相关产品推荐
相关产品推荐

