如何搭建一个生产级 WebSocket 服务?这 7 个问题必须要解决

做过后端和实时业务开发的同学应该都有同感:WebSocket 本地调试真的太简单了。随便写几行前后端代码,跑个 Demo 就能实现双向通信,看起来零难度、零 bug。但只要一上生产环境,各种离谱问题就接踵而至,让人防不胜防。
要么是连接毫无征兆自动断开,要么是用户收不到推送消息,高峰期还容易出现服务卡顿、OOM 内存溢出,集群部署后消息广播更是直接失效。很多时候,本地跑得稳稳当当的功能,上线就频繁翻车,排查日志、调试网络、复盘配置,折腾大半天也找不到根因。
其实 WebSocket 本地开发拼的是基础语法,生产落地拼的是工程化细节。握手校验、网关代理、心跳保活、并发安全、消息可靠性、集群扩容、性能防护,每一个环节都是生产环境的必经关卡。
我在多次实时大屏、IM 即时通讯、在线协同项目的落地过程中,踩遍了 WebSocket 线上的各类坑。今天就以 Vue3+Spring Boot 主流技术栈为例,把生产环境必须搞定的 7 个核心问题,结合实战代码、底层原理和解决方案一次性讲透,帮大家彻底避开线上隐患,实现稳定、高可用的 WebSocket 生产级服务。
一、握手阶段的安全和校验——如何防止非法连接?
WebSocket 的连接是始于一次 HTTP 升级请求。但是生产环境不能允许任何客户端随意连接,必须在握手阶段完成身份校验。如果等到连接建立后(如 @OnOpen)再校验,非法连接已经占用了服务端资源。
解决方案:使用 Spring 的 HandshakeInterceptor,在握手发起前拦截请求,解析 Token 并验证身份。同时,浏览器原生 WebSocket API 不支持自定义 Header,因此 Token 通常通过 URL 参数传递。
后端:配置类和握手拦截器
配置 WebSocket 入口:
@Configuration
@EnableWebSocket
public class WebSocketConfig implements WebSocketConfigurer {
@Autowired
private MyWebSocketHandler myWebSocketHandler;
@Autowired
private WebSocketAuthInterceptor webSocketAuthInterceptor;
@Override
public void registerWebSocketHandlers(WebSocketHandlerRegistry registry) {
registry.addHandler(myWebSocketHandler, "/ws")
.setAllowedOrigins("*") // 生产环境需限制具体域名
.addInterceptors(webSocketAuthInterceptor);
}
}
握手拦截器实现:
@Component
public class WebSocketAuthInterceptor implements HandshakeInterceptor {
@Override
public boolean beforeHandshake(ServerHttpRequest request, ServerHttpResponse response,
WebSocketHandler wsHandler, Map<String, Object> attributes) {
if (request instanceof ServletServerHttpRequest) {
HttpServletRequest servletRequest = ((ServletServerHttpRequest) request).getServletRequest();
// 1. 从 URL 参数获取 Token
String token = servletRequest.getParameter("token");
// 2. 验证 Token 逻辑 (解析 JWT 或查询 Redis)
String userId = JwtUtil.validateAndGetUserId(token);
if (userId == null) {
response.setStatusCode(HttpStatus.UNAUTHORIZED);
return false; // 拒绝握手
}
// 3. 把用户信息存入 WebSocket Session 的 attributes 中
// ⚠️ 关键:WebSocket Session 和 HTTP Session 是独立的,必须在这里绑定用户信息
attributes.put("userId", userId);
return true;
}
return false;
}
@Override
public void afterHandshake(ServerHttpRequest request, ServerHttpResponse response, WebSocketHandler wsHandler, Exception exception) {
}
}
前端:建立连接时附带 Token
const userToken = "eyJhbGciOi..."; // 从本地存储获取
const socket = new WebSocket(`wss://api.your-domain.com/ws?token=${userToken}`);
二、网关层代理——为什么 Nginx 总是断开连接?
如果 WebSocket 服务部署在 Nginx 后面,默认配置会导致连接在短时间内被断开。原因在于:Nginx 默认会把 Upgrade 和 Connection 头丢弃,后端收不到协议升级指示;同时 Nginx 的默认读超时(60s)很短,长连接会被主动掐断。
解决方案:配置生产级 Nginx 代理规则,显式转发协议升级头并调大超时时间。
upstream ws_backend {
server 127.0.0.1:8080;
# 生产环境多实例可配置负载均衡策略
}
server {
listen 443 ssl;
server_name ws.your-domain.com;
location /ws/ {
proxy_pass http://ws_backend;
# 核心:显式转发 Upgrade 头
proxy_http_version 1.1;
proxy_set_header Upgrade $http_upgrade;
proxy_set_header Connection "upgrade";
# 超时设置——必须大于服务端心跳间隔(如心跳 25s,这里设 120s)
proxy_read_timeout 120s;
proxy_connect_timeout 60s;
proxy_set_header X-Real-IP $remote_addr;
proxy_set_header X-Forwarded-For $proxy_add_x_forwarded_for;
}
}
proxy_read_timeout是关键。如果在这个时间内没有数据收发,Nginx 会主动断开连接。所以它必须大于前端的心跳间隔。
三、心跳保活和指数退避重连——如何防止静默断开?
TCP 连接可能因 NAT 超时、防火墙策略等被中间设备回收,而应用层毫无感知(静默断开)。同时,比如用户进电梯网络被瞬间切断时,如果没有重连机制,用户就会错过所有推送消息。如果断线后瞬间疯狂重连,又会引发服务器惊群效应。
解决方案:双端 Ping/Pong 心跳保活 + 前端指数退避重连。
前端实现 (Vue 3 封装)
封装一个具备心跳和指数退避重连的 Composable:
// composables/useWebSocket.ts
import { ref, onUnmounted } from 'vue'
export function useWebSocket(url: string) {
const ws = ref<WebSocket | null>(null)
const isConnected = ref(false)
let heartbeatTimer: number | null = null
let reconnectTimer: number | null = null
let reconnectCount = 0
const MAX_RECONNECT = 6
const connect = () => {
ws.value = new WebSocket(url)
ws.value.onopen = () => {
isConnected.value = true
reconnectCount = 0 // 重置重连次数
startHeartbeat()
// 连接成功后,主动拉取离线消息
fetchOfflineMessages()
}
ws.value.onmessage = (event) => {
const data = JSON.parse(event.data)
if (data.type === 'pong') return // 收到心跳回复
if (data.type === 'ack') {
// 处理 ACK 确认...
return
}
// 处理业务消息,并回复 ACK
handleBusinessMessage(data)
ws.value.send(JSON.stringify({ type: 'ack', messageId: data.messageId }))
}
ws.value.onclose = () => {
isConnected.value = false
stopHeartbeat()
reconnect() // 触发重连
}
ws.value.onerror = () => {
ws.value?.close()
}
}
const startHeartbeat = () => {
heartbeatTimer = setInterval(() => {
if (ws.value?.readyState === WebSocket.OPEN) {
ws.value.send(JSON.stringify({ type: 'ping' }))
}
}, 25000) // 25s 心跳
}
const reconnect = () => {
if (reconnectCount < MAX_RECONNECT) {
reconnectCount++
// 指数退避: 1s, 2s, 4s, 8s... 防止服务重启时惊群效应
const delay = Math.pow(2, reconnectCount - 1) * 1000
reconnectTimer = setTimeout(connect, delay)
}
}
const stopHeartbeat = () => {
if (heartbeatTimer) clearInterval(heartbeatTimer)
if (reconnectTimer) clearTimeout(reconnectTimer)
}
onUnmounted(() => {
stopHeartbeat()
ws.value?.close()
})
return { connect, isConnected, ws }
}
2. 后端实现:心跳响应 + 定时清理死连接
在 MyWebSocketHandler 中处理心跳,并增加定时清理任务:
@Component
public class MyWebSocketHandler extends TextWebSocketHandler {
@Override
protected void handleTextMessage(WebSocketSession session, TextMessage message) throws Exception {
String payload = message.getPayload();
JSONObject obj = JSON.parseObject(payload);
if ("ping".equals(obj.getString("type"))) {
// 收到 Ping,立即回复 Pong 并更新心跳时间
session.sendMessage(new TextMessage("{\"type\":\"pong\"}"));
WebSocketSessionManager.updateHeartbeat(session.getId());
return;
}
// ... 处理业务消息和 ACK 逻辑
}
}
// 定时清理超过 90s 未心跳的死连接
@Scheduled(fixedDelay = 30000)
public void cleanIdleSessions() {
long now = System.currentTimeMillis();
WebSocketSessionManager.sessions.values().removeIf(session -> {
Long lastPing = WebSocketSessionManager.heartbeatMap.get(session.getId());
return lastPing != null && (now - lastPing) > 90_000;
});
}
- 心跳保活:客户端每 25 秒发一个极小的 Ping 包,服务端秒回 Pong。如果服务端连续 90 秒没听到动静,直接判定死亡,主动清理连接,绝不留着吃内存。
- 指数退避重连: 网络瞬断时,如果所有客户端瞬间疯狂重连,会直接把服务器打死(惊群效应)。我们采用“指数退避”:第 1 次断开等 1 秒重试,第 2 次等 2 秒,第 4 次等 4 秒……越来越慢,直到上限。这样既保证了用户网络恢复时能自动连上,又给了服务器喘息的时间。
四、消息发送的线程安全性——为什么并发推送会报错?
javax.websocket.Session不是线程安全的。当多个业务线程(比如系统通知、聊天消息并发)同时调用session.sendMessage()时,会抛出IllegalStateException: The remote endpoint was in state [TEXT_PARTIAL_WRITING]。
解决方案:对单个 Session 加锁,串行化发送。同时,设计支持同一用户多端登录的 Session 管理器。
@Component
public class WebSocketSessionManager {
// 支持多端登录:一个用户可能有多个 Session (PC、App 同时在线)
private static final ConcurrentMap<String, Set<WebSocketSession>> userSessions = new ConcurrentHashMap<>();
public static final ConcurrentMap<String, Long> heartbeatMap = new ConcurrentHashMap<>();
public static void add(String userId, WebSocketSession session) {
userSessions.computeIfAbsent(userId, k -> ConcurrentHashMap.newKeySet()).add(session);
heartbeatMap.put(session.getId(), System.currentTimeMillis());
}
public static void remove(String userId, WebSocketSession session) {
Set<WebSocketSession> sessions = userSessions.get(userId);
if (sessions != null) {
sessions.remove(session);
heartbeatMap.remove(session.getId());
}
}
public static void sendMessageToUser(String userId, String message) {
Set<WebSocketSession> sessions = userSessions.get(userId);
if (sessions == null || sessions.isEmpty()) return;
sessions.removeIf(session -> !session.isOpen());
for (WebSocketSession session : sessions) {
safeSend(session, message);
}
}
private static void safeSend(WebSocketSession session, String message) {
// 锁的粒度:仅锁当前 Session,不同用户互不影响
synchronized (session) {
if (session.isOpen()) {
try {
session.sendMessage(new TextMessage(message));
} catch (IOException e) {
log.error("发送失败", e);
}
}
}
}
}
五、消息可靠性和离线消息(ACK + 事务解耦)
断线期间服务端推的消息会丢失;此外,很多开发者习惯在数据库事务内进行 WebSocket 推送,这会导致网络 IO 阻塞长期占用数据库连接,甚至推送异常导致事务回滚,数据丢失。
解决方案:ACK 确认机制 + 持久化兜底 + 事务和推送解耦。
后端:事务和推送严格解耦
@Service
public class MessageService {
@Autowired
private MessageDao messageDao;
public void saveAndPush(Message msg) {
// 1. 事务方法:仅负责数据落库
saveInTransaction(msg);
// 2. 事务外执行推送,推送失败不影响数据落库
try {
WebSocketSessionManager.sendMessageToUser(msg.getReceiverId(), JSON.toJSONString(msg));
} catch (Exception e) {
log.error("推送失败,等待用户上线后拉取", e);
}
}
@Transactional(rollbackFor = Exception.class)
public void saveInTransaction(Message msg) {
messageDao.insert(msg);
}
}
后端:ACK 确认和离线队列处理
// 在 Handler 中处理 ACK
@Override
protected void handleTextMessage(WebSocketSession session, TextMessage message) throws Exception {
JSONObject obj = JSON.parseObject(message.getPayload());
if ("ack".equals(obj.getString("type"))) {
String messageId = obj.getString("messageId");
// 收到 ACK,从 Redis 离线队列中删除该消息,防止重复推送
redisTemplate.opsForHash().delete("ws:offline:" + userId, messageId);
return;
}
}
// 推送时的逻辑:如果用户不在线,先存 Redis 再推
public void pushMessage(String userId, Message msg) {
String msgJson = JSON.toJSONString(msg);
// 1. 先存入 Redis 离线队列 (以 messageId 为 key)
redisTemplate.opsForHash().put("ws:offline:" + userId, msg.getId(), msgJson);
// 2. 尝试推送
WebSocketSessionManager.sendMessageToUser(userId, msgJson);
// 如果推送成功,客户端会回 ACK,触发 Redis 删除;
// 如果客户端不在线,推送失败,消息留在 Redis 中,等待上线拉取。
}
前端:上线拉取离线消息
ws.value.onopen = async () => {
startHeartbeat();
// 1. 连接成功后,先通过 HTTP 请求拉取所有离线消息
const res = await fetchOfflineMessages();
res.data.forEach(msg => {
handleBusinessMessage(msg);
// 2. 对拉取的离线消息也发送 ACK,让服务端清理 Redis
ws.value.send(JSON.stringify({ type: 'ack', messageId: msg.id }));
});
};
六、水平扩展——多实例部署时如何广播?
生产环境几乎都是集群部署。用户 A 连在节点 1,用户 B 连在节点 2。节点 1 内部广播消息时,节点 2 上的 B 根本收不到。
解决方案:引入 Redis Pub/Sub,打通多节点的消息流转。
Redis 配置和监听器
@Configuration
public class RedisConfig {
@Bean
public RedisMessageListenerContainer redisMessageListenerContainer(
RedisConnectionFactory connectionFactory, RedisMessageSubscriber subscriber) {
RedisMessageListenerContainer container = new RedisMessageListenerContainer();
container.setConnectionFactory(connectionFactory);
// 订阅 WebSocket 广播频道
container.addMessageListener(subscriber, new ChannelTopic("ws:broadcast"));
return container;
}
}
跨节点消息分发逻辑
@Service
public class RedisMessageSubscriber implements MessageListener {
@Override
public void onMessage(Message message, byte[] pattern) {
String payload = new String(message.getBody());
WebSocketMessageDTO dto = JSON.parseObject(payload, WebSocketMessageDTO.class);
// 收到 Redis 广播,推送给本实例连接的用户
if (dto.getTargetUserId() != null) {
WebSocketSessionManager.sendMessageToUser(dto.getTargetUserId(), dto.getContent());
} else {
// 全局广播
WebSocketSessionManager.userSessions.keySet().forEach(uid ->
WebSocketSessionManager.sendMessageToUser(uid, dto.getContent())
);
}
}
}
业务层调用方式:当需要推送时,业务代码只需向 Redis 发布消息,不再直接调用本地 Session 推送:
public void sendCrossNodeMessage(String userId, String content) {
WebSocketMessageDTO dto = new WebSocketMessageDTO(userId, content);
redisTemplate.convertAndSend("ws:broadcast", JSON.toJSONString(dto));
}
七、性能、资源治理和业务保护
高并发连接耗尽内存;恶意客户端疯狂建连或发消息打死服务;同账号多端登录冲突。
解决方案:多维度资源治理和保护机制。
单点登录踢出(同端互斥)
同账号新设备登录时,主动踢掉旧连接,避免消息多端重复消费:
public static void add(String userId, WebSocketSession session) {
Set<WebSocketSession> existingSessions = userSessions.get(userId);
if (existingSessions != null) {
for (Session oldSession : existingSessions) {
if (oldSession.isOpen()) {
// 通知客户端被踢下线
safeSend(oldSession, "{\"type\":\"KICK_OUT\"}");
oldSession.close(CloseStatus.NORMAL);
}
}
existingSessions.clear();
}
userSessions.computeIfAbsent(userId, k -> ConcurrentHashMap.newKeySet()).add(session);
}
限流策略防恶意刷屏
使用 Guava RateLimiter 限制单客户端发送消息频率:
// 在 Handler 中维护每个 Session 的限流器
private final ConcurrentMap<String, RateLimiter> limiters = new ConcurrentHashMap<>();
@Override
protected void handleTextMessage(WebSocketSession session, TextMessage message) throws Exception {
RateLimiter limiter = limiters.computeIfAbsent(session.getId(), k -> RateLimiter.create(5)); // 每秒 5 条
if (!limiter.tryAcquire()) {
safeSend(session, "{\"type\":\"ERROR\",\"msg\":\"请求过于频繁\"}");
// 触发限流,直接关闭恶意连接
session.close(CloseStatus.POLICY_VIOLATION);
return;
}
// ... 继续处理正常消息
}
虚拟线程撑住海量连接 (Spring Boot 3.2+)
传统平台线程每个连接占用约 1MB,Tomcat 默认 200 线程很容易耗尽。启用虚拟线程后,单机可轻松支撑数万连接:
# application.yml spring: threads: virtual: enabled: true
⚠️ 注意:synchronized 块会导致虚拟线程被固定在载体线程上,失去调度优势。如果使用虚拟线程,问题四中的 safeSend 建议改用 ReentrantLock:
private static final ConcurrentMap<String, ReentrantLock> sessionLocks = new ConcurrentHashMap<>();
private static void safeSend(WebSocketSession session, String message) {
ReentrantLock lock = sessionLocks.computeIfAbsent(session.getId(), k -> new ReentrantLock());
lock.lock();
try {
if (session.isOpen()) {
session.sendMessage(new TextMessage(message));
}
} catch (IOException e) {
log.error("发送失败", e);
} finally {
lock.unlock();
}
}
内存泄漏防御
在 afterConnectionClosed 中必须严格执行清理逻辑:
@Override
public void afterConnectionClosed(WebSocketSession session, CloseStatus status) throws Exception {
String userId = (String) session.getAttributes().get("userId");
if (userId != null) {
WebSocketSessionManager.remove(userId, session);
limiters.remove(session.getId()); // 清理限流器
sessionLocks.remove(session.getId()); // 清理锁
}
}
在实际落地时,如果你的业务非常复杂(如大型多人在线游戏、复杂 IM),建议可以直接考虑使用 Netty 底层框架,或者使用专为 WebSocket 设计的网关组件。但对于大多数 Spring Boot + Vue 的常规业务系统,按照咱们这篇文章的方案,已经完全足以撑起生产环境的高并发和高可用需求。
以上关于如何搭建一个生产级 WebSocket 服务?这 7 个问题必须要解决的文章就介绍到这了,更多相关内容请搜索码云笔记以前的文章或继续浏览下面的相关文章,希望大家以后多多支持码云笔记。
如若内容造成侵权/违法违规/事实不符,请将相关资料发送至 admin@mybj123.com 进行投诉反馈,一经查实,立即处理!
重要:如软件存在付费、会员、充值等,均属软件开发者或所属公司行为,与本站无关,网友需自行判断
码云笔记 » 如何搭建一个生产级 WebSocket 服务?这 7 个问题必须要解决
微信
支付宝