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/