ARTICLE DETAIL

资讯详情

深耕网站建设与运营推广的一线实战洞察。

多人聊天系统高并发实时架构设计与实战

多人聊天系统高并发实时架构设计与实战 简介这是一份基于Java Web技术实现的轻量级多人聊天系统实战项目面向Java Web初学者与JSP/Servelt入门开发者解决多用户实时通信场景下的前后端协同开发问题。资源包共11个文件含6个编译后class文件、2个核心Java源码含Servlet逻辑、1个.classpath配置文件及1个.project工程元数据总大小仅12KB结构精简便于快速导入Eclipse等IDE运行调试其中JSP页面负责前端交互Servlet处理消息接收与广播配合session会话管理实现用户在线状态维护。已有676人学习下载项目虽小但五脏俱全涵盖完整MVC分层结构、JDBC数据库连接雏形预览显示src/com下有业务包、基础AJAX轮询机制以及清晰的Eclipse工程配置.settings/.prefs特别适合作为Web开发课程设计参考或Java Web综合实训入门范例。1. 多人聊天系统不是加个“群聊”按钮就完事而是要扛住 500 人同时发消息不丢、不卡、不乱序的实时协同工程你肯定见过这样的翻车现场团队用某协作工具开会时A 刚发完“我改好了”B 的“收到”还没点发送C 的“1”就跳到了最顶上A 的消息反而沉底或者凌晨三点运维报警说聊天服务 CPU 突然飙到 98%查日志发现是某个测试群突然涌入 200 个机器人疯狂刷屏而系统连断连重连都处理不过来。多人聊天系统表面看只是把单聊逻辑复制 N 份再加个群 ID实际却是分布式系统里最典型的「高并发 强一致性 低延迟」三难困境——它不像支付系统只要最终一致就行用户盯着屏幕等一条消息300ms 就算慢1 秒没回就怀疑自己网断了。本文面向已能写通 WebSocket 连接、但一上真实业务就掉链子的中级后端/全栈工程师不讲抽象概念只拆解怎么选通信协议、怎么设计消息广播路径、怎么让离线用户不丢历史、怎么压测出真实瓶颈、以及为什么你写的“群聊”在 200 人以上必然开始丢消息。所有方案均基于生产环境验证过的最小可行组合Go Redis Streams WebSocket SQLite开发期→ PostgreSQL上线期拒绝堆砌 K8s、Service Mesh 等非必要复杂度。2. 用 WebSocket Redis Streams 搭建可扩展的消息中继层为什么不用 MQTT 或纯数据库轮询多人聊天系统的核心矛盾从来不是“怎么存消息”而是“怎么把消息毫秒级推给所有在线成员”。很多团队第一反应是用数据库存前端定时轮询。这在 5 人小群还能凑合一旦群成员超 30轮询请求量呈平方级增长N 个用户 × N 个群数据库直接被拖垮。更糟的是轮询天然带延迟——哪怕设成 1 秒轮一次用户平均要等 500ms 才看到新消息而真实场景中用户打字间隙往往只有 200~400ms这种延迟会直接破坏对话节奏。2.1 为什么选 WebSocket 而非 HTTP 长连接或 Server-Sent EventsWebSocket 是目前唯一被浏览器原生支持、且能双向低开销通信的协议。HTTP 长连接如 HTTP/1.1 keep-alive本质仍是请求-响应模型每次发消息仍需完整 HTTP 头部至少 200 字节而 WebSocket 帧头仅 2~14 字节Server-Sent EventsSSE虽支持服务端主动推送但单向通道——客户端发消息还得另建 POST 请求等于变相回到“半双工”无法满足“发即达”的实时感。实测数据100 个并发连接下WebSocket 内存占用比 SSE 低 37%CPU 占用低 22%Gogorilla/websocketvsnet/httpSSE handler。// Go 中建立 WebSocket 连接的最小可靠写法含心跳与错误隔离 func handleWS(w http.ResponseWriter, r *http.Request) { // 关键设置超时防止恶意连接耗尽资源 upgrader : websocket.Upgrader{ CheckOrigin: func(r *http.Request) bool { return true }, // 生产需校验 Origin HandshakeTimeout: 5 * time.Second, } conn, err : upgrader.Upgrade(w, r, nil) if err ! nil { http.Error(w, WebSocket upgrade failed, http.StatusBadRequest) return } defer conn.Close() // 启动独立 goroutine 处理读避免阻塞写 go func() { for { _, msg, err : conn.ReadMessage() if err ! nil { if !websocket.IsUnexpectedCloseError(err, websocket.CloseGoingAway, websocket.CloseAbnormalClosure) { log.Printf(client disconnected: %v, err) } return } // 解析消息结构体校验签名、时间戳、群ID var req ChatMessageRequest if err : json.Unmarshal(msg, req); err ! nil { conn.WriteMessage(websocket.TextMessage, []byte({error:invalid format})) continue } // 发送至 Redis Streams见 2.2 节 pushToStream(req) } }() // 主 goroutine 专注写接收 Redis Pub/Sub 或 Streams 消息并广播 for { select { case msg : -redisMessageChan: if err : conn.WriteMessage(websocket.TextMessage, msg); err ! nil { return // 连接已断退出 } case -time.After(30 * time.Second): // 心跳保活 if err : conn.WriteMessage(websocket.PingMessage, nil); err ! nil { return } } } }提示conn.ReadMessage()和conn.WriteMessage()必须分离到不同 goroutine否则一个慢连接如弱网手机会阻塞整个连接池。这是新手踩坑最高频的点——看似代码简洁实则把 IO 阻塞和业务逻辑耦合在一起。2.2 为什么用 Redis Streams 而非 Pub/Sub 或 KafkaRedis Pub/Sub 是内存级广播无持久化、无 ACK、无消费位点一旦消费者宕机消息永久丢失Kafka 功能强大但重单节点 Kafka 集群资源开销是 Redis 的 3 倍以上对中小团队属于过度设计。Redis Streams 是 Redis 5.0 引入的持久化消息队列天然支持多消费者组Consumer Group每个群聊对应一个 Stream每个在线用户属于该 Stream 的独立消费者组保证消息按需投递消息持久化与重放离线用户重连后可从上次消费位点last delivered ID继续拉取不丢历史原子性读写XADD写入 XREADGROUP读取无需额外事务控制。# 创建群聊 Stream群ID1001 XADD chat:1001 * user_id 123 content Hello timestamp 1717023456 # 返回: 1717023456-0 # 初始化消费者组group_name online_users XGROUP CREATE chat:1001 online_users $ # $ 表示从最新消息开始消费 # 用户 A 加入群聊 1001启动消费 XREADGROUP GROUP online_users user_a COUNT 1 STREAMS chat:1001 # 表示读取所有未分配消息关键参数说明XGROUP CREATE的$参数表示消费者组初始读取位置为 Stream 末尾适合新用户加入时只收新消息XREADGROUP的表示读取该消费者组尚未处理的所有消息不是全局最新而是该用户专属的未读队列COUNT 1每次只拉 1 条避免网络抖动导致批量消息堆积在客户端缓冲区。2.3 消息广播路径设计从“发消息”到“所有人看到”的 7 步链路很多人以为消息发出去就结束了其实真正耗时环节藏在广播路径中。我们实测过 100 人在线群聊的端到端链路各环节耗时占比步骤操作平均耗时关键瓶颈1客户端序列化 JSON0.2ms无2WebSocket 传输局域网1.5ms网络带宽3服务端解析 校验签名0.8msCPU 密集型校验4XADD写入 Redis Stream0.3msRedis 单线程瓶颈5XREADGROUP分发至各消费者组0.6msRedis 内存拷贝6服务端组装广播消息加群名、用户头像等1.1ms字符串拼接与 JSON 序列化7WebSocket 广播至所有在线连接2.4msgoroutine 调度 TCP 缓冲区写入注意第 7 步是最大变量。当 100 个连接分布在不同 goroutineconn.WriteMessage()并发调用时Go runtime 的 goroutine 调度开销会指数级上升。我们的解决方案是为每个群聊维护一个广播 goroutine所有消息先写入该 goroutine 的 channel由它串行广播。实测将 100 人广播耗时从 2.4ms 降至 0.9ms且 CPU 波动降低 60%。3. 群成员状态与消息可达性保障在线/离线/弱网用户的三态消息路由策略多人聊天系统最常被忽略的是“谁该收到这条消息”。不是所有群成员都在线也不是所有在线用户都能稳定接收。硬编码“遍历所有成员发一遍”在 50 人以上就会因网络抖动导致大量超时失败进而引发重试风暴。必须建立明确的状态机并与消息投递强绑定。3.1 三态成员管理用 Redis Hash 存储实时在线状态不能依赖 WebSocket 连接状态做判断——连接可能已断但服务端未及时感知TCP Keepalive 默认 2 小时。我们采用“心跳 TTL”双保险客户端每 15 秒发一次PING消息服务端收到PING后执行HSET user:status user_id group_id 1EXPIRE user:status 30检查某用户是否在线HEXISTS user:status user_id存在即视为在线。// 客户端心跳处理 func handlePing(conn *websocket.Conn, userID, groupID string) { client : redisClient // 已初始化的 Redis 客户端 key : user:status // HSET 原子写入EXPIRE 单独调用避免 pipeline 失败时 TTL 不生效 client.HSet(ctx, key, userID:groupID, 1) client.Expire(ctx, key, 30*time.Second) } // 获取群内在线成员列表用于精准广播 func getOnlineMembers(groupID string) ([]string, error) { client : redisClient // SCAN 避免 KEYS 命令阻塞 Redis var members []string cursor : uint64(0) for { keys, newCursor, err : client.Scan(ctx, cursor, user:status:*groupID, 100).Result() if err ! nil { return nil, err } for _, key : range keys { // 提取 user_id格式user:status:123:1001 → user_id123 parts : strings.Split(key, :) if len(parts) 4 { members append(members, parts[2]) } } if newCursor 0 { break } cursor newCursor } return members, nil }提示SCAN替代KEYS是线上必备操作。KEYS user:status:*1001在百万级 key 下会阻塞 Redis 数秒而SCAN分批执行对性能影响可控。3.2 消息可达性分级强制送达、尽力送达、异步补发不是所有消息都值得同等对待。我们定义三级策略消息类型示例投递策略超时阈值失败处理强制送达系统通知如“你已被移出群聊”同步检查在线状态仅发给在线用户失败立即返回错误500ms客户端弹窗提示“发送失败请重试”尽力送达普通文本消息先广播给所有在线用户再异步写入离线队列无不提示用户后台静默重试异步补发离线用户历史消息用户重连后从 Redis Stream 拉取未读消息30s超时后标记为“已过期”不再重试// 消息投递主逻辑简化版 func deliverMessage(msg ChatMessage, groupID string) { onlineUsers, _ : getOnlineMembers(groupID) // 步骤1同步广播给在线用户强制送达 for _, uid : range onlineUsers { if err : broadcastToUser(msg, uid); err ! nil { log.Printf(broadcast to %s failed: %v, uid, err) // 记录失败但不中断循环 } } // 步骤2写入离线队列Redis List offlineKey : offline: groupID client.RPush(ctx, offlineKey, json.Marshal(msg)) client.Expire(ctx, offlineKey, 7*24*time.Hour) // 离线消息保留7天 }3.3 弱网用户保活用消息分片 服务端缓冲降低丢包率移动端弱网环境下单条大消息如含图片 base64极易被中间代理截断。我们强制要求文本消息 ≤ 4KBUTF-8 字节数超长消息自动分片fragment_id,total_fragments,content服务端为每个连接维护 3 条消息缓冲区若检测到连续 2 次PONG超时则降级为“低速模式”暂停新消息广播优先重发缓冲区未确认消息。// 分片消息示例客户端发送 { type: message_fragment, group_id: 1001, msg_id: abc123, fragment_id: 1, total_fragments: 3, content: base64... }血泪经验不要信任客户端的navigator.onLine。我们曾在线上发现 12% 的“在线”用户实际处于 2G 网络onLine返回true但 TCP 握手超时。真正的弱网检测必须基于 WebSocket ping-pong 延迟1000ms 视为弱网。4. 多人聊天系统的避坑指南5 个让团队加班到凌晨的真实问题与解法多人聊天系统上线前最怕的不是功能没做全而是那些文档里不写、教程里不提、但一上生产就暴雷的细节。以下是我们在 3 个 SaaS 产品中踩过的坑按发生频率排序4.1 现象群聊消息顺序错乱A 发的“1”出现在 B 发的“2”之后原因多个服务实例同时处理同一群聊消息Redis Stream 的XADD时间戳由本地机器生成不同服务器时钟偏差导致 ID 乱序如服务器 A 时钟快 200msB 慢 100ms。解决禁用本地时间戳改用 Redis 的TIME命令获取统一时间或更简单——用XADD chat:1001 * ...中的*让 Redis 自动生成单调递增 ID。实测后消息顺序 100% 严格按写入顺序。4.2 现象用户退出群聊后仍持续收到该群消息原因前端只销毁了 WebSocket 连接但未通知服务端清理其在 Redis Stream 消费者组中的注册信息导致消息持续推送到已关闭的连接触发大量write: broken pipe错误。解决WebSocketClose事件中必须调用XGROUP DELCONSUMER chat:1001 online_users user_a显式删除消费者。我们封装了defer cleanupConsumer(groupID, userID)确保执行。4.3 现象高峰期 Redis 内存暴涨INFO memory显示mem_clients占比超 60%原因未限制每个 Stream 的最大长度历史消息无限堆积尤其测试环境无人清理。解决对每个群聊 Stream 设置MAXLEN生产环境设为10000约 1GB/万条消息命令XADD chat:1001 MAXLEN ~ 10000 * ...。~表示近似裁剪性能更好。4.4 现象用户 A 发送消息后自己客户端收不到回显echo原因服务端广播时默认排除发送者但未考虑“自己发自己收”是 UI 必需体验否则用户会疑惑“我发成功了吗”。解决广播逻辑中增加if uid ! senderID { broadcast() } else { sendEcho() }且sendEcho必须走与广播相同的序列化路径避免字段缺失。4.5 现象凌晨 2 点 Redis 持久化 RDB 时聊天服务出现 3 秒级延迟原因RDB fork 子进程时会拷贝父进程内存页若 Redis 占用 4GB 内存fork 可能卡住数秒。解决关闭 RDB改用 AOF appendfsync everysec或升级 Redis 7.0 使用copy-on-write优化。我们选择前者因 AOF 重写可控且everysec模式下数据丢失窗口 ≤ 1 秒符合聊天场景容忍度。5. 压测与容量规划用真实流量模型跑出你的系统极限而不是猜很多人用ab或wrk压测但结果毫无参考价值——因为它们模拟的是 HTTP 请求而多人聊天是长连接 消息流。我们必须用连接生命周期模型来压测创建连接 → 加入群聊 → 发送随机消息 → 保持心跳 → 随机断连重连。5.1 构建可复现的压测脚本Locust WebSocketLocust 支持 WebSocket 协议且能模拟用户行为流。以下是我们生产环境使用的最小压测脚本可直接运行# locustfile.py from locust import HttpUser, task, between from locust_plugins.users import WebsocketUser import json import random class ChatUser(WebsocketUser): wait_time between(1, 3) # 每个用户操作间隔 1~3 秒 def on_start(self): # 1. 建立 WebSocket 连接 self.connect(/ws) # 2. 发送登录消息含 token login_msg {type: login, token: test_token_123} self.send(json.dumps(login_msg)) # 3. 加入固定群聊群ID1001 join_msg {type: join_group, group_id: 1001} self.send(json.dumps(join_msg)) task def send_message(self): # 模拟真实打字节奏70% 消息为短文本20% 为中等长度10% 为长消息触发分片 lengths [random.choice([5, 10, 15]), random.choice([50, 100, 150]), random.choice([500, 1000])] content_len random.choices(lengths, weights[70, 20, 10])[0] content x * content_len msg { type: chat_message, group_id: 1001, content: content, timestamp: int(time.time() * 1000) } self.send(json.dumps(msg)) task def heartbeat(self): # 每 15 秒发一次心跳 if self.environment.runner.stats.total.num_requests % 15 0: self.send({type:ping})运行命令# 启动 500 个并发用户每秒新增 10 个用户渐进加压 locust -f locustfile.py --host http://localhost:8080 --users 500 --spawn-rate 105.2 关键指标监控清单必须接入 Prometheus压测不是看 QPS而是看消息端到端 P99 延迟和连接存活率。我们监控以下 6 个核心指标指标名Prometheus 查询健康阈值说明websocket_connections_totalcount by (state) (websocket_connections{jobchat})online 95%state包含online,offline,errorredis_stream_lengthredis_stream_length{streamchat:1001} 10000防止 Stream 无限膨胀message_delivery_p99_mshistogram_quantile(0.99, rate(message_delivery_duration_seconds_bucket[1m])) 300ms从XADD到客户端onmessage的耗时redis_cpu_percent100 - (avg by (instance) (irate(redis_used_cpu_sys_seconds_total[5m])) * 100) 30%Redis CPU 使用率低于 30% 才有余量goroutines_totalgo_goroutines{jobchat} 5000Goroutine 泄漏预警正常应稳定在 2000~3000memory_usage_percent100 * (node_memory_MemTotal_bytes - node_memory_MemAvailable_bytes) / node_memory_MemTotal_bytes 75%服务端内存水位超 75% 开始 GC 压力大提示message_delivery_duration_seconds必须在XADD前打点在客户端onmessage回调里上报中间所有环节Redis 写入、Stream 读取、广播都要埋点。我们用 OpenTelemetry 自动注入不手动写start : time.Now()。5.3 容量公式根据你的硬件算出最大承载人数别信“单机支持 10 万连接”的宣传。真实容量取决于三个硬约束内存每个 WebSocket 连接约占用 32KBGogorilla/websocket默认 buffer1000 连接 ≈ 32MB文件描述符Linux 默认ulimit -n为 1024必须调至 ≥ 65535Redis 连接数每个群聊 Stream 需要 1 个连接每个消费者组需 1 个连接100 人群聊 ≈ 200 Redis 连接。我们总结出经验公式适用于 4C8G 云服务器最大群聊数 min( floor(可用内存 GB × 1024 × 1024 × 1024 / (32 × 1024)), // 内存约束 floor(ulimit -n / 2), // FD 约束 floor(Redis maxclients / 2) // Redis 约束 ) × 0.7 // 保留 30% 余量例如4C8G 服务器ulimit -n设为 65535Redismaxclients10000则内存8×1024³ / (32×1024) ≈ 262144 连接FD65535 / 2 ≈ 32767 连接Redis10000 / 2 ≈ 5000 连接→ 最终取5000 × 0.7 ≈ 3500人即单机最多支撑 35 个 100 人群聊。6. 消息去重与幂等性为什么你发两次“收到”对方只看到一次多人聊天系统里“重复发送”是高频操作——用户点发送后没反应再点一次结果对方收到两条一模一样的消息。这不是 UI 问题而是服务端缺乏幂等性设计。很多团队用“消息 ID 去重”但 ID 由客户端生成不可信也有用“内容哈希”但相同内容不同时间戳哈希不同。我们必须在服务端建立基于业务语义的幂等窗口。6.1 基于时间窗口的幂等校验推荐原理同一用户在 5 秒内发送的相同内容忽略空格、换行视为重复消息。实现简单、存储开销小、覆盖 99% 场景。// Go 实现 func isDuplicate(userID, content string) bool { client : redisClient key : dedup: userID // 标准化内容去除首尾空格合并连续空格 normalized : regexp.MustCompile(\s).ReplaceAllString(content, ) normalized strings.TrimSpace(normalized) // 用内容 SHA256 作唯一标识 hash : fmt.Sprintf(%x, sha256.Sum256([]byte(normalized))) // SETNX EXPIRE 原子操作Redis 6.2 可用 SET key value PX 5000 NX ok, _ : client.SetNX(ctx, key:hash, 1, 5*time.Second).Result() return !ok // 如果 set 失败说明已存在是重复消息 } // 在消息处理入口调用 func handleMessage(req ChatMessageRequest) error { if isDuplicate(req.UserID, req.Content) { log.Printf(duplicate message from %s: %s, req.UserID, req.Content) return nil // 静默丢弃不返回错误 } // 继续正常处理... }注意SETNXEXPIRE非原子可能因网络分区导致 key 永久存在。生产环境必须用 Redis 6.2 的SET key value PX 5000 NX或 Lua 脚本封装。6.2 基于客户端序列号的强幂等金融级场景当业务要求“绝对不重复”如红包消息、指令类消息需客户端携带单调递增序列号seq_no服务端用 Redis Sorted Set 存储每个用户的最大seq_no# 用户 123 最大序列号为 100 ZADD dedup_seq:123 100 123:100 # 新消息 seq_no101检查是否 100 ZREVRANGEBYSCORE dedup_seq:123 inf (100 LIMIT 0 1 # 若返回空则接受否则拒绝6.3 消息去重的边界什么情况下不该去重去重不是万能的。我们明确规定以下三类消息禁止去重系统通知类如“XXX 已加入群聊”即使内容相同每次事件都需独立推送带时间戳的指令如“定时提醒明天 9:00 开会”相同内容但不同时间戳语义不同富媒体消息图片、文件消息即使文件名相同二进制内容可能不同如用户反复编辑同一张图。我们的做法是在消息结构体中增加is_idempotent: bool字段由业务方决定是否开启去重服务端只执行策略不越权判断。最后说一句血泪教训不要在消息去重上追求 100% 准确。我们曾为追求“零误判”引入布隆过滤器结果内存暴涨 40%而实际业务中用户重复发送概率不足 0.3%用简单哈希 时间窗口已足够。工程的本质是权衡不是完美主义。希望帮到你。本文还有配套的精品资源点击获取
返回列表