news 2026/9/29 16:59:25

WebSocket+Redis构建高可靠多人聊天系统:路由、去重与压测

作者头像

张小明

前端开发工程师

1.2k 24
文章封面图
WebSocket+Redis构建高可靠多人聊天系统:路由、去重与压测

简介:这是一份基于 JSP/Servlet 与 JavaScript 实现的多人聊天系统项目源码,面向正在学习 Java Web 开发、实时通信或准备课堂实训的开发者。压缩包共 11 个文件,以编译后的 class 文件与 Java 源码为主,另含 classpath、project 等工程配置,整体仅 12KB,体积精简,适合直接导入开发环境阅读核心逻辑。已有 676 人学习下载。通过该工程可了解 JSP 页面如何借助 session 等内置对象维护用户会话,Servlet 怎样接收并广播聊天消息,以及 AJAX/WebSocket 思路下的前端交互设计;源码中同时涉及 JDBC 数据库连接与聊天记录存储,便于梳理从请求处理到数据落地的完整链路。对想快速上手多人聊天场景、理解 JSP 动态页面与服务器端协作机制的读者来说,是一份紧凑可用的参考样本。

1. 多人聊天系统:从“能说话”到“人多了不崩”的坎

一套聊天源码买回来,上线两周,200 人在线还挺流畅;等注册用户冲到 2000,事情开始不对劲——消息偶尔乱序,频繁掉线,重连上来又看到重复的消息。这不是某一个人的运气问题,而是“多人聊天系统”这个标题背后的真问题:要处理的不只是消息怎么发,而是长连接管理、消息路由、可靠投递和水平扩展这几件事同时成立。我这里以 WebSocket + Redis 为主线,讲一个能落到生产环境的实现方案,适合正在评估现成源码的团队,也适合决定自己写聊天模块的开发者。后面所有命令和参数都是从真实项目里整理出来的血泪经验,照着做能少走不少弯路。

2. 选型:为什么 WebSocket + Redis 是中小团队最稳的组合

多人聊天系统的第一步不是写代码,是把传输方案定下来。这里有个很容易踩的思维惯性:把聊天当成普通接口,客户端每隔几秒 GET 一次拿新消息。200 人以内确实行,但一旦同时在线数量上去,HTTP 的重复握手和头部开销会让网关 CPU 先撑不住。

2.1 通信方式横评:轮询、SSE、WebSocket 各自的边界

短轮询是最容易想到的方案。客户端定时发 HTTP 请求,服务端把这段时间的新消息返回。实现成本最低,但它有两个硬伤:轮询间隔决定了实时性,间隔短了就把服务端变成定时炸弹;每个请求都要建立 TCP 连接,连接建立、TLS 握手、HTTP 头部的开销加在一起,会让机器在用户量不大时就先打满。

做个算术题你就明白为什么轮询撑不住:假设 2000 人在线,每人每 2 秒轮询一次,每秒就是 1000 个 HTTP 请求,每个请求即使只算 1KB 的头部,每秒也要搬运 1MB 的纯开销。而 WebSocket 保持连接后,同样 2000 人,空闲时几乎零负担,只有真正发消息时才产生流量。这个对比在向团队解释选型时非常有用。

长轮询把请求挂住,直到有新消息或超时才返回,能减少一部分空转请求,但代价是服务端要为每个挂住的请求维持一个上下文,中间还有一层 Nginx 默认 60 秒超时等着你。SSE 是单向通道,服务端可以持续推送,浏览器原生支持自动重连,做一些公告、行情推送很省事,但它天然不支持客户端往服务端发消息,聊天这种双向通信需要另搭一条上行链路,等于两套机制并行。

WebSocket 是这几种方案里唯一一个“连接建立后,双向传输都在同一条 TCP 长连接上”的。帧头最小只要 2 字节,相比 HTTP 请求头动辄几百字节,同样消息量的带宽开销能差一个量级。它带来的额外工作是:心跳、断线重连、消息背压要自己做。下表把这四种方案的关键差异列在一起,方便对比后直接拍板:

