ARTICLE DETAIL

资讯详情

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

WebSocket + Redis:从零构建万级并发实时消息系统的完整方案

WebSocket + Redis:从零构建万级并发实时消息系统的完整方案 早些年我维护过一个日活几十万的社区App最头疼的不是业务功能而是消息模块。最开始用HTTP轮询每3秒拉一次新消息用户量一上来服务器CPU直接打满手机电量也撑不住。后来切到WebSocket心跳、重连、在线状态这些问题一个个冒出来真正让我把并发稳定推上万级的是引入Redis做状态和路由协调之后。这篇文章就把这套方案的完整设计和实操经验拆开讲清楚包括为什么选WebSocket Redis、连接管理的几个硬骨头、压测数据和性能瓶颈以及几个典型的线上事故复盘。如果你正打算从零做一个实时消息系统或者已经在用WebSocket但并发上不去、经常出现连接异常和消息丢失这篇内容应该能帮你少走不少弯路。1. 为什么实时消息系统绕不开 WebSocket Redis 这个组合1.1 从轮询到长连接你等消息的方式决定了架构的复杂度先看最直观的问题HTTP轮询为什么扛不住万级并发假设1万用户在线每3秒轮询一次每秒大概产生3300个HTTP请求。这些请求大部分没有新消息但服务器依旧要完成TCP建连、HTTP头解析、业务查询、响应返回再销毁连接。这不仅仅是CPU开销的问题频繁建连和断连还会占用大量文件描述符和内核资源用户越多浪费越明显。WebSocket解决的就是这个核心矛盾一次HTTP握手升级成TCP长连接后续双向通信都走这条连接不再重复握手。1万在线用户对应1万条常驻连接消息量不大的时候服务器负载其实很轻。这就是数量级的差异而不是简单的优化。但长连接也带来了新问题连接是长久的状态的维护、断开的检测、服务端主动推送、多节点部署后的路由这些都成了必须处理的工程问题。这也是为什么很多项目用了WebSocket反而更麻烦——协议本身很简单难的是围绕连接的配套系统。1.2 Redis在链路里真正补上的空缺WebSocket管的是一条连接但一个大型实时系统真正需要的是一套连接的管理系统用户A在哪个节点上用户B在线不在线离线消息存在哪里这些信息不可能只放在某个进程的内存里——一旦服务多实例部署A在node1、B在node2的情况就处理不了。Redis在这里的角色说白了就是一个所有实例都能访问的共享状态层。它把连接路由表、在线状态、离线消息从单机内存里挪出来用统一的数据结构Hash、List、Stream等让每个节点都能实时读写。再加上Redis单线程模型下小对象操作极其快内存操作微秒级非常适合这种高频低数据量的实时场景。另外一个实际考虑是落库分离。聊天记录要长期保存一般写进MySQL或MongoDB但实时推送走MySQL完全不现实。Redis天然充当了热数据这层缓冲在线用户读路由、离线用户收消息都先走Redis只有需要持久化的数据才异步落库。1.3 万级并发背后的两个隐性瓶颈第一是连接本身的消耗。每个WebSocket连接都占一个socket、一部分内核内存和用户态buffer。默认的ulimit往往只有1024不开大连接数到几千就崩就算调大文件描述符单机支撑的连接数也有上限通常4核8G的云主机在8000到10000连接时CPU和内存已经很紧张。解决方向不是硬怼单机而是多节点横向扩展——这又回到了必须有Redis做路由协调的结论。第二是消息扇出fan-out。私聊是一对一投递压力很小群聊尤其全员群一条消息可能要推给几千上万个连接。如果所有推送都在一个进程里循环发送慢客户端还会拖慢后来的消息。扇出的瓶颈往往不在Redis而在应用层代码怎么控制批量和背压。2. 消息系统的三层架构连接、路由、消息各管一摊2.1 一次私聊消息的完整旅行先画清楚数据流。假设用户A要给用户B发一条文本消息链路如下A的客户端把消息通过WebSocket发送到A所在节点的网关服务。网关节点解析出消息内容、目标用户B调用消息服务接口。消息服务先查Redis路由表确认B当前在哪个节点ws:route哈希表。如果B在线且在同一节点直接往本地连接池里找到B的连接并推送如果B在别的节点则通过Redis Pub/Sub或内部RPC把消息转发给对应节点。如果B不在线消息进离线列表Redis List或Stream等B上线后按需拉取。这条链路里WebSocket网关只负责维持连接和收发消息业务逻辑比如敏感词过滤、消息落库全在消息服务层完成。分层的好处是每一层都能独立扩缩网关节点可以随时加减路由和离线数据都在Redis里不会因为某个节点宕机导致整个系统不可用。2.2 单机到分布式session信息从内存搬到Redis的关键一步很多人一开始的写法是这样的启动时建立一个全局的map[userId]*websocket.Conn消息来了一查map直接推。这种单机模型在500并发以内非常舒服代码少、延迟低。但一旦部署第二个节点问题就出现了A在node1登录B在node2登录A给B发消息时node1的本地map里根本查不到B的连接。我见过不少团队在这里硬扛做法是让所有连接都连到同一个节点或者靠Nginx的ip_hash把同一个用户钉在一个节点上——但这只是暂缓节点一挂整个区域用户全部掉线。正确的演进是把连接对象和连接路由信息分离连接对象*websocket.Conn仍然留在各节点内存里因为只有本节点才能真正往这个连接写入数据。路由信息用户ID - 节点ID外置到Redis所有节点共享。每次连接建立时写入路由心跳续期断开或超时后删除。这样加节点、踢节点都只需要调整Redis里的路由记录客户端无感知。这是从千级走向万级并发最关键的一步。2.3 为什么没有用MQ当核心总线有人会问节点间转发消息用Kafka或RabbitMQ不是更稳吗我聊一下实际取舍。消息系统对实时性要求很高用户发出一条消息期望对方1秒内收到而Kafka这类消息队列更擅长削峰填谷、批量吞吐单条消息的端到端延迟相对更高运维成本也更大。对于社交App的私聊和中小群聊场景消息量并没有那么大Redis Pub/Sub的广播能力和Stream的持久化能力已经足够。只有当你有海量群聊、需要大量消息堆积和回放时才值得引入Kafka。我现在的做法是实时链路用Redis消息落库和审计走MySQL如果日后群消息量提升可以在消息服务层再加一层MQ做削峰而不是一开始就上重武器。3. WebSocket连接管理的三个硬骨头鉴权、心跳、重连3.1 握手鉴权在协议升级之前就把非法连接挡掉WebSocket的握手本质是一次HTTP GET请求所以鉴权完全可以在Upgrade之前完成。流程是客户端在握手URL里带上短期有效的token比如JWT有效期10分钟服务端收到请求后先解析token、查出用户ID校验通过才调用Accept()升级协议校验失败直接返回403或关闭连接。这里有个细节很多浏览器WebSocket API不支持自定义Header常见的方案是把token放在QueryString里或放在Sec-WebSocket-Protocol子协议字段里。QueryString的问题是会留在Nginx或网关的访问日志里所以token一定要短期有效别把用户的长效凭据直接放在URL里。我当时用JWT 10分钟过期 每次重连重新获取既安全又不用引入复杂的临时凭证逻辑。# FastAPI websockets 风格示意鉴权流程 async def ws_endpoint(websocket: WebSocket): token websocket.query_params.get(token) user_id verify_jwt_token(token) # 校验失败返回None if not user_id: await websocket.close(code4401, reasonunauthorized) return await websocket.accept() # 建立连接后写Redis路由表 await register_route(user_id, NODE_ID, websocket)3.2 心跳机制参数怎么定才能既不误杀也不拖死WebSocket协议本身有Ping/Pong帧但实际项目我建议在应用层做一套业务心跳原因有两个一是应用层心跳可以附带额外数据客户端时间、网络类型等方便排查问题二是很多负载均衡器和Nginx对空闲连接有超时回收机制底层Ping帧未必能阻止中间链路断连。我的默认参数客户端每30秒发送一个{type:ping,ts:...}服务端记录last_seen服务端每60秒扫描一次连接池超过90秒没收到心跳就判定离线。心跳包用JSON的话一个包大概几十字节1万连接每秒的心跳处理量约几百条对Redis和网关CPU的影响很小。这里有一个我在生产环境踩过的坑Nginx对上游WebSocket连接的proxy_read_timeout默认是60秒如果客户端心跳间隔正好是60秒服务端会在超时边缘被Nginx掐断表现为线上用户周期性掉线。解决方法是客户端心跳调成30秒或者把Nginx的proxy_read_timeout调到75秒以上。两种我都试过双管齐下最稳。3.3 断线重连与消息补发别让用户感知到掉线移动端网络抖动是常态Wi-Fi切4G、地铁进隧道、App切后台这些都会导致连接断开。客户端要做的是自动重连但别用固定间隔疯狂重试。我采用指数退避加随机抖动1秒、2秒、4秒、8秒最大30秒每次重试加一个0到1000毫秒的随机偏移。目的是防止大量客户端同时掉线后在同一秒集体重连把服务端冲垮。消息补发比重连更难。解决思路是客户端本地维护一个自增序号last_seq每次收到消息更新重连成功后把last_seq带给服务端服务端从Redis里的消息序列中查出缺失部分补发。离线期间的消息数量通常不大用Redis List的一个小范围区间就能搞定。如果消息量大可以在消息服务里加一个专门的补发队列逐条推送而不是一次全量塞给连接避免瞬间打爆慢客户端。# 伪代码重连后按seq补发 async def on_reconnect(user_id, last_seq): missing get_messages_since(user_id, last_seq) # Redis List range操作 for msg in missing: await push_to_user(user_id, msg)4. Redis在消息链路里的三种角色在线状态、连接路由、离线暂存4.1 在线状态Hash 过期时间的组合用法最简单直接的做法是每个用户一个带过期时间的keySETEX user:online:{userId} 90 {nodeId}心跳每次续期服务端清扫任务扫描过期的key清理离线状态。这个方案代码好写但用户量大时key数量太多Redis内存碎片多而且逐key续期会产生大量写操作。我后来改成了Hash结构HSET ws:online {userId} {lastHeartbeatTs}。心跳只更新field的值不需要建新key后台扫描时HGETALL整个Hash剔除超过90秒的field。这样1万用户的在线状态就是一个Hash对象内存占用小管理也方便。要注意的是Hash大key在热点扫描时会有点开销所以按用户ID分片比如拆成10个Hash可以进一步降低单key压力。在线状态还有个细节前端展示的在线不应该只看WebSocket是否连着还要结合用户是否在活跃使用。很多App把在线定义为最近5分钟有操作而WebSocket连接可能还挂着App在后台两套数据要分开维护。用Redis存两个维度连接状态实时和用户活跃状态分钟级展示层取活跃状态推送层取连接状态。4.2 连接路由表userId到nodeId的映射怎么维护才不出错路由表是节点间转发的依据核心设计是一个HashHSET ws:route {userId} {nodeId}。连接建立时写入连接关闭时删除。但这里有一个隐藏的并发问题用户设备A断开的同时用户可能在另一台设备B上建立了新连接如果断开的旧连接清理逻辑先执行再写入新连接倒还好就怕后到的旧连接清理把新连接的路由也删了。解决办法是给每次连接分配一个connIdUUID或自增路由表里存的是nodeId : connId。删除前先比对connId只有匹配才删除。这样即使旧的关闭事件晚到也不会误删新连接的路由。我在线上遇到过几次用户突然收不到消息查了半天就是这种旧连接的defer在关闭时覆盖了新连接的路由后来加了connId比对就再没出过这类问题。4.3 离线消息与消息补偿List和Stream到底怎么选离线消息最朴素的实现是LPUSH user:{userId}:offline msg用户上线后BRPOP或LRANGE拉取。这个方案够用但有两个问题第一List没有消费确认机制如果客户端拉取后处理失败消息就丢了第二List不能做条件读取想拉取某个时间点之后的消息要自己遍历。Redis 5.0引入的Stream弥补了这些缺陷。Stream支持消费者组、消息ACK、按消息ID范围读取非常适合做离线消息和消息回放。我的实践是普通私聊场景用List就够了重要消息比如转账、订单通知走Stream。群聊的离线消息量更大也更需要进程组来分担消费压力Stream是更稳妥的选择。不管用哪种结构都得给离线消息设上限。我当时的做法是每条离线消息列表最多保留100条、超过则压缩成一条你有99条未读消息的聚合通知用户点进会话再从数据库拉完整记录。这样既控制了Redis内存也避免了用户上线时一次性拉取几百条消息造成的流量高峰。4.4 Redis数据结构选型小结场景数据结构理由在线状态HashfielduserId, value时间戳/节点ID批量管理key数量少连接路由HashfielduserId, valuenodeId:connId高频读写O(1)操作离线消息List或StreamList简单Stream支持ACK和范围读取节点间广播Pub/Sub或Stream实时性强按需使用消费者组未读计数String/HINCRBY原子自增避免并发覆盖5. 压测实测万级并发的上限在哪里瓶颈怎么找5.1 压测环境和工具怎么搭压测WebSocket系统和压测HTTP接口不太一样得先有能建立大量长连接的客户端。我用的压测工具是自写的Golang脚本因为gorilla/websocket库比较成熟开几百个goroutine模拟持续连接和收发消息。生产环境跑过几个公网工具但连接数上万时控制力和数据采集反而不如自写脚本。推荐压测时至少采三组数据并发连接数、消息收发吞吐每秒处理的消息条数、消息延迟P50/P95/P99。延迟这个指标最容易忽略很多人只盯着连接数结果1万连接建起来了但消息要几百毫秒才送达用户体验照样差。5.2 我拿到的实测数据我的压测环境网关节点为4核8G云主机内存8G、Redis独立部署在2核4G上客户端模拟用一台8核16G的压测机。技术栈是Golang gorilla/websocket go-redis。测试规模从2000连接逐步加到10000结果如下并发连接数网关CPU网关内存Redis CPUP99消息延迟200018%0.5G5%5ms400035%1.2G8%7ms600055%2.1G12%9ms800078%3.0G18%14ms10000因文件描述符耗尽连接失败3.8G22%无法统计可以看到单节点在8000连接附近还能工作但CPU已经接近80%延迟明显上升——注意延迟不只是CPU高还有内存分配的GC压力。而在10000连接时连接失败原因是我当时没调整ulimit -n系统默认的进程文件描述符上限直接卡住了。5.3 系统参数调优和瓶颈定位万级连接的单节点调优必须做这几件事ulimit -n调到65535甚至更高。net.ipv4.ip_local_port_range调大给大量短连接腾出本地端口。net.core.somaxconn调大避免握手时backlog溢出。tcp_tw_reuse开启减少TIME_WAIT对连接建立的消耗。调完这些参数单节点上万连接不是神话但要稳定支撑还得看应用层。压测时最容易出现的瓶颈有三个一是心跳扫描的定时器精度如果用标准库的time.Sleep循环连接多了会产生调度偏移我改成了time.Ticker二是消息序列化JSON解析在高频推送下CPU消耗非常明显后来把内部通信改成了MessagePack整体CPU降低了十几个百分点三是Redis操作串行化如果每个消息收发都同步等待Redis响应延迟会被拉高正确的做法是批量管道操作或异步处理。6. 复盘真实线上踩过的坑以及对应的解决方案6.1 连接数只涨不跌服务器内存越用越多现象很典型监控里在线人数没有明显增长但服务器内存和句柄数缓慢上升两三天后必须重启网关。一开始以为是用户真的变多了查了Redis在线状态发现并不是。后来扒代码发现客户端非正常断连比如App直接被杀、网络切换时服务端的ReadMessage会返回错误并退出循环但负责关闭连接的defer conn.Close()没有执行——于是连接对象永远留在本地连接池里。修补方案每次Read循环结束必须走统一的清理逻辑同时用心跳扫描作为兜底。我在每次心跳扫描时遍历本地连接池检查最后活跃时间超过阈值直接强制关闭并从池中移除。清理逻辑千万别只依赖客户端的Close帧因为它可能永远到不了服务端。6.2 心跳误判引发的重连风暴有段时间线上频繁出现一种雪崩某个区域网络抖动大量用户连接短暂中断客户端按照固定1秒间隔集体重连网关瞬间涌入几千个握手请求直接打满CPU然后触发健康检查失败负载均衡又把流量切到别的节点整体雪崩。复盘下来有两个教训。第一心跳超时阈值定得太紧当时心跳间隔30秒、超时判定45秒一次网络抖动超过15秒就会被误杀但实际上这些连接可能马上恢复。后来调整成心跳50秒、超时180秒宁可多等一会儿也不主动断开。第二重连没有加随机抖动所有客户端都在同一秒发起重连。修复就是在重连策略里加入0到5秒的随机延迟让请求平铺开。6.3 Pub/Sub丢消息的真相不是Redis不稳定而是设计没闭环用Redis Pub/Sub做节点间消息转发时我一度以为消息丢了是Redis的问题。实测发现Pub/Sub就是发后即忘如果某个节点在消息发布时短暂不可用重启、网络抖动它订阅的消息就永久丢失了。Redis本身没有能力补偿这些消息。解决思路分两类一是对关键消息走Redis Stream的消费者组读取后ACK处理失败的进入重试流程二是客户端重连时用last_seq做可靠性兜底。我的最终方案是私聊消息继续走Pub/Sub实时性好、代码简单群聊和系统通知走Stream需要可靠性和回放。不要试图让一套机制同时满足所有场景。6.4 慢消费者拖垮整条链路背压处理没做好这个坑花了我最多时间。某个群里有用户网络状况极差TCP发送缓冲区被占满服务端往这个连接写消息时阻塞住导致整个WebSocket网关的写循环卡死其他用户的消息也跟着延迟。WebSocket库一般都有写超时设置我当时没配后来给每次WriteMessage加上10秒超时超过就关闭连接并标记这个用户离线。同时给每个连接维护一个发送队列队列长度超过100就丢弃最老的消息并给客户端发一个消息过快请重连补拉的通知。这条策略上线后网关的P99延迟从100多毫秒降回20毫秒以内。// Go示例带write deadline的发送 conn.SetWriteDeadline(time.Now().Add(10 * time.Second)) err : conn.WriteMessage(websocket.TextMessage, payload) if err ! nil { // 关闭连接清理路由表 closeConnAndCleanRoute(userId, connId) }最后再说一点经验之谈实时消息系统最容易翻车的地方往往不在技术选型而在你根本没想清楚连接断开了之后该怎么办。与其追求一上来就完美不如先把链路跑通——单机版的内存map方案先验证业务逻辑等第二台机器真需要加进来时再把状态外置到Redis。这样每一步都有明确的驱动力代码也不会过度设计。如果你正在搭自己的消息系统希望这篇能帮你绕过我当年踩过的那些大坑。
返回列表