资讯详情

资讯详情

WebSocket实时聊天系统从零搭建与高可用实践

简介这是一份面向计算机专业本科生的毕业设计/课程设计实战项目基于Python与Vue.js实现轻量级实时在线聊天系统解决传统HTTP轮询在即时通信场景下的高延迟与低效问题。资源包共31个文件涵盖15个JavaScript逻辑文件含WebSocket连接管理与消息处理、2个Vue组件文件构建响应式聊天界面、2个Stylus样式文件、3个SVG图标资源及配套配置文件.babelrc、.editorconfig等整体仅134KB结构精简、开箱即用。目前已有35人学习下载适合初学者理解全双工通信原理与前后端协同开发流程。读者可直接运行服务端index.js与前端Vue项目获得完整可交互的聊天室实例代码分层清晰含路由配置、Webpack构建脚本、环境变量管理及README说明便于快速复现、调试与二次开发。1. 为什么用 WebSocket 做实时在线聊天系统不是轮询、长连接或 SSE你写过一个“每秒发一次 AJAX 请求查新消息”的聊天界面吗用户刚打完字后端还没存进数据库前端就弹出“收到 0 条新消息”——这种体验不是延迟是幻觉。真实项目里轮询在 300 用户并发时 CPU 就飙到 95%SSE 在 Chrome 旧版和 iOS Safari 上断连不重连而 WebSocket 从 TCP 握手完成那一刻起就建立了一条双向、低开销、服务端可主动推的全双工通道。这不是“更时髦的选择”而是当消息到达时效要求 ≤800ms、单日活跃用户 ≥5000、需支持离线消息补推与多端状态同步时工程上唯一能稳住的协议底座。本方案不依赖任何 SaaS 推送服务不封装黑盒 SDK所有逻辑可见、可调、可压测从握手鉴权、心跳保活、消息广播策略到断线重连的退避算法与会话恢复机制全部基于原生 WebSocket API Node.jsv18.17 Redis7.2落地。适合正在从 HTTP 短连接迁移到实时通信的中后台工程师、独立开发者以及需要把“在线状态”“已读回执”“消息撤回”真正跑通的 ToB SaaS 团队。2. 从零搭起 WebSocket 服务端握手鉴权、路由分发与连接池管理WebSocket 不是“开个端口就能聊”它本质是 HTTP 升级协议第一次请求必须走标准 HTTP 握手且必须校验来源、身份、权限。很多翻车都发生在这一秒——比如没校验 Origin 导致跨域劫持或把 token 放在 URL 里被 CDN 缓存泄露。我们用ws库v8.16.0而非socket.io就是为了剥离抽象层看清每一帧怎么来、怎么走。2.1 握手阶段强制校验三要素Origin、Token 与 User-Agent// server.js import { WebSocketServer } from ws; import { createHash } from crypto; const wss new WebSocketServer({ port: 8080, verifyClient: (info, done) { const { origin, req } info; const url new URL(req.url, http://localhost); // 1. 拦截非法 Origin防 CSRF if (![https://your-app.com, http://localhost:3000].includes(origin)) { return done(false, 403, Forbidden Origin); } // 2. 提取并校验 JWT放在 Sec-WebSocket-Protocol 头比 URL 安全 const protocol req.headers[sec-websocket-protocol]; if (!protocol || !protocol.startsWith(auth-)) { return done(false, 401, Missing auth header); } const token protocol.split(auth-)[1]; try { // 这里调用你的 JWT 验证函数含签名校验过期检查 const payload verifyJWT(token); if (!payload.userId || !payload.exp || payload.exp Date.now() / 1000) { return done(false, 401, Invalid or expired token); } // 把用户信息挂到 info 上供后续 use info.userId payload.userId; info.deviceId payload.deviceId || web- createHash(md5).update(origin).digest(hex).slice(0, 8); done(true); } catch (err) { done(false, 401, Token verification failed); } } });提示Sec-WebSocket-Protocol是 WebSocket 标准头浏览器端可通过new WebSocket(url, [auth-xxx])传入比拼接?tokenxxx更安全——CDN 和代理不会缓存该头且服务端可统一拦截未带该头的连接。2.2 连接池用 Map 管理在线用户支持按用户 ID、设备 ID、房间 ID 三维度索引光有连接对象不够你得知道“张三的 iPhone 和 Chrome 同时在线谁该收这条消息”。我们不用全局广播而是建三层索引userMap: MapuserId, SetWebSocket—— 用户所有在线设备deviceMap: MapdeviceId, WebSocket—— 单设备精准控制用于踢下线、推送设备专属通知roomMap: MaproomId, SetWebSocket—— 房间级广播群聊// connection-manager.js export class ConnectionManager { constructor() { this.userMap new Map(); this.deviceMap new Map(); this.roomMap new Map(); } add(ws, userId, deviceId, roomId null) { // 1. 绑定用户 if (!this.userMap.has(userId)) this.userMap.set(userId, new Set()); this.userMap.get(userId).add(ws); // 2. 绑定设备强覆盖同一设备重复登录则踢旧 this.deviceMap.set(deviceId, ws); // 3. 绑定房间可多次加入不同房间 if (roomId) { if (!this.roomMap.has(roomId)) this.roomMap.set(roomId, new Set()); this.roomMap.get(roomId).add(ws); } // 4. 绑定关闭钩子 ws.on(close, () this.remove(ws, userId, deviceId, roomId)); } remove(ws, userId, deviceId, roomId) { if (this.userMap.has(userId)) { this.userMap.get(userId).delete(ws); if (this.userMap.get(userId).size 0) this.userMap.delete(userId); } this.deviceMap.delete(deviceId); if (roomId this.roomMap.has(roomId)) { this.roomMap.get(roomId).delete(ws); if (this.roomMap.get(roomId).size 0) this.roomMap.delete(roomId); } } broadcastToUser(userId, data) { const connections this.userMap.get(userId); if (!connections) return; connections.forEach(ws { if (ws.readyState WebSocket.OPEN) { ws.send(JSON.stringify(data)); } }); } broadcastToRoom(roomId, data) { const connections this.roomMap.get(roomId); if (!connections) return; connections.forEach(ws { if (ws.readyState WebSocket.OPEN) { ws.send(JSON.stringify(data)); } }); } } export const connMgr new ConnectionManager();参数说明ws.readyState必须显式判断因为close事件触发前连接可能已断但对象未销毁broadcastToUser用于私聊/系统通知broadcastToRoom用于群聊二者不可混用——群聊消息若误走用户广播会导致消息重复一人多设备时每个设备都收一遍。2.3 消息路由解析客户端帧分发到业务处理器客户端发来的不是裸文本而是结构化 JSON含type、payload、seq消息序号用于去重与排序。服务端不做业务逻辑只做路由// message-router.js wss.on(connection, (ws, req) { const { userId, deviceId } req; ws.on(message, (data) { let msg; try { msg JSON.parse(data.toString()); } catch (e) { ws.send(JSON.stringify({ type: error, code: INVALID_JSON, message: Malformed JSON })); return; } // 统一添加上下文 msg.context { userId, deviceId, timestamp: Date.now() }; // 路由分发 switch (msg.type) { case JOIN_ROOM: handleJoinRoom(ws, msg.payload); break; case SEND_MESSAGE: handleSendMessage(ws, msg.payload); break; case HEARTBEAT: handleHeartbeat(ws); break; case SET_STATUS: handleSetStatus(ws, msg.payload); break; default: ws.send(JSON.stringify({ type: error, code: UNKNOWN_TYPE, message: Unsupported type: ${msg.type} })); } }); });血泪经验message事件回调里不能 await 异步操作再 send否则会阻塞整个连接的消息处理队列。所有耗时操作如 DB 写入、Redis 查询必须process.nextTick(() {...})或setTimeout(..., 0)脱离当前事件循环保证ws.send()实时响应。3. 客户端 WebSocket 封装自动重连、心跳保活与消息队列浏览器 WebSocket 是“脆弱”的网络抖动、休眠唤醒、Chrome 后台节流都会导致 silent disconnect。直接ws.onclose监听远远不够——你得让前端像 TCP 一样“感知断连、主动重试、平滑恢复”。3.1 重连策略指数退避 最大重试次数 状态锁定// client/ws-client.js class WsClient { constructor(url, options {}) { this.url url; this.options { maxReconnectAttempts: 10, baseDelayMs: 1000, maxDelayMs: 30000, reconnectOnClose: true, ...options }; this.ws null; this.reconnectAttempts 0; this.isConnecting false; this.messageQueue []; // 断线期间发的消息暂存 } connect() { if (this.isConnecting || this.ws?.readyState WebSocket.OPEN) return; this.isConnecting true; this.ws new WebSocket(this.url, [auth-${this.getToken()}]); this.ws.onopen () { console.log([WS] Connected); this.reconnectAttempts 0; this.isConnecting false; // 重发断线期间积压的消息 this.flushQueue(); // 启动心跳 this.startHeartbeat(); }; this.ws.onmessage (event) { const data JSON.parse(event.data); this.handleMessage(data); }; this.ws.onclose (event) { console.warn([WS] Closed: ${event.code} ${event.reason}); this.isConnecting false; if (this.options.reconnectOnClose this.reconnectAttempts this.options.maxReconnectAttempts) { const delay Math.min( this.options.baseDelayMs * Math.pow(2, this.reconnectAttempts), this.options.maxDelayMs ); setTimeout(() this.connect(), delay); this.reconnectAttempts; } }; this.ws.onerror (error) { console.error([WS] Error:, error); this.isConnecting false; }; } send(data) { if (this.ws?.readyState WebSocket.OPEN) { this.ws.send(JSON.stringify(data)); } else { this.messageQueue.push(data); } } flushQueue() { while (this.messageQueue.length 0) { this.send(this.messageQueue.shift()); } } startHeartbeat() { this.heartbeatInterval setInterval(() { if (this.ws?.readyState WebSocket.OPEN) { this.ws.send(JSON.stringify({ type: HEARTBEAT, timestamp: Date.now() })); } }, 30000); // 30s 发一次 } handleMessage(data) { switch (data.type) { case PONG: // 心跳响应更新最后活动时间 this.lastPong Date.now(); break; case MESSAGE: // 业务消息交由上层处理 this.onMessage?.(data); break; default: this.onUnknown?.(data); } } }注意reconnectAttempts计数必须在onopen里清零否则重连成功后下次断连会立刻用最大延迟messageQueue是救命稻草——用户快速输入 5 条消息后断网重连后这 5 条必须按序发出不能丢也不能乱序。3.2 心跳机制实现客户端发 PING服务端回 PONG超时即断连WebSocket 没有内置心跳必须自己定义帧格式。关键点在于服务端不能只靠ping事件响应必须记录每个连接的最后 PONG 时间并在定时器里主动检查。// server.js 补充 const PONG_TIMEOUT_MS 45000; // 45s 内没收到 PONG视为断连 const pongTimers new Map(); // MapWebSocket, Timeout function handleHeartbeat(ws) { // 收到 PING立即回 PONG ws.send(JSON.stringify({ type: PONG, timestamp: Date.now() })); // 清除旧 timer设新 timer if (pongTimers.has(ws)) { clearTimeout(pongTimers.get(ws)); } const timer setTimeout(() { console.warn([WS] No PONG from ${ws._socket?.remoteAddress}, closing); ws.terminate(); // 强制终止不走 close 流程 }, PONG_TIMEOUT_MS); pongTimers.set(ws, timer); } // 连接关闭时清理 timer ws.on(close, () { if (pongTimers.has(ws)) { clearTimeout(pongTimers.get(ws)); pongTimers.delete(ws); } });玄学细节ws.terminate()比ws.close()更彻底——它直接关底层 socket不等缓冲区清空避免因对方卡死导致连接假活PONG_TIMEOUT_MS必须 客户端心跳间隔30s留 15s 容忍网络抖动但不能 60s否则 Nginx 默认proxy_read_timeout会先断你。4. 消息可靠性保障服务端去重、离线存储与已读回执WebSocket 是“尽力而为”不是“确保送达”。用户切后台、手机锁屏、WiFi 切 4G都会导致消息丢失。要达到“微信级”体验必须叠加三重保险服务端幂等去重、Redis 离线队列、客户端已读状态同步。4.1 服务端消息去重用 Redis SETNX 做分布式幂等锁客户端可能因重连、网络重传把同一条消息发两次。服务端必须识别并丢弃// server/handlers/send-message.js import { redisClient } from ../redis.js; async function handleSendMessage(ws, payload) { const { messageId, roomId, content, senderId } payload; // 1. 用 messageId senderId 生成唯一 key const dedupeKey msg:dedupe:${senderId}:${messageId}; // 2. Redis SETNX存在则返回 0已处理不存在则设值并返回 1首次处理 const isDuplicate await redisClient.set(dedupeKey, 1, EX, 3600, NX); if (!isDuplicate) { console.warn([DEDUPE] Duplicate message ignored: ${messageId} from ${senderId}); return; } // 3. 存入消息主表MySQL/PostgreSQL await db.query( INSERT INTO messages (id, room_id, sender_id, content, created_at) VALUES (?, ?, ?, ?, NOW()), [messageId, roomId, senderId, content] ); // 4. 广播给房间内所有在线用户 connMgr.broadcastToRoom(roomId, { type: MESSAGE, payload: { id: messageId, roomId, content, senderId, timestamp: Date.now() } }); // 5. 推送到离线队列见 4.2 await pushToOfflineQueue(roomId, messageId, senderId); }参数说明EX 3600表示 key 1 小时后自动过期避免 Redis 内存无限增长NX是原子操作无竞态messageId必须由客户端生成推荐uuidv4或snowflake服务端不生成——否则重连后客户端无法保证重发消息 ID 不变。4.2 离线消息队列用 Redis Sorted Set 实现按时间排序的拉取队列当用户不在线消息不能丢但也不能全量推——要支持“上线后只拉取未读的 50 条”。Sorted Set 天然支持按 score时间戳范围查询// redis.js export async function pushToOfflineQueue(roomId, messageId, senderId) { const key offline:${roomId}; const score Date.now(); // 用时间戳作 score天然有序 await redisClient.zAdd(key, { score, value: messageId }); // 设置过期时间避免长期堆积 await redisClient.expire(key, 24 * 3600); // 24小时 } export async function getOfflineMessages(roomId, lastReadId null, limit 50) { const key offline:${roomId}; let start 0; let end limit - 1; if (lastReadId) { // 查找 lastReadId 的 rank从下一位开始取 const rank await redisClient.zRank(key, lastReadId); if (rank ! null) { start rank 1; end start limit - 1; } } return await redisClient.zRange(key, start, end, { BY: SCORE, REV: true }); } // 拉取后从队列中移除已读消息幂等 export async function markOfflineAsRead(roomId, messageId) { const key offline:${roomId}; await redisClient.zRem(key, messageId); }避坑zRange默认按字典序必须加{ BY: SCORE, REV: true }才按时间倒序zRem要在客户端确认“已展示”后再调不能在onmessage里立刻删——否则用户切后台再切回消息就没了。4.3 已读回执客户端上报 服务端聚合统计“对方正在输入…” 是心理暗示“已读”是确定性反馈。实现分两步客户端在消息 DOM 渲染完成后上报READ_ACK服务端聚合并广播// client function markAsRead(messageId) { ws.send(JSON.stringify({ type: READ_ACK, payload: { messageId, roomId, userId: currentUser.id } })); } // server const readReceipts new Map(); // MaproomId, MapmessageId, SetuserId function handleReadAck(ws, payload) { const { messageId, roomId, userId } payload; if (!readReceipts.has(roomId)) { readReceipts.set(roomId, new Map()); } const msgMap readReceipts.get(roomId); if (!msgMap.has(messageId)) { msgMap.set(messageId, new Set()); } msgMap.get(messageId).add(userId); // 广播已读人数变化仅通知房间内其他成员 connMgr.broadcastToRoom(roomId, { type: READ_RECEIPT_UPDATE, payload: { messageId, readCount: msgMap.get(messageId).size, readerIds: Array.from(msgMap.get(messageId)) } }); }提示已读回执不要求 100% 实时可以加 500ms 延迟上报防滚动时频繁触发用setTimeoutclearTimeout做防抖。5. 避坑指南WebSocket 实时聊天系统 5 个高频翻车点线上环境永远比本地开发残酷。以下是我在 3 个百万级 DAU 项目中踩过的真坑每一条都附带复现方式、根因分析和可落地的修复代码。5.1 现象Chrome 浏览器切后台 5 分钟后WebSocket 自动断开且onclose事件不触发原因Chrome 为省电对后台标签页的 WebSocket 连接执行 aggressive throttlingsetInterval心跳失效且不触发标准close事件连接进入“半死”状态readyState 1但实际不可用。解决客户端增加visibilitychange监听切后台时主动ws.close()切回前台时强制重连服务端配合ping/pong超时检测。// client document.addEventListener(visibilitychange, () { if (document.hidden) { console.log([VISIBILITY] Tab hidden, closing WS); ws.close(); // 主动关闭触发 onclose } else { console.log([VISIBILITY] Tab visible, reconnecting); wsClient.connect(); // 调用你的重连方法 } });5.2 现象高并发下服务端 CPU 100%ws.send()大量报错WebSocket is not open原因ws.send()是同步非阻塞调用但底层 socket 缓冲区满时会排队当大量连接同时发送Node.js 事件循环被send调用占满close事件无法及时处理导致ws.readyState滞后于真实状态。解决对ws.send()做节流 错误兜底失败时降级为setTimeout(send, 0)function safeSend(ws, data) { try { if (ws.readyState WebSocket.OPEN) { ws.send(data); } else { console.warn([SAFE_SEND] WS not open, state${ws.readyState}); setTimeout(() safeSend(ws, data), 10); } } catch (err) { console.error([SAFE_SEND] Send failed:, err); setTimeout(() safeSend(ws, data), 100); } }5.3 现象iOS Safari 上 WebSocket 频繁断连错误码 1006原因iOS Safari 对 WebSocket 连接数有限制通常 6 个且页面内存超限时会静默 kill 连接1006是连接异常终止无close事件。解决客户端限制单页最多 1 个 WebSocket 连接禁用多实例服务端增加Connection: keep-alive响应头并在握手响应中加Sec-WebSocket-Extensions: permessage-deflate启用压缩减少流量// server.js verifyClient 中追加 done(true, { headers: { Connection: keep-alive, Sec-WebSocket-Extensions: permessage-deflate } });5.4 现象Redis 内存暴涨offline:*key 占用 20GB原因离线队列未设置 TTL或markOfflineAsRead调用失败如客户端崩溃导致消息永久堆积。解决所有离线队列 key 必须EXPIRE且服务端启动时扫描过期 key 并清理# 启动脚本中加一行 redis-cli --scan --pattern offline:* | xargs -I {} redis-cli expire {} 864005.5 现象用户 A 撤回消息用户 B 却看到“消息已撤回”延迟 3 秒原因撤回指令走 WebSocket 广播但用户 B 的客户端未做消息本地缓存收到撤回指令时原消息 DOM 已被 GC无法定位修改。解决客户端维护messageCache: Mapid, Element所有消息渲染后存入撤回时通过messageCache.get(id)?.remove()直接操作 DOM// client const messageCache new Map(); function renderMessage(msg) { const el document.createElement(div); el.dataset.messageId msg.id; el.innerHTML span classcontent${msg.content}/span; messageList.appendChild(el); messageCache.set(msg.id, el); // 关键存 DOM 引用 } function handleRevoke(msg) { const el messageCache.get(msg.messageId); if (el) { el.innerHTML span classrevoked[消息已撤回]/span; } }6. 生产验证技巧用 wrk 自定义脚本压测 WebSocket 连接与消息吞吐写完代码不压测等于没写。HTTP 压测工具如 ab、jmeter不支持 WebSocket必须用wrk配 Lua 脚本。以下是我验证 10K 并发连接 5K 消息/秒的完整流程所有命令可直接复制运行。6.1 准备压测环境部署最小化服务与监控先确保服务监听0.0.0.0:8080并暴露/health接口返回{ status: ok, connections: 1234 }从wss.clients.size获取。用pm2启动并监控内存# pm2 start ecosystem.config.js # ecosystem.config.js module.exports { apps: [{ name: chat-ws, script: ./server.js, instances: 2, // 启动 2 进程利用多核 exec_mode: cluster, watch: false, max_memory_restart: 512M, env: { NODE_ENV: production, REDIS_URL: redis://127.0.0.1:6379 } }] };6.2 编写 wrk WebSocket 脚本模拟真实用户行为wrk本身不支持 WebSocket但可通过--script加载 Lua 脚本实现。创建ws-bench.lua-- ws-bench.lua local wrk require(wrk) local websocket require(websocket) -- 全局连接池 local connections {} -- 初始化建立 1000 个连接可调 function setup(thread) local clients {} for i 1, 1000 do local ws websocket.new(ws://localhost:8080, { protocols {auth- .. math.random(10000, 99999)} }) table.insert(clients, ws) end thread:set(clients, clients) end -- 运行每个连接每秒发 5 条消息 function init(args) requests 0 end function request() local clients wrk.thread:get(clients) local idx math.random(#clients) local ws clients[idx] if ws:connected() then ws:send({type:HEARTBEAT}) requests requests 1 if requests % 5 0 then ws:send(({type:SEND_MESSAGE,payload:{roomId:room1,content:msg-%d}}):format(requests)) end end end -- 关闭退出时清理 function done(summary, latency, requests) print(Total requests:, summary.requests) print(Requests/sec:, summary.requests / summary.duration) end注意math.random(10000, 99999)模拟不同用户 token避免服务端鉴权缓存requests % 5控制消息频率避免压垮服务端广播逻辑。6.3 执行压测与关键指标解读# 安装 wrkmacOS brew install wrk # 安装 lua-websocket需提前装好 Lua 5.1 luarocks install lua-websocket # 运行压测100 并发线程持续 60 秒使用自定义脚本 wrk -t100 -d60s -s ws-bench.lua http://localhost:8080 # 同时监控服务端 pm2 monit # 看 CPU、内存、进程数 redis-cli info memory | grep used_memory_human # 看 Redis 内存 curl http://localhost:8080/health # 看实时连接数关键指标阈值16C32G 服务器指标健康值预警值危险值wss.clients.size≤ 80008000–10000 10000OOM 风险Redisused_memory_human≤ 2G2–4G 4G需扩容或调 TTLwrkRequests/sec≥ 40003000–4000 3000服务端瓶颈平均延迟latency≤ 150ms150–300ms 300ms用户可感知卡顿我一般会在 5K 连接、3K 消息/秒下稳定运行 24 小时期间观察 GC 次数node --trace-gc、Redis 连接数redis-cli client list \| wc -l和TIME_WAIT连接数netstat -an \| grep :8080 \| grep TIME_WAIT \| wc -l。如果TIME_WAIT 30000说明内核net.ipv4.tcp_fin_timeout太小需调大至 30。最后说句实在的WebSocket 聊天系统最难的从来不是“连上”而是“连得久、发得准、断得明、恢复快”。你不需要堆砌 fancy 功能先把握手鉴权、心跳保活、消息去重、离线队列这四块地基夯死剩下的——已读回执、消息搜索、语音转文字都是站在巨人肩膀上的迭代。希望帮到你。本文还有配套的精品资源点击获取
觉得有用,分享给同行:

为您的企业打造数字门面

稳重轻奢商务风格,端正雅致视觉,长效耐看不易过时。

立即咨询 →