方案连接开销实时性服务端压力适用边界
短轮询每次请求重建连接取决于轮询间隔CPU/带宽压力大低频通知,间隔 ≥ 5s
长轮询连接挂在服务端秒级大量挂起连接占内存老系统改造,不推荐新做
SSE一条长连接毫秒级服务端出口带宽压力大单向推送、行情、公告
WebSocket一条长连接毫秒级需要心跳与背压管理聊天、IM、实时协作

我一般会把 WebSocket 作为主链路,SSE 只留作“系统公告”这种单向推送的副链路。真正要动手前,先把 WebSocket 的网关骨架搭起来,这是后面所有消息逻辑的承载层。

2.2 用 Node.js + ws 起一个最小网关:心跳与断线重连

Node.js 处理大量长连接的优势在于事件驱动,单个线程可以容纳每连接很低的内存开销。用 ws 库起一个最简网关,代码如下:

const WebSocket = require('ws'); const { v4: uuidv4 } = require('uuid'); const wss = new WebSocket.Server({ port: 8080, maxPayload: 64 * 1024, // 单条消息上限 64KB,防止恶意大包 clientTracking: true, }); // 每个连接的初始状态 wss.on('connection', (ws, req) => { ws.userId = parseUserIdFromReq(req); // 从 URL/Header 解析登录态 ws.sessionId = uuidv4(); ws.isAlive = true; // 30 秒一次心跳探测,收到 pong 就把状态拉回 true const heartbeatTimer = setInterval(() => { if (!ws.isAlive) { ws.terminate(); // 连续两个周期没回应,直接断开 return; } ws.isAlive = false; ws.ping(); }, 30000); ws.on('pong', () => { ws.isAlive = true; }); ws.on('message', (data) => { const buf = Buffer.isBuffer(data) ? data : Buffer.from(data); if (buf.length === 0) return; const msg = JSON.parse(buf.toString('utf8')); routeMessage(ws, msg); }); ws.on('close', () => { clearInterval(heartbeatTimer); releaseOnlineState(ws.userId, ws.sessionId); }); });

这段代码里有几个参数要专门说明。maxPayload设成 64KB,是给“聊天消息基本不会超过这个体积”留的余量,同时挡住有人把一张图片 base64 塞进 JSON 导致内存暴涨。心跳间隔 30 秒,配合 60 秒的 Nginx 默认 read timeout:两个探测周期内收不到回应就terminate(),把死连接及时清掉。

clientTracking打开后,wss.clients可以拿到当前所有连接,这在调试连接数时很有用,但生产环境不要靠它做用户态路由,因为wss.clients跨网关节点不可见。parseUserIdFromReq这里只是一个占位函数,真实环境里建议在握手阶段的 token 校验完成后,把 userId 写入连接对象,而不是在 JSON 消息里反复传身份。

断线重连和网关的关系是另一个需要注意的点。断线重连本质上发生在客户端,但服务端要配合:连接断开时,服务端要把该用户从在线状态表里移除;客户端则用指数退避重连。指数退避一般从 1 秒开始,每次乘 2,最大到 30 秒,再往上就会让用户觉得“卡死了”。

WebSocket 还有一个 HTTP 轮询没有的问题:对端消费速度比生产速度慢时,消息会在发送缓冲区堆积。ws 库提供bufferedAmount表示未发送字节数,我一般在 send 之前判断这个值,超过 1MB 就主动断开这个慢连接,让客户端重连后走增量补拉。这比硬塞内存到 OOM 强得多:

ws.on('message', (data) => { // 发送前检查缓冲区,超过 1MB 判定为慢连接 if (ws.bufferedAmount > 1 * 1024 * 1024) { ws.close(1011, 'send buffer full'); return; } // ... 正常消息处理 });

网关骨架有了,下一步才是真正让消息“找到人”的路由层。单机时可以遍历wss.clients找到目标连接,但你迟早会遇到两个以上网关节点的情况。这就是第 3 章要解决的问题。

3. 消息路由与在线状态:单机能跑,集群就抓瞎的地方

多人聊天系统的架构演进有一个很清晰的信号:当两个网关实例同时运行,用户 A 连在网关 1,用户 B 连在网关 2,A 给 B 发消息时,网关 1 要怎么知道 B 在网关 2?答案只有一个:把在线状态从进程内存里搬出去。

3.1 在线状态不能存内存:Redis Key 设计与过期策略

先解释为什么不能只靠wss.clients。单机场景下,wss.clients遍历一遍就能找到目标连接,似乎很方便。但多网关场景下,A 在 gw-01,B 在 gw-02,wss.clients只在 gw-01 内可见,gw-01 根本看不到 B。这时必须有一个人人都能读的状态中心,Redis 是常用的选择。

在线状态表的设计我一般用一个简单 Key:

online:{userId} -> gw-01|sessionId TTL 60 秒

为什么要把gatewayId放进 value?因为消息到达网关后,需要判断“目标用户是否在我这个节点上”。如果 value 只有“在线”两个字的标志,网关还是要广播各种事件去猜;放上gatewayId后,每个网关只看与自己 ID 匹配的 value,就能决定是亲手投递还是忽略。

TTL 的作用是兜底。进程被 kill -9 时,close事件不一定来得及触发,Redis 里的在线状态会残留;设 60 秒 TTL,配合每 30 秒一次的心跳续约,最多 60 秒后残留状态会被自动清掉。写入与续约的伪代码:

const key = `online:${userId}`; const ttlSeconds = 60; // 连接建立时注册 await redis.setex(key, ttlSeconds, `${gatewayId}|${sessionId}`); // 收到客户端心跳时续约 await redis.expire(key, ttlSeconds);
# 排查时直接查某个用户在哪台网关 redis-cli GET "online:100023" # 期望输出:gw-01|fd4a9c-7b3e

这里有个容易被忽略的细节:续约要复用expire而不是setex。setex会把 value 重置,如果并发的地方正在更新同一 Key,可能覆盖掉 sessionId 信息;expire只改 TTL 不动值,安全得多。另一个常见做法是把 session 也存一份,session:{userId} -> {lastSeen, gatewayId},用于登录态校验和退出登录时的精准清理,但小团队第一版可以先不拆这么细,一个online:{userId}足够。

网关节点的gatewayId必须是全局唯一的常量,不能随机生成,否则跨网关路由时无法比对。我一般直接用“gw-01、gw-02”这种配置项,写在每个节点的启动配置里,不要用 IP 端口,因为 IP 会变。

上线和下线的事件处理同样要小心。用户断开连接时,不能直接del掉在线状态——用户可能同时在手机和电脑登录,手机断线时电脑还活着。所以删除前要校验 sessionId:

async function onConnect(userId, sessionId, gatewayId) { await redis.setex(`online:${userId}`, 60, `${gatewayId}|${sessionId}`); pub.publish('chat:presence', JSON.stringify({ userId, status: 'online' })); } async function onDisconnect(userId, sessionId) { // 只清理属于本次会话的在线状态,防止新连接被覆盖 const cur = await redis.get(`online:${userId}`); if (cur && cur.endsWith(sessionId)) { await redis.del(`online:${userId}`); } pub.publish('chat:presence', JSON.stringify({ userId, status: 'offline' })); }

chat:presence是上下线事件频道,好友列表模块订阅它就能实时刷新在线状态。把“消息路由”和“上下线事件”拆成两个频道的好处是:presence 频道可以单独降级——好友列表模块暂时出问题,取消订阅即可,不影响核心消息链路。

3.2 消息转发:用 Redis Pub/Sub 把消息写给所有网关

在线状态解决了“用户在哪”,接下来要解决“消息怎么送到那个网关”。最直接的做法是:发布者把消息写进 Redis Pub/Sub 频道,所有网关节点订阅同一个频道,拿到消息后只对属于自己节点的连接做投递。

const redis = require('redis'); // 每个网关节点启动时订阅同一个频道 const sub = redis.createClient({ url: 'redis://127.0.0.1:6379/0' }); sub.subscribe('chat:route'); sub.on('message', (channel, message) => { const envelope = JSON.parse(message); // envelope: { targetUserId, roomId, payload } const conn = findLocalConnection(envelope.targetUserId); if (conn && conn.readyState === WebSocket.OPEN) { conn.send(JSON.stringify(envelope.payload)); } }); // 发送消息时,发布到频道 function publishToRouter(envelope) { const pub = redis.createClient({ url: 'redis://127.0.0.1:6379/0' }); pub.publish('chat:route', JSON.stringify(envelope)); }

这段代码的逻辑线很清楚:不管消息从哪个网关进来,最终都publish到同一个频道;每个网关都从频道里拿全量消息,再按targetUserId判断是否与自己有关。单聊时只有目标网关会命中连接;群聊时目标网关会命中多个连接,在当前节点内逐个投递。

参数上需要关注两个点。一是 Redis 的subscribe客户端不要复用做普通读写。Redis 文档里写明,订阅模式下的客户端进入“订阅态”,只能执行SUBSCRIBE/UNSUBSCRIBE等少数命令;如果你用同一个连接去get在线状态,会直接报错。我在真实项目里被这个坑过一次,后来规约就是:pub、sub、普通读写各用一个连接。

二是 Pub/Sub 不积压的脾气要到集群阶段才显山露水。某个网关节点处理不过来,或连接突然断掉,频道里已有的消息不会为它排队,直接丢。所以 Pub/Sub 适合做“大家都醒着、消息量可控”的内部总线,不适合承担可靠投递职责。聊天的可靠投递要靠第 4 章的 ACK 机制接住。

群聊场景下,Pub/Sub 广播的写法和单聊还有一点不同:群聊需要先拿到房间成员列表,再逐个publish。一个可接受的优化是把“获取成员列表”和“广播”放到 Lua 脚本里原子执行,避免消息发出时成员列表已变:

-- room:1001 -> set of userId local members = redis.call('SMEMBERS', KEYS[1]) for _, uid in ipairs(members) do redis.call('PUBLISH', KEYS[2], ARGV[1] .. '|' .. uid) end return #members

这样能保证“往频道里推送”这一个动作不会因并发导致漏消息。至于成员列表的缓存,我一般会给room:{roomId}:members设 5 秒 TTL,有人进出房间时主动刷新,避免每次发消息都打一次全量查询。

路由层到这里就通了。可是 Pub/Sub 不积压、网关可能重启、客户端可能断网,这些都会导致消息在中间环节丢失。下一章把这部分补上。

4. 消息可靠性:ACK、重发、去重三层防线

聊天消息看起来是“发出去就行”的,但真实场景里,“发出去”和“收到”之间隔着一堆不可靠环节:客户端断网、网关重启、Redis 抖动、Nginx 断连。要做到“消息不丢、不重、不乱”,三层防线是基础:客户端回执、服务端重发、幂等去重。

4.1 客户端 ACK 与服务端重发队列

先明确一个原则:服务端把消息放到 TCP 发送缓冲区不等于“送达”。TCP 只能保证字节流有序,不能保证对端应用层已经处理。真正可靠的交付标记,是客户端明确回一个ack。

服务端这边的实现思路是:发送消息时把它放进一个pending队列,启动一个定时器;收到ack就把这条消息从队列里移除;超时未确认则重发,重发超过次数上限就转入离线消息表。代码:

const pending = new Map(); // msgId -> { ws, content, retryCount, timer } const ACK_TIMEOUT = 3000; // 毫秒 const MAX_RETRY = 3; function sendWithAck(ws, content) { const msgId = uuidv4(); const item = { ws, content: { ...content, msgId }, retryCount: 0, timer: setTimeout(() => resend(msgId), ACK_TIMEOUT), }; pending.set(msgId, item); ws.send(JSON.stringify({ type: 'chat', ...item.content })); } function handleAck(ws, msgId) { const item = pending.get(msgId); if (!item || item.ws !== ws) return; clearTimeout(item.timer); pending.delete(msgId); } function resend(msgId) { const item = pending.get(msgId); if (!item) return; if (item.retryCount >= MAX_RETRY) { pending.delete(msgId); saveOfflineMessage(item.ws.userId, item.content); return; } item.retryCount += 1; item.ws.send(JSON.stringify({ type: 'chat', ...item.content, resent: true })); item.timer = setTimeout(() => resend(msgId), ACK_TIMEOUT * 2); // 退避 }

参数上,ACK_TIMEOUT设 3000 毫秒,是因为多数内网环境下聊天消息的往返可以在 50–200 毫秒内完成;3 秒是一个能容忍弱网又不至于让用户等待感爆棚的值。MAX_RETRY设 3,考虑到重试间隔会翻倍,三次重试大约覆盖 9 秒的窗口,足够应对大多数瞬时抖动。再失败的进入离线表,由用户下次上线时拉取。

客户端这边的 ACK 实现很简单,但有一个关键点:渲染成功后立刻回执,不要把ack放到 setTimeout 里拖延,否则 pending 队列会越堆越长:

ws.addEventListener('message', (event) => { const msg = JSON.parse(event.data); if (msg.type === 'chat') { renderMessage(msg); // 回执只带 msgId,让服务端尽快清掉 pending ws.send(JSON.stringify({ type: 'ack', msgId: msg.msgId })); } });

saveOfflineMessage的落库我一般用一张简单的表:

CREATE TABLE offline_message ( id BIGINT AUTO_INCREMENT PRIMARY KEY, user_id BIGINT NOT NULL, msg_id VARCHAR(64) NOT NULL, content VARCHAR(2048) NOT NULL, created_at DATETIME NOT NULL, UNIQUE KEY uk_msg_id (msg_id), KEY idx_user_created (user_id, created_at) );

msg_id上的唯一键是防重复落库的关键。同一消息如果重发两次、数据库连接抖动一次,就可能被插入两次;有了唯一键,第二次插入会直接失败,保证离线消息本身不重复。

这里有一个容易忽略的问题:客户端断线时,resend还在定时器里,消息会继续进入saveOfflineMessage,但同一消息可能已经在上次发送时被对方收到了,只是ack没回来。所以服务端的重发必须带一个“内容不变、msgId 不变”的标记,让客户端通过 msgId 去重。这自然引到下一节。

4.2 幂等:让同一条消息不被消费两次

消息重发机制让“同一条消息到达客户端两次”成为常态,而不是异常。幂等就是解决这个常态的:客户端按msgId判断,已经处理过的消息直接丢弃。

客户端维护一个滑动窗口去重表:

class DuplicateFilter: def __init__(self, window_size=200): self.seen = set() self.window = [] def is_duplicate(self, msg_id): if msg_id in self.seen: return True self.seen.add(msg_id) self.window.append(msg_id) if len(self.window) > self.window_size: old = self.window.pop(0) self.seen.discard(old) return False

窗口大小 200 是我根据“聊天消息密集场景下,200 条大约覆盖数秒流量”设的经验值。消息频率高的群聊可以把窗口调大到 500,但要警惕内存:每个 msgId 是 36 字符的 UUID,500 条的内存占用在几 KB 级,完全可接受。窗口太小的后果是:一条重发消息到达时,如果它早已经被窗口滑出,去重失败,客户端就会把它当成新消息插入会话,造成“幽灵消息”。

服务端同样要做幂等,但位置不同。离线消息补拉时,客户端用最后一条已确认消息的seq或用msgId作为游标请求增量;服务端按游标查库,天然幂等。另一个服务端幂等点是“入队”环节:客户端发消息时带一个客户端生成的clientMsgId,服务端收到后先查这个 ID 是否已经处理:

// 服务端对客户端消息做幂等 const dedupKey = `clientMsg:${clientMsgId}`; const already = await redis.set(dedupKey, '1', 'EX', 300, 'NX'); if (!already) { // 说明这条消息已经处理过,只回执不转发 ackOnly(ws, clientMsgId); return; }

SET NX的返回值是是否首次设置:返回OK表示第一次处理,返回 null 说明之前已经处理过。TTL 设 300 秒足够覆盖消息的整个生命周期。这个方案和服务端滑动窗口是互补的:一个挡“重发导致的重复”,一个挡“双击导致的重复”,合起来才把“不重”这件事兜住。

到这里,“不丢、不重”基本有了框架,但实践中还有一箩筐的坑。下一章把这些年在实时聊天上踩过的具体故障列成清单,按现象、原因、解决三步给出来。

5. 多人聊天系统避坑:5 个高频翻车点与排查清单

聊天的故障不会写成报错抛给你,它更像玄学:消息乱序、“在线”却收不到、偶尔重复、CPU 莫名打满。这一章我把高频翻车点按“现象 → 原因 → 解决”列出来,每条都是可以对照复现的。

5.1 消息乱序:一条 TCP 有序不代表全链路有序

现象:用户 A 连续发两条消息,用户 B 看到的却是第二条先出现。单机单连接时基本不会发生,一旦上了多网关、Redis Pub/Sub 路由,乱序就变常见。

原因:同一 TCP 连接内,数据一定有序;但两条消息若是不同客户端、不同网关发出,它们在不同连接里到达 Redis 的先后是没有全局保证的。Pub/Sub 本身也不承诺跨发布者的顺序。

解决:为每个会话/房间维护一个单调递增的seq。以房间维度为例,消息进入房间时用 Redis INCR 取号:

const seq = await redis.incr(`chat:seq:${roomId}`); msg.seq = seq;

客户端收到消息后不立刻渲染,而是进本地排序缓冲:

class SeqBuffer: def __init__(self, max_delta=50): self.buffer = {} self.max_seq = 0 def push(self, msg): if msg.seq <= self.max_seq: return # 重复或过期消息,丢弃 if msg.seq - self.max_seq <= 1: self.max_seq = msg.seq self._flush() # 连续,直接上抛 else: self.buffer[msg.seq] = msg if len(self.buffer) > self.max_delta: self._fill_gap() # 缺口太大,触发增量拉取

max_seq初始是最后一条已上抛消息的序号;收到序号等于max_seq + 1时直接上抛;大于则进缓冲。缓冲超过 50 条说明中间缺口太大,优先触发一次增量拉取补齐历史消息。这个方案把“顺序保证”从传输层移到了应用层,所有乱序都在进客户端 UI 之前被理平。

5.2 心跳假死:头像还在线,消息已经进不来了

现象:联系人在线列表里亮着绿灯,但发出去的消息没有回应,过一会儿才掉线。运维侧看连接数没有明显下降,但活跃消息量骤降。

原因:连接没有真正断开,而是处在“假死”状态——比如用户把笔记本合上、Wi-Fi 切网,TCP 层没有任何 FIN/RST,服务端不知道连接已死。Nginx 等中间链路也可能在空闲超时后静默切断连接,服务端收不到通知。

解决:服务端要主动探测,而不是等出错。第 2 章里的心跳代码就是一个可落地的实现:30 秒 ping 一次,连续两个周期没收到 pong 就terminate()并清理在线状态。这里有一个补充参数:心跳超时可以收紧到 15 秒。对于移动端弱网频繁的场景,15 秒能让假死连接更快被清掉;相应地,online:{userId}的 TTL 要跟节奏调成 40 秒,避免合法连接因心跳和 TTL 不同步被误踢。

事故现场可以用ss快速确认是不是假死连接太多:

# 看当前 TCP 连接里有多少半开连接 ss -tnp | grep 8080 | awk '{print $1}' | sort | uniq -c | sort -rn # 如果 ESTAB 远大于实际在线用户数,多半是假死连接没被清掉

5.3 连接被中间链路悄悄断开

现象:客户端没有主动断开,但消息出现“过一会儿就断”的规律,每次都在相同的时间点附近。

原因:Nginx 作为反代时,默认proxy_read_timeout是 60 秒。WebSocket 是一旦建立就长期挂着的连接,读操作可能在 60 秒内没有任何数据——如果没有心跳,Nginx 到上游的连接就会被清掉。这是心跳之外另一个必须改的配置:

location /ws/ { proxy_pass http://chat_upstream; proxy_http_version 1.1; proxy_set_header Upgrade $http_upgrade; proxy_set_header Connection "upgrade"; proxy_read_timeout 3600s; proxy_send_timeout 3600s; }

3600s是“1 小时无数据才断”的兜底。只调 Nginx 不放心跳,长连接依然可能因为 TCP 层占用超时被系统清理;只放心跳不调 Nginx,60 秒无数据的连接照断不误。两个都要动。验证这个故障最直接的证据是 Nginx 错误日志:

grep "upstream timed out" /var/log/nginx/error.log | tail -20

5.4 聊天记录双写不一致

现象:数据库里消息齐全,但用户端看到的历史消息缺几条,或顺序和实时消息不一致。

原因:常见实现里“写库”和“发消息”是两步,遇到数据库慢查询或网络抖动,先写库后广播或先广播后写库都可能只成功一半。另一个典型场景:离线消息补拉与实时消息同时到达,客户端没有对消息源做合并排序,实时消息就把补拉消息顶乱了。

解决:把“落库”和“进路由”串成同一条链路。消息进网关后先落库拿到自增 ID,再用这个 ID 作为seq进 Pub/Sub 广播。这样实时消息和历史记录共享同一套 ID 秩序,客户端按 ID 排序即可。补拉接口以“最后一条已收消息 ID”为游标:小于等于游标的不返回,大于游标的按 ID 升序返回。排查时看一眼消息表有没有重复的 msg_id:

SELECT msg_id, COUNT(*) c FROM message GROUP BY msg_id HAVING c > 1;

如果查询结果不为空,说明建表时漏了唯一键,或者入队逻辑没有走幂等。

5.5 广播风暴与连接数打满

现象:群里有人频繁发言,消息量不大但网关 CPU 很高,Redis 的 pub/sub 频道每秒收到大量重复推送;单机到一万连接后,新连接开始连接失败。

原因:广播逻辑把“房间内每个成员”当成一条独立消息publish,500 人群里发一条消息会产生 500 次 Redis PUBLISH,三个网关就是 1500 次;成员列表每次都查库不做缓存,放大倍数更夸张。连接数打满则是 Linux 文件描述符限制和内存叠加的结果。

解决:群聊广播改成“房间维度单次发布”,网关订阅后在本机再按房间分发。一条群聊消息只publish一次,而不是按成员publishN 次。成员列表查询做 5 秒缓存。

连接数限制先确认系统层面的值再调:

ulimit -n cat /proc/sys/fs/file-max

ulimit -n是进程级 fd 上限,file-max是系统级上限。1 万连接至少需要 1 万个 fd,把 ulimit 调到 655350 是常见做法。内存方面,一个 WebSocket 连接大约占用几十 KB,一万连接就是几百 MB,单机内存预算要按这个量级留,别按“一个连接 1KB”这种乐观值算。

6. 从单机到集群的验证:先证明它扛得住,再谈优化

前面的章节把方案搭了出来,但一个聊天系统值不值得投入,最后要拿数据说话。我习惯在写业务逻辑前先做一轮压测,重点不是刷高并发数字,而是找到“哪个环节先死”。

6.1 压测脚本与四个核心指标

用 Python 压测可以快速验证网关的吞吐边界。下面这个脚本模拟 50 个客户端,每个客户端发送 100 条消息并等待回执,记录回环延迟:

import asyncio, json, time import websockets async def client(idx, url, total): async with websockets.connect(url) as ws: for i in range(total): t0 = time.perf_counter() await ws.send(json.dumps({ "type": "chat", "roomId": "room1", "content": f"msg-{idx}-{i}", })) resp = await ws.recv() # 等待服务端回执 dt = (time.perf_counter() - t0) * 1000 print(f"cli-{idx} msg-{i} rtt={dt:.1f}ms") await asyncio.sleep(0.02) # 20ms 间隔,模拟真实输入 async def main(): url = "ws://127.0.0.1:8080" await asyncio.gather(*[client(i, url, 100) for i in range(50)]) asyncio.run(main())

压测重点看四个指标:连接数上限、消息吞吐、回环延迟、内存增长。脚本里 20ms 的sleep是模拟真实用户打字节奏,如果去掉这个 sleep,结果会严重高估网关能力,因为真实聊天不可能每个人都在无限速地发消息。压测时客户端和服务端最好分两台机器,同机压测会把网络栈的竞争也算进去,导致数字失真。

6.2 上线前的演练清单

压测通过后,我还会把下面这几项当作上线前的固定科目:杀掉一个网关进程,确认在线状态在 TTL 内被清理、连接自动迁移;断网 30 秒再恢复,确认重连后没有消息丢失;Redis 重启一次,观察发布订阅重挂是否正常。这些动作看起来简单,做一轮往往能暴露出比压测更多的问题。

我在真实项目里的一个习惯是:每次上线前挑凌晨做一次一小时的断流演练,把某个网关节点的进程直接 kill -9,然后盯着两个数字——在线状态恢复时间和消息补拉成功率。顺手也会看一眼 Redis 里有没有残留的online:{userId}僵尸 Key,有就说明close事件和 TTL 机制里至少有一个没生效。这个习惯帮我挡住过好几次线上故障的蔓延。希望帮到你。

本文还有配套的精品资源,点击获取

版权声明: 本文来自互联网用户投稿,该文观点仅代表作者本人,不代表本站立场。本站仅提供信息存储空间服务,不拥有所有权,不承担相关法律责任。如若内容造成侵权/违法违规/事实不符,请联系邮箱:809451989@qq.com进行投诉反馈,一经查实,立即删除!
网站建设 2026/9/29 16:57:20

数据要素流通下的“可用不可见“:安当DBG 如何支撑隐私计算与数据交易的字段级加密与可控脱敏

一、数据要素流通为什么需要"可用不可见" 数据作为生产要素&#xff0c;其价值只有在流通、汇聚、联合计算时才真正释放。银行要联合运营商做风控&#xff0c;医疗机构要联合药企做流行病学统计&#xff0c;政务数据要对社会力量开放用于惠民应用——这些场景的共性需…

作者头像 李华
网站建设 2026/9/29 16:56:37

从零构建AI工程能力:数据管道、训练流程与部署监控全链路指南

1. 从零搭建AI工程能力&#xff0c;为什么大多数人卡在第一步就放弃了“ai-engineering-from-scratch”这个标题&#xff0c;我第一次看到的时候心里咯噔了一下。不是因为它有多高深&#xff0c;而是因为它精准戳中了一个普遍困境&#xff1a;想学AI工程&#xff0c;但不知道从…

作者头像 李华
网站建设 2026/9/29 16:56:24

基于Django的Bilibili青少年模式使用情况数据分析系统

最近后台私信里好几个同学都在问同一个题目&#xff1a;列表里写着“基于Django的Bilibili青少年模式使用情况的数据分析系统”&#xff0c;但网页标题却顶着“Java毕设选题推荐”的标签&#xff0c;还附带“源码、mysql、文档、调试代码讲解全bao”。说实话&#xff0c;这个标…

作者头像 李华
网站建设 2026/9/29 16:55:35

C语言读文件实战指南:从fopen到fread的避坑手册

从“读不出来”说起&#xff1a;一次深夜调bug的经历先讲一件我自己的事。去年帮朋友调试一个跨平台的小工具&#xff0c;程序在某台Windows机器上怎么都读不出配置文件里的中文路径&#xff0c;返回的全是乱码。折腾到凌晨&#xff0c;最后发现是打开文件时没指定二进制模式&a…

作者头像 李华
网站建设 2026/9/29 16:55:04

srecord合并HEX文件:量产烧录的地址偏移与避坑指南

简介&#xff1a;这是面向嵌入式与微控制器开发者的 srecord-1.65.0 Windows 64 位版本&#xff0c;核心用途是把 KEIL MDK 等环境生成的多个 HEX 文件合并为单一烧录文件。合并过程中会自动核对各文件记录的地址、纠正地址顺序&#xff0c;并对冲突或重复数据做处理&#xff0…

作者头像 李华
网站建设 2026/9/29 16:54:51

Windows平台RTMP低延迟推流实战:SmartMediaKit与编码调优

做流媒体开发这几年&#xff0c;我在 Windows 平台上做 RTMP 推流验证的次数多得数不清。SmartMediaKit 是我后来用得越来越顺手的一套开源流媒体工具包&#xff0c;它本身是一套完整的媒体服务框架&#xff0c;内置了 RTMP、RTSP、HLS、GB28181 等协议支持&#xff0c;在 Wind…

作者头像 李华