WebSocket实时通信实战:连接管理与心跳保活机制设计

WebSocket协议与实时通信场景适配

WebSocket是HTML5标准中定义的全双工通信协议,通过单条TCP连接实现服务器与客户端的双向实时通信。与传统HTTP轮询和长轮询相比,WebSocket的握手开销仅发生在连接建立时,后续数据帧直接通过TCP通道传输,延迟低至网络RTT级别。后端开发中WebSocket广泛应用于即时通讯、实时协作、股票行情推送和在线游戏等场景。Spring Boot、Go和Node.js均提供成熟的WebSocket支持,高并发设计的关键在于连接管理、心跳保活和消息广播效率。

Spring Boot WebSocket服务端实现

Spring Boot通过spring-boot-starter-websocket提供WebSocket支持,结合STOMP协议可实现基于主题/队列的消息路由。对于简单的点对点通信场景,原生WebSocket Endpoint即可满足需求。

// WebSocket配置类
@Configuration
@EnableWebSocket
public class WebSocketConfig implements WebSocketConfigurer {
    
    @Override
    public void registerWebSocketHandlers(WebSocketHandlerRegistry registry) {
        registry.addHandler(new RealTimeHandler(), "/ws/realtime")
                .setAllowedOrigins("https://app.example.com")
                .addInterceptors(new AuthHandshakeInterceptor());
        
        registry.addHandler(new RealTimeHandler(), "/ws/realtime")
                .setAllowedOrigins("https://app.example.com")
                .withSockJS();
    }
}

// 握手认证拦截器
public class AuthHandshakeInterceptor implements HandshakeInterceptor {
    
    @Override
    public boolean beforeHandshake(ServerHttpRequest request,
            ServerHttpResponse response, WebSocketHandler wsHandler,
            Map attributes) {
        String token = extractToken(request);
        if (token == null || !validateToken(token)) {
            response.setStatusCode(HttpStatus.UNAUTHORIZED);
            return false;
        }
        UserInfo user = getUserInfo(token);
        attributes.put("userId", user.getId());
        attributes.put("username", user.getName());
        return true;
    }
    
    @Override
    public void afterHandshake(ServerHttpRequest request,
            ServerHttpResponse response, WebSocketHandler wsHandler,
            Exception exception) {
    }
}

连接管理与消息处理器

WebSocket连接的生命周期管理是实时通信系统的核心。TextWebSocketHandler提供了连接建立、消息接收和连接关闭的回调接口。连接池使用ConcurrentHashMap保证线程安全,每个连接关联用户ID便于定向推送。

public class RealTimeHandler extends TextWebSocketHandler {
    
    private final ConcurrentHashMap 
        sessions = new ConcurrentHashMap<>();
    
    private final ConcurrentHashMap>
        roomMembers = new ConcurrentHashMap<>();
    
    @Override
    public void afterConnectionEstablished(WebSocketSession session) {
        String userId = (String) session.getAttributes().get("userId");
        sessions.put(userId, session);
        broadcastUserStatus(userId, "online");
    }
    
    @Override
    protected void handleTextMessage(WebSocketSession session, 
            TextMessage message) throws Exception {
        String userId = (String) session.getAttributes().get("userId");
        JsonNode payload = objectMapper.readTree(message.getPayload());
        String type = payload.get("type").asText();
        
        switch (type) {
            case "chat" -> handleChatMessage(userId, payload);
            case "join_room" -> handleJoinRoom(userId, payload);
            case "leave_room" -> handleLeaveRoom(userId, payload);
            case "heartbeat" -> handleHeartbeat(session, payload);
            default -> sendError(session, "未知消息类型: " + type);
        }
    }
    
    @Override
    public void afterConnectionClosed(WebSocketSession session, 
            CloseStatus status) {
        String userId = (String) session.getAttributes().get("userId");
        sessions.remove(userId);
        roomMembers.values().forEach(members -> members.remove(userId));
        broadcastUserStatus(userId, "offline");
    }
    
