
简介实时通信是现代Web应用的核心需求其技术演进经历了从HTTP轮询到长轮询再到WebSocket协议的过程。WebSocket协议基于TCP连接通过在握手阶段升级HTTP连接实现了真正的全双工、低延迟通信解决了传统轮询方式的高延迟与资源浪费问题。这一技术为在线聊天、协同编辑、实时通知等场景提供了基础支撑。在工程实践中构建高可用聊天系统需关注连接管理、心跳机制、消息协议设计等核心环节并借助Redis实现会话状态外置与消息路由结合消息队列解耦业务逻辑以支持水平扩展。本文以Node.js和ws库为例详细解析了WebSocket网关、业务服务与缓存组件的协作并针对连接稳定性、消息有序性及海量并发等常见问题提供了解决方案。1. 项目概述从HTTP轮询到WebSocket的必然选择聊到实时在线聊天这几乎是每个现代Web应用都想拥有的功能。从早期的网页QQ到现在的各种客服系统、协同办公工具背后都离不开一个核心需求消息要快要实时。我最早接触这类需求时用的还是最“朴素”的HTTP轮询客户端每隔几秒就问一次服务器“有新消息吗”服务器说“没有。”过几秒又问一遍。这种方式简单粗暴但问题一大堆延迟高、浪费带宽、服务器压力大。后来出现了长轮询Long Polling算是有点进步但本质上还是“披着羊皮的HTTP”连接管理复杂状态维护困难。直到WebSocket协议的出现才真正为Web实时通信打开了新世界的大门。它允许在单个TCP连接上进行全双工通信服务器可以主动向客户端推送数据这才是“实时”该有的样子。这次我们要聊的“基于WebSocket的实时在线聊天系统”就是利用这个协议构建一个从零到一、稳定可用的聊天架构。这个系统不仅仅是能发消息那么简单它涉及到连接管理、消息路由、状态同步、异常处理等一系列工程问题。无论你是想做一个内部团队工具还是一个面向公众的社交应用这里面的核心思路和踩坑经验都是相通的。2. 核心架构设计与技术选型2.1 为什么是WebSocket协议对比与选型逻辑在决定用WebSocket之前我们必须清楚它解决了什么问题以及它的“竞品”们有哪些短板。除了刚才提到的HTTP轮询和长轮询还有一个常被拿来比较的技术是Server-Sent EventsSSE。SSE允许服务器向客户端单向推送数据对于只需要服务器推送的场景比如新闻推送、股价更新很合适但它不支持客户端向服务器的双向通信。聊天场景下客户端既要收消息也要发消息SSE就不够用了。WebSocket协议RFC 6455在握手阶段借用了HTTP/1.1的Upgrade机制一旦握手成功连接就升级为全双工的WebSocket连接后续的数据帧传输就与HTTP无关了。这带来了几个核心优势低延迟消息无需等待客户端请求服务器可以立即推送。低开销每个消息帧的头部很小最小仅2字节远小于HTTP头部。双向通信客户端和服务器可以随时互发消息适合对话场景。技术选型上对于聊天系统WebSocket几乎是唯一正确的选择。但具体到实现我们还需要选择服务端的技术栈。Node.js的ws库、Socket.IO 框架Java的Netty、Spring WebSocketGo的gorilla/websocket等都是热门选择。我的经验是如果你的团队熟悉JavaScript且追求快速原型开发Socket.IO是个不错的选择它自动处理了降级在不支持WebSocket的环境下回退到HTTP轮询和重连。但如果你追求极致的性能和可控性像ws或 Netty 这样的底层库会更合适它们更轻量但需要你自己处理更多细节比如心跳、重连、消息编解码。2.2 系统整体架构图与组件职责一个健壮的聊天系统不能只有一个WebSocket服务。一个典型的架构会包含以下组件客户端Web前端、移动端App等。负责建立并维持WebSocket连接发送和渲染消息。WebSocket网关/服务器这是系统的核心负责维护所有活跃的WebSocket连接。它处理连接建立、关闭、消息的路由和广播。为了支持水平扩展这个服务通常是无状态的或者将会话状态外置到Redis等存储中。业务逻辑服务处理聊天相关的业务逻辑如消息的持久化存储存入数据库、敏感词过滤、消息推送逻辑如某人的判断等。WebSocket网关接收到消息后通常会通过消息队列如RabbitMQ、Kafka或RPC调用将消息转发给业务逻辑服务处理。状态与缓存服务主要是Redis。用于存储在线用户列表、用户的连接信息哪个用户连接到了哪个WebSocket网关实例、未读消息计数、临时消息缓存等。这是实现多网关实例协同工作的关键。数据库如MySQL、PostgreSQL或MongoDB用于永久存储聊天记录、用户信息、群组信息等。消息队列用于解耦WebSocket网关和业务逻辑服务确保消息的可靠异步处理尤其在流量高峰时起到削峰填谷的作用。整个数据流大致是这样的用户A发送一条消息 - 前端通过WebSocket连接发送到网关 - 网关解析消息可能进行初步验证 - 网关将消息发布到消息队列的某个主题Topic - 业务逻辑服务消费该消息进行业务处理并持久化 - 业务逻辑服务根据处理结果如确定要发送给用户B和C向消息队列发布一条“推送任务” - WebSocket网关订阅“推送任务”根据任务中指定的用户ID查询Redis找到对应用户当前连接在哪个网关实例上然后将消息通过对应的WebSocket连接推送给用户B和C的前端。3. 核心细节解析与实操要点3.1 连接建立、维持与优雅关闭连接的生命周期管理是WebSocket编程中最基础也最容易出问题的一环。连接建立握手前端使用new WebSocket(ws://your-domain.com/chat)发起连接。服务端需要正确处理HTTP Upgrade请求。这里有个关键点务必验证Origin头如果是在浏览器环境中以防止跨站WebSocket劫持攻击。虽然WebSocket协议本身不受同源策略限制但服务端应该检查Origin是否在白名单内。连接维持心跳网络环境复杂中间可能有防火墙或代理会关闭长时间空闲的连接。因此必须实现心跳机制Heartbeat/Ping-Pong。WebSocket协议本身定义了Ping和Pong帧可以用于此目的。服务端应定时如每30秒向客户端发送一个Ping帧客户端收到后自动回复Pong帧。如果连续多次未收到Pong回复则可以认为连接已失效主动关闭它并清理相关资源。// 服务端Node.js ws库心跳示例 const WebSocket require(ws); const wss new WebSocket.Server({ port: 8080 }); wss.on(connection, function connection(ws) { ws.isAlive true; ws.on(pong, () { ws.isAlive true; }); // 设置一个定时器每隔30秒检查一次 const interval setInterval(() { if (ws.isAlive false) { clearInterval(interval); return ws.terminate(); // 终止连接 } ws.isAlive false; ws.ping(); // 发送Ping帧 }, 30000); ws.on(close, () { clearInterval(interval); // 清理定时器 }); });优雅关闭连接关闭时状态码Code和原因Reason很重要。正常关闭应使用状态码1000CLOSE_NORMAL或1001CLOSE_GOING_AWAY。异常关闭如服务端内部错误可以使用1011INTERNAL_ERROR。前端需要监听onclose事件并根据状态码决定是否重连。特别注意状态码1006这是一个特殊的代码表示连接异常关闭但具体原因未知通常发生在网络突然中断、浏览器标签页关闭等场景。处理1006错误时前端应尝试自动重连。3.2 消息协议设计从JSON到二进制WebSocket传输的是帧Frame帧里是二进制数据。我们需要定义一套应用层协议让客户端和服务端能理解彼此发送的“消息”是什么。最常用、最方便的是JSON。我们可以定义一个简单的消息格式{ type: chat_message, // 消息类型chat_message, system_notice, heart_beat等 sender: user123, recipient: room456, // 或 user456 content: { text: 你好, timestamp: 1627891234567 }, seq: 42 // 可选消息序列号用于去重或排序 }JSON的好处是可读性好易于调试前端直接JSON.parse()就能用。缺点是体积相对较大尤其是传输大量小消息或需要传输二进制数据如图片、文件时效率不高。对于性能要求极高的场景可以考虑使用二进制协议如Protobuf、MessagePack或自定义的二进制格式。它们能显著减少消息体积加快序列化/反序列化速度。但代价是开发调试复杂度增加需要预先定义严格的.proto文件或格式规范。我的建议是对于绝大多数聊天应用初期使用JSON完全足够。当消息量非常大成为性能瓶颈时再考虑优化协议。你可以先为JSON消息添加一个简单的“压缩”标志服务端和客户端约定好如果消息体超过一定大小如1KB就先用gzip或deflate压缩一下再传输这是一个性价比很高的优化。3.3 用户状态、会话管理与多节点扩展单机WebSocket服务只能支撑有限的连接数通常受限于端口和内存。要支持海量用户在线必须让WebSocket服务能水平扩展。核心思路是让WebSocket网关无状态化将状态外置。会话存储当用户连接成功时生成一个唯一的会话IDSession ID。将这个会话ID、用户ID、当前连接的网关实例IPInstance ID的映射关系存储到Redis中并设置一个过期时间如连接超时时间的两倍。消息路由当需要给某个用户发消息时业务逻辑服务不直接发而是向消息队列发布一个“推送指令”指令里包含目标用户ID和消息内容。所有的WebSocket网关实例都订阅这个消息队列。网关寻址每个网关实例消费到“推送指令”后去Redis里查一下目标用户ID当前连接在哪个网关实例上。如果查到的Instance ID正是自己那么就直接通过本地维护的连接对象把消息发出去。如果不是自己则忽略这条指令或者有一种更高效的设计是用Redis的Pub/Sub让持有连接的网关订阅以用户ID为名的频道业务服务直接向该频道发布消息。这就带来了另一个问题如何让用户始终连接到“正确”的网关这通常需要一个负载均衡器。可以使用Nginx或HAProxy做TCP层的负载均衡因为WebSocket建立在TCP之上。更现代的做法是使用云服务商提供的负载均衡器或者使用应用层网关如Kong、Apisix它们对WebSocket有更好的支持并能做更复杂的路由判断比如基于用户ID的粘性会话。注意使用Nginx做WebSocket代理时必须配置几个关键参数来支持长连接proxy_http_version 1.1; proxy_set_header Upgrade $http_upgrade; proxy_set_header Connection upgrade; proxy_read_timeout 3600s; # 设置一个较长的超时时间否则连接可能会被意外断开。4. 实操过程与核心环节实现4.1 搭建基础WebSocket服务以Node.js ws为例我们先从最简单的单机服务开始。使用Node.js和ws库可以快速搭建一个原型。首先初始化项目并安装依赖mkdir websocket-chat-server cd websocket-chat-server npm init -y npm install ws redis ioredis express // 我们加上redis和express以备后用创建一个server.js文件const WebSocket require(ws); const http require(http); const url require(url); // 创建HTTP服务器WebSocket服务器将附加其上 const server http.createServer(); const wss new WebSocket.Server({ noServer: true }); // 不立即监听端口 // 用于存储连接用户映射单机内存存储仅演示生产环境用Redis const userConnections new Map(); wss.on(connection, function connection(ws, request) { // 从请求URL中解析用户ID实际应从认证token中获取 const queryParams url.parse(request.url, true).query; const userId queryParams.userId; if (!userId) { ws.close(1008, Missing user identification); // 1008: Policy Violation return; } console.log(用户 ${userId} 已连接); // 存储连接 userConnections.set(userId, ws); // 通知用户连接成功 ws.send(JSON.stringify({ type: system, content: 连接成功 })); // 监听消息 ws.on(message, function incoming(message) { console.log(收到来自 %s 的消息: %s, userId, message); try { const data JSON.parse(message); // 处理不同类型的消息 handleMessage(userId, data, ws); } catch (e) { ws.send(JSON.stringify({ type: error, content: 消息格式错误 })); } }); // 监听关闭 ws.on(close, function close() { console.log(用户 ${userId} 断开连接); userConnections.delete(userId); // 这里可以广播用户下线通知 broadcastSystemMessage(${userId} 离开了聊天室); }); // 错误处理 ws.on(error, console.error); }); // 简单的消息处理器 function handleMessage(senderId, data, ws) { switch (data.type) { case chat: // 假设data.recipient是接收者IDdata.content是内容 const recipientWs userConnections.get(data.recipient); if (recipientWs recipientWs.readyState WebSocket.OPEN) { recipientWs.send(JSON.stringify({ type: chat, sender: senderId, content: data.content, timestamp: Date.now() })); // 可选发送回执给发送者 ws.send(JSON.stringify({ type: ack, msgId: data.msgId })); } else { // 接收者不在线可以存入离线消息库 ws.send(JSON.stringify({ type: error, content: 用户不在线 })); } break; case heartbeat: ws.send(JSON.stringify({ type: heartbeat, echo: data.echo })); break; default: ws.send(JSON.stringify({ type: error, content: 未知的消息类型 })); } } // 广播系统消息 function broadcastSystemMessage(content) { const message JSON.stringify({ type: system, content }); for (const [userId, ws] of userConnections) { if (ws.readyState WebSocket.OPEN) { ws.send(message); } } } // 处理HTTP服务器升级请求 server.on(upgrade, function upgrade(request, socket, head) { // 这里可以添加身份验证逻辑比如验证JWT Token const pathname url.parse(request.url).pathname; if (pathname /chat) { wss.handleUpgrade(request, socket, head, function done(ws) { wss.emit(connection, ws, request); }); } else { socket.destroy(); // 拒绝非/chat路径的升级请求 } }); server.listen(8080, function() { console.log(WebSocket server is listening on port 8080); });这个简单的服务器实现了用户连接、点对点聊天、心跳响应和系统广播。但它把所有状态都存在内存里一旦服务重启所有连接和状态都会丢失。接下来我们要引入Redis来解决这个问题。4.2 集成Redis管理在线状态与消息中转我们使用ioredis库来连接Redis。主要做两件事1. 用户上线/下线时在Redis中记录/清除其连接信息2. 通过Redis的Pub/Sub功能实现跨进程/跨服务器的消息广播。首先修改连接处理逻辑将用户连接信息存入Redisconst Redis require(ioredis); const redis new Redis(); // 默认连接本地6379端口 // 在connection事件中 wss.on(connection, async function connection(ws, request) { const userId getUserIdFromRequest(request); // 假设这是一个提取用户ID的函数 // 存储用户-网关映射。key: ws:user:${userId}, value: gateway:${instanceId} const instanceId process.env.INSTANCE_ID || gateway_1; // 网关实例标识 await redis.set(ws:user:${userId}, instanceId, EX, 7200); // 过期时间2小时 // 将用户加入在线集合方便统计和广播 await redis.sadd(ws:online_users, userId); // ... 其余连接逻辑 ws.on(close, async function close() { // 用户断开时删除映射和在线状态 await redis.del(ws:user:${userId}); await redis.srem(ws:online_users, userId); }); });然后实现一个基于Redis Pub/Sub的广播通道。每个网关实例启动时都订阅一个公共频道如ws:broadcast和一个自己实例的专属频道如ws:gateway:${instanceId}。const subRedis new Redis(); // 专门用于订阅的Redis连接 // 订阅公共广播频道 subRedis.subscribe(ws:broadcast, (err, count) { if (err) console.error(订阅失败, err); else console.log(已订阅广播频道当前频道数: ${count}); }); // 订阅本实例专属频道用于接收定向推送 subRedis.subscribe(ws:gateway:${instanceId}); // 监听订阅频道的消息 subRedis.on(message, async (channel, message) { try { const msg JSON.parse(message); if (channel ws:broadcast) { // 广播消息发给所有本地连接的用户 broadcastToAllLocal(msg); } else if (channel ws:gateway:${instanceId}) { // 定向推送消息msg.targetUserId 指定了目标用户 const targetUserId msg.targetUserId; const targetWs getLocalConnection(targetUserId); // 从本地Map查找连接 if (targetWs) { targetWs.send(JSON.stringify(msg.payload)); } } } catch (e) { console.error(处理Redis订阅消息出错, e); } }); // 当需要给特定用户发消息时例如从业务逻辑服务发出 async function sendMessageToUser(targetUserId, messagePayload) { // 1. 查一下用户在哪个网关 const userGateway await redis.get(ws:user:${targetUserId}); if (!userGateway) { // 用户不在线存入离线消息 await storeOfflineMessage(targetUserId, messagePayload); return; } // 2. 向该网关的专属频道发布消息 const pubMessage JSON.stringify({ targetUserId: targetUserId, payload: messagePayload }); await redis.publish(ws:gateway:${userGateway}, pubMessage); }这样我们就实现了一个支持水平扩展的、状态外置的WebSocket消息路由骨架。业务逻辑服务只需要调用sendMessageToUser函数而无需关心用户具体连接在哪台机器上。4.3 前端连接管理与消息收发前端部分我们使用原生WebSocket API并封装一些重连和状态管理逻辑。class WebSocketClient { constructor(url, options {}) { this.url url; this.options options; this.ws null; this.reconnectAttempts 0; this.maxReconnectAttempts options.maxReconnectAttempts || 5; this.reconnectDelay options.reconnectDelay || 3000; this.messageHandlers new Map(); // 按消息类型存储处理器 this.isConnected false; this.connect(); } connect() { try { // 构建带认证参数的URL实际中更常用在header中传token这里演示用query const fullUrl ${this.url}?userId${this.options.userId}token${this.options.token}; this.ws new WebSocket(fullUrl); this.ws.onopen () { console.log(WebSocket连接已建立); this.isConnected true; this.reconnectAttempts 0; this.options.onConnected this.options.onConnected(); // 开始心跳 this.startHeartbeat(); }; this.ws.onmessage (event) { try { const message JSON.parse(event.data); const handler this.messageHandlers.get(message.type); if (handler) { handler(message); } else { console.warn(未注册处理器的消息类型: ${message.type}, message); } } catch (e) { console.error(解析消息失败:, e, event.data); } }; this.ws.onclose (event) { console.log(连接关闭代码: ${event.code}, 原因: ${event.reason}); this.isConnected false; this.stopHeartbeat(); this.options.onDisconnected this.options.onDisconnected(event); // 非正常关闭且未超过重试次数则尝试重连 if (event.code ! 1000 this.reconnectAttempts this.maxReconnectAttempts) { this.scheduleReconnect(); } }; this.ws.onerror (error) { console.error(WebSocket错误:, error); this.options.onError this.options.onError(error); }; } catch (error) { console.error(创建WebSocket连接失败:, error); } } scheduleReconnect() { this.reconnectAttempts; const delay this.reconnectDelay * Math.pow(1.5, this.reconnectAttempts - 1); // 指数退避 console.log(将在 ${delay}ms 后尝试第 ${this.reconnectAttempts} 次重连...); setTimeout(() this.connect(), delay); } startHeartbeat() { this.heartbeatInterval setInterval(() { if (this.isConnected this.ws.readyState WebSocket.OPEN) { const heartbeatMsg { type: heartbeat, echo: Date.now() }; this.ws.send(JSON.stringify(heartbeatMsg)); } }, 25000); // 25秒发送一次心跳 } stopHeartbeat() { if (this.heartbeatInterval) { clearInterval(this.heartbeatInterval); this.heartbeatInterval null; } } sendMessage(type, payload) { if (this.isConnected this.ws.readyState WebSocket.OPEN) { const message { type, ...payload, timestamp: Date.now() }; this.ws.send(JSON.stringify(message)); return true; } else { console.error(发送失败连接未就绪); // 可以将消息加入发送队列等待重连后发送 return false; } } onMessage(type, handler) { this.messageHandlers.set(type, handler); } close() { this.stopHeartbeat(); if (this.ws) { this.ws.close(1000, 用户主动关闭); // 正常关闭 } } } // 使用示例 const client new WebSocketClient(ws://localhost:8080/chat, { userId: alice123, token: your_auth_token_here, onConnected: () { console.log(已连接); }, onDisconnected: (event) { console.log(连接断开, event); } }); // 注册消息处理器 client.onMessage(chat, (msg) { console.log(收到来自 ${msg.sender} 的消息: ${msg.content.text}); // 更新UI显示消息 }); client.onMessage(system, (msg) { console.log(系统通知: ${msg.content}); }); // 发送消息 document.getElementById(send-btn).addEventListener(click, () { const text document.getElementById(message-input).value; client.sendMessage(chat, { recipient: bob456, content: { text } }); });这个前端封装类处理了连接建立、自动重连使用指数退避策略避免重连风暴、心跳维持和消息分发是一个比较健壮的实践。5. 常见问题与排查技巧实录5.1 连接不稳定与1006错误排查问题现象客户端频繁断开连接控制台看到onclose事件中event.code为1006event.reason为空。排查思路检查网络环境1006通常意味着底层TCP连接异常断开。可能是用户网络不稳定、移动网络切换、或者浏览器标签页进入后台被冻结。这是客户端环境问题服务端能做的不多。检查代理与中间件如果使用了Nginx、HAProxy或云负载均衡器确保其配置正确支持WebSocket长连接如前文提到的proxy_read_timeout,Upgrade,Connection头部。一个常见的坑是负载均衡器的空闲超时时间设置过短如60秒而你的心跳间隔是70秒那么连接就会被代理主动掐断。确保代理的超时时间远大于心跳间隔。检查服务端资源服务端是否达到了文件描述符上限内存是否耗尽使用netstat或ss命令查看连接数。使用ulimit -n检查并调整进程可打开的文件数限制。检查防火墙与安全组云服务器安全组是否放行了WebSocket使用的端口如8080、443/wss本地防火墙是否阻止了连接客户端增加诊断在前端的onclose和onerror事件中详细记录错误信息、时间戳和网络状态如果有navigator.onLine并尝试上报到服务端帮助定位问题。应对策略前端实现健壮的重连机制并给用户友好的提示如“网络连接已断开正在尝试重连...”。服务端优化心跳间隔比如从30秒调整为25秒确保在代理超时前有数据交换。考虑使用WebSocket over TLSWSS特别是在使用某些代理和移动网络时WSS连接更不容易被中间设备干扰。5.2 多节点下的消息重复与乱序问题场景在分布式环境下用户A发送一条消息由于网络延迟或业务逻辑服务多实例消费消息队列可能导致用户B收到两条一模一样的消息或者后发的消息先被收到。解决方案消息去重为每条消息生成一个全局唯一的ID如UUID或雪花算法生成的ID。在业务逻辑服务处理消息并准备持久化或转发时先检查Redis中是否存在该消息ID使用SET key message_id NX EX 60命令设置60秒过期。如果已存在说明是重复消息直接丢弃。这个ID可以由客户端生成并在发送时带上。消息有序性对于单聊可以在消息体中增加一个严格递增的序列号seq由发送方客户端维护每个对话一个独立的seq。接收方前端根据seq进行排序和展示。服务端不需要保证全局有序只需保证转发不丢失。对于群聊/广播保证绝对有序非常困难且成本高。一个折中方案是在消息体中加入服务器端的毫秒级时间戳前端收到后按时间戳排序显示。由于时钟可能不同步可以允许小范围如500ms内的消息由前端根据一些规则如发送者、消息类型进行微调。对于强有序要求的场景如协同编辑可能需要引入更复杂的版本向量或CRDT数据结构。5.3 海量连接下的性能优化当在线用户数达到万级甚至十万级时单纯的“一个连接一个线程/进程”模型会崩溃。优化方向使用高性能网络库在Node.js中ws库本身基于uWebSockets的C绑定性能已经很好。在Java领域Netty是公认的高性能异步事件驱动框架专门为高并发网络应用设计。Go语言的gorilla/websocket或nhooyr.io/websocket也非常高效。连接优化减少内存占用每个连接都是一个对象。确保不在连接对象上挂载过大的数据如完整的用户信息。只存储必要的连接标识如用户ID、会话ID其他数据从Redis等外部存储按需读取。使用连接池对于需要访问数据库或其他外部服务的操作一定要使用连接池避免为每个请求创建新连接。流量控制与背压客户端可能以极快的速度发送消息如拖拽事件流。服务端需要对每个连接设置接收缓冲区上限并在达到上限时暂停读取在Node.js中可监听ws._socket的drain事件或者直接断开恶意连接。服务端向客户端推送时也要注意。如果客户端网络慢服务端一直发会导致内核缓冲区积压最终耗尽内存。需要监听send方法的回调或drain事件实现简单的背压控制。水平扩展如前所述使用无状态网关 Redis 消息队列的架构。通过负载均衡将连接分散到多个网关实例上。确保你的Redis和消息队列集群也能承受相应的压力。5.4 安全与认证考量WebSocket连接一旦建立就是一个长久的通道认证和授权至关重要。连接认证不要在URL参数中传递明文Token。推荐的做法是在建立WebSocket连接前先通过一个普通的HTTP接口完成登录获取一个短期有效的、针对WebSocket连接的Ticket或Token。然后在建立WebSocket连接时将这个Token放在标准的Authorization头中虽然WebSocket握手是HTTP但自定义头部在某些代理环境下可能被过滤更稳妥的做法是放在一个约定的协议头中或者作为第一个握手后的数据帧发送给服务端进行验证。消息验证即使连接认证通过对每条消息也要进行业务层面的权限验证。例如用户A是否真的有权限向群组B发送消息这条消息是否包含恶意脚本服务端在处理每条消息时都需要重新从Redis或数据库加载用户的权限上下文进行校验。输入清洗与防注入对接收到的消息内容进行严格的清洗和转义防止XSS攻击。如果消息内容要存入数据库还要防止SQL注入。对于JSON消息使用安全的JSON解析库避免解析畸形JSON导致的服务崩溃。限制与监控对每个连接的消息发送频率进行限制限流防止恶意用户刷屏或发起DoS攻击。监控异常连接如每秒发送消息数异常高、连接存活时间极短等并自动加入黑名单。6. 进阶话题从单聊到群聊与聊天室我们之前主要讨论了点对点单聊。扩展到群聊或聊天室核心变化在于消息的路由逻辑。群聊消息路由用户A在群G中发送一条消息。WebSocket网关收到后将消息发布到消息队列主题为chat:group:${groupId}。业务逻辑服务消费该消息进行持久化并查询群G的在线成员列表可以从Redis的集合group:${groupId}:members和全局在线用户集合ws:online_users取交集得到。对于每一个在线成员业务服务调用sendMessageToUser函数即我们之前实现的通过Redis Pub/Sub路由到具体网关。如果群成员非常多如超大群全量遍历在线成员并逐个发送效率低下。此时可以优化让网关也订阅群主题。业务服务只需向chat:group:${groupId}发布一条消息所有网关实例都收到然后每个网关在自己本地维护的连接中找出属于该群的用户进行推送。这要求网关在用户连接时不仅记录用户ID还要记录用户加入了哪些群并维护一个群ID到本地用户连接列表的倒排索引。聊天室广播室 聊天室通常不需要持久化每一条消息且消息是广播给房间内所有人的。实现更简单用户加入房间时服务端将其连接加入一个“房间连接集合”在单机内存中或使用Redis的Set。有用户发言时服务端遍历这个集合向其中每一个连接发送消息。对于分布式场景可以使用Redis的Pub/Sub。每个房间对应一个频道如room:${roomId}。用户连接的网关订阅其所在房间的频道。当有消息需要广播时发布到该频道所有订阅了该频道的网关实例都会收到并转发给本地属于该房间的连接。离线消息处理 对于单聊和群聊当目标用户不在线时消息需要被存储起来待其上线后推送。可以设计一张offline_messages表字段包括id,recipient_id,sender_id,type单聊/群聊,related_id对话ID/群ID,content,created_at。当用户上线时业务服务查询其所有未读的离线消息按会话合并或逐条推送给用户并在推送成功后标记为已读或删除。为了减轻数据库压力对于非常活跃的用户也可以将最近的离线消息缓存在Redis中。本文还有配套的精品资源点击获取