You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

基于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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.08.14 15:20:22