    @Override
    public void handleTransportError(WebSocketSession session, 
            Throwable exception) {
        try {
            session.close(CloseStatus.SERVER_ERROR);
        } catch (IOException e) {
            log.error("关闭连接失败", e);
        }
    }
    
    public void sendToUser(String targetUserId, Object message) {
        WebSocketSession session = sessions.get(targetUserId);
        if (session != null && session.isOpen()) {
            try {
                session.sendMessage(new TextMessage(
                    objectMapper.writeValueAsString(message)));
            } catch (IOException e) {
                sessions.remove(targetUserId);
            }
        }
    }
    
    public void broadcastToRoom(String roomId, Object message, 
            String excludeUserId) {
        Set members = roomMembers.get(roomId);
        if (members == null) return;
        String json = objectMapper.writeValueAsString(message);
        members.stream()
            .filter(uid -> !uid.equals(excludeUserId))
            .map(sessions::get)
            .filter(session -> session != null && session.isOpen())
            .forEach(session -> {
                try {
                    session.sendMessage(new TextMessage(json));
                } catch (IOException e) {
                    log.error("广播失败", e);
                }
            });
    }
}

心跳保活与断线重连机制

TCP连接的Keep-Alive机制在NAT环境和负载均衡器后不可靠,应用层心跳是检测连接健康状态的必要手段。心跳间隔通常设置为30秒,超过3次未收到响应判定连接断开。客户端断线后实现指数退避重连,避免大量客户端同时重连导致服务端连接风暴。

class WebSocketClient {
    constructor(url, options = {}) {
        this.url = url;
        this.heartbeatInterval = options.heartbeatInterval || 30000;
        this.reconnectMinDelay = options.reconnectMinDelay || 1000;
        this.reconnectMaxDelay = options.reconnectMaxDelay || 30000;
        this.reconnectAttempts = 0;
        this.maxReconnectAttempts = options.maxReconnect || 10;
        this.ws = null;
        this.heartbeatTimer = null;
        this.lastHeartbeatTime = 0;
        this.heartbeatTimeout = this.heartbeatInterval * 3;
    }
    
    connect() {
        this.ws = new WebSocket(this.url);
        
        this.ws.onopen = () => {
            this.reconnectAttempts = 0;
            this.startHeartbeat();
        };
        
        this.ws.onmessage = (event) => {
            const data = JSON.parse(event.data);
            if (data.type === 'heartbeat_ack') {
                this.lastHeartbeatTime = Date.now();
            } else {
                this.onMessage?.(data);
            }
        };
        
        this.ws.onclose = (event) => {
            this.stopHeartbeat();
            if (event.code !== 1000 && this.reconnectAttempts < this.maxReconnectAttempts) {
                this.scheduleReconnect();
            }
        };
    }
    
    startHeartbeat() {
        this.lastHeartbeatTime = Date.now();
        this.heartbeatTimer = setInterval(() => {
            if (Date.now() - this.lastHeartbeatTime > this.heartbeatTimeout) {
                this.ws.close(4000, 'heartbeat timeout');
                return;
            }
            if (this.ws.readyState === WebSocket.OPEN) {
                this.ws.send(JSON.stringify({ type: 'heartbeat', ts: Date.now() }));
            }
        }, this.heartbeatInterval);
    }
    
    stopHeartbeat() {
        if (this.heartbeatTimer) {
            clearInterval(this.heartbeatTimer);
            this.heartbeatTimer = null;
        }
    }
    
    scheduleReconnect() {
        this.reconnectAttempts++;
        const delay = Math.min(
            this.reconnectMinDelay * Math.pow(2, this.reconnectAttempts - 1),
            this.reconnectMaxDelay
        ) + Math.random() * 1000;
        this.reconnectTimer = setTimeout(() => this.connect(), delay);
    }
    
    send(data) {
        if (this.ws?.readyState === WebSocket.OPEN) {
            this.ws.send(JSON.stringify(data));
        }
    }
}

服务端心跳检测与空闲连接清理

服务端需要定时检测空闲连接,清理异常断开的Socket避免资源泄漏。Spring Boot中通过定时任务扫描连接池,对超时未发送心跳的连接强制关闭。

@Component
public class ConnectionCleaner {
    
    private final RealTimeHandler handler;
    private final ConcurrentHashMap lastActiveTime;
    private static final long MAX_IDLE_TIME = 90_000;
    
    @Scheduled(fixedRate = 30000)
    public void cleanIdleConnections() {
        long now = System.currentTimeMillis();
        
        handler.getSessions().forEach((userId, session) -> {
            Long lastTime = lastActiveTime.get(userId);
            if (lastTime == null) {
                lastActiveTime.put(userId, now);
                return;
            }
            
            if (now - lastTime > MAX_IDLE_TIME) {
                try {
                    session.close(CloseStatus.SESSION_NOT_RELIABLE);
                } catch (IOException e) {
                    log.error("关闭空闲连接失败", e);
                }
                lastActiveTime.remove(userId);
            }
        });
    }
    
    public void updateActiveTime(String userId) {
        lastActiveTime.put(userId, System.currentTimeMillis());
    }
}

水平扩展与消息广播优化

单机WebSocket连接数上限受限于JVM内存和文件描述符限制,通常单机支撑5万到10万连接。超出单机容量时需要水平扩展,多实例间消息同步通过Redis Pub/Sub或消息中间件实现。客户端连接经过负载均衡器分发到不同实例,消息广播通过Redis转发到所有实例。

@Component
public class WebSocketMessageBroker {
    
    private final StringRedisTemplate redisTemplate;
    private final RealTimeHandler handler;
    private final ChannelTopic broadcastTopic = new ChannelTopic("ws:broadcast");
    
    @PostConstruct
    public void subscribe() {
        MessageListener listener = (message, pattern) -> {
            String body = new String(message.getBody());
            WsBroadcastMessage msg = new ObjectMapper().readValue(body, 
                WsBroadcastMessage.class);
            
            switch (msg.getScope()) {
                case "user" -> handler.sendToUser(msg.getTarget(), msg.getPayload());
                case "room" -> handler.broadcastToRoom(
                    msg.getTarget(), msg.getPayload(), msg.getExclude());
                case "all" -> handler.getSessions().forEach((uid, session) -> {
                    if (session.isOpen()) {
                        session.sendMessage(new TextMessage(msg.getPayloadJson()));
                    }
                });
            }
        };
        
        RedisMessageListenerContainer container = new RedisMessageListenerContainer();
        container.addMessageListener(listener, broadcastTopic);
        container.start();
    }
    
    public void publish(String scope, String target, Object payload, 
            String excludeUserId) {
        WsBroadcastMessage msg = new WsBroadcastMessage();
        msg.setScope(scope);
        msg.setTarget(target);
        msg.setExclude(excludeUserId);
        msg.setPayload(payload);
        
        redisTemplate.convertAndSend(broadcastTopic.getTopic(), 
            new ObjectMapper().writeValueAsString(msg));
    }
}

WebSocket连接的负载均衡需要使用Sticky Session或IP Hash策略,确保同一客户端的连接始终路由到同一后端实例。Nginx和HAProxy均支持WebSocket协议升级,配置proxy_set_header Upgrade和proxy_set_header Connection “upgrade”即可透传WebSocket握手请求。连接重建时如果原实例已下线,负载均衡器将连接路由到新实例,客户端通过重连机制恢复通信。

原创文章,作者:小编,如若转载,请注明出处:https://www.yunthe.com/websocket-shi-shi-tong-xin-shi-zhan-lian-jie-guan-li-yu-xin/

(0)
小编小编
上一篇 1小时前
下一篇 1小时前

相关推荐