news 2026/8/19 20:11:52

Go实战:实时聊天系统架构

作者头像

张小明

前端开发工程师

1.2k 24
文章封面图
Go实战:实时聊天系统架构

Go实战:实时聊天系统架构

摘要: 本篇讲解Go实时聊天系统架构,实现WebSocket Hub管理连接,Redis pub/sub做多节点消息广播,消息持久化到Redis List,在线状态用Redis维护,分享百万连接导致内存溢出的踩坑经验,对比单机Hub、Redis pub/sub、Kafka消息队列三种广播方案。

开篇故事

今年3月我们上线了一个在线教育直播聊天功能,单场直播同时在线5000人,老师发一条消息要推送给所有在线用户。上线第一天就收到投诉,用户说消息延迟十几秒才收到,有的根本收不到。

排查发现我们只部署了一个服务实例,5000个WebSocket连接全部集中在一台机器上。老师发消息时,服务端要遍历5000个连接逐个推送,CPU直接打满,消息排队延迟严重。更麻烦的是直播高峰同时在线能到5万,单机根本扛不住。

这次我把聊天系统的架构设计写清楚,重点讲WebSocket Hub怎么管理连接,多节点场景怎么用Redis pub/sub做消息广播。

一、WebSocket Hub管理连接

WebSocket连接建立后需要管理起来,谁在线、给谁推送消息。Hub就是连接管理器,维护所有活跃的WebSocket连接,负责消息分发。

packagechatimport("log""net/http""sync""github.com/gorilla/websocket")// Hub 连接管理器// 维护所有在线WebSocket连接,负责消息广播typeHubstruct{mu sync.RWMutex clientsmap[string]*Client// 在线用户连接,key为userIDregisterchan*Client// 注册通道unregisterchan*Client// 注销通道broadcastchan[]byte// 广播通道}// Client 单个WebSocket连接typeClientstruct{IDstring// 用户IDconn*websocket.Conn// WebSocket连接sendchan[]byte// 发送缓冲通道}// NewHub 创建HubfuncNewHub()*Hub{return&Hub{clients:make(map[string]*Client),register:make(chan*Client,100),unregister:make(chan*Client,100),broadcast:make(chan[]byte,256),}}// Run 启动Hub,处理注册注销和广播// 单独goroutine运行,是Hub的主循环func(h*Hub)Run(){for{select{caseclient:=<-h.register:// 新连接注册,加入clients maph.mu.Lock()h.clients[client.ID]=client h.mu.Unlock()log.Printf("用户上线: %s, 在线: %d",client.ID,len(h.clients))caseclient:=<-h.unregister:// 连接断开,从clients移除h.mu.Lock()if_,ok:=h.clients[client.ID];ok{delete(h.clients,client.ID)close(client.send)// 关闭发送通道}h.mu.Unlock()log.Printf("用户下线: %s, 在线: %d",client.ID,len(h.clients))casemessage:=<-h.broadcast:// 广播消息给所有在线用户h.mu.RLock()for_,client:=rangeh.clients{select{caseclient.send<-message:default:// 发送缓冲满了,说明客户端消费慢// 被动断开慢消费者}}h.mu.RUnlock()}}}// upgrader WebSocket升级器varupgrader=websocket.Upgrader{CheckOrigin:func(r*http.Request)bool{returntrue// 允许跨域,生产环境需校验},}// HandleWebSocket 处理WebSocket连接请求func(h*Hub)HandleWebSocket(w http.ResponseWriter,r*http.Request){// 升级HTTP为WebSocketconn,err:=upgrader.Upgrade(w,r,nil)iferr!=nil{log.Printf("WebSocket升级失败: %v",err)return}// 从URL参数获取用户IDuserID:=r.URL.Query().Get("uid")client:=&Client{ID:userID,conn:conn,send:make(chan[]byte,64),}// 注册到Hubh.register<-client// 启动读写goroutinegoh.readPump(client)goh.writePump(client)}// readPump 读取客户端消息func(h*Hub)readPump(client*Client){deferfunc(){h.unregister<-client client.conn.Close()}()for{_,message,err:=client.conn.ReadMessage()iferr!=nil{break// 读取失败,连接断开}// 收到消息,交给广播通道h.broadcast<-message}}// writePump 向客户端推送消息func(h*Hub)writePump(client*Client){deferclient.conn.Close()for{select{casemessage,ok:=<-client.send:if!ok{return// 通道关闭,退出}iferr:=client.conn.WriteMessage(websocket.TextMessage,message,);err!=nil{break}}}}

Hub用channel做注册、注销、广播的通信,避免直接操作map加锁的麻烦。发送缓冲满了的客户端会被跳过,防止一个慢消费者拖垮整个广播。

二、Redis pub/sub多节点广播

单机Hub只解决了一台机器内的消息分发。多节点部署时,用户A连在节点1,用户B连在节点2,用户A发消息用户B收不到,因为节点2不知道节点1有什么消息。Redis pub/sub解决这个问题,所有节点订阅同一个频道,消息发到Redis,Redis广播给所有订阅节点。

packagechatimport("context""encoding/json""log""time""github.com/redis/go-redis/v9")// Message 聊天消息typeMessagestruct{Fromstring`json:"from"`// 发送者IDContentstring`json:"content"`// 消息内容Roomstring`json:"room"`// 聊天室IDTimeint64`json:"time"`// 时间戳}// RedisBridge Redis pub/sub桥接器// 负责跨节点消息广播和消息持久化typeRedisBridgestruct{client*redis.Client hub*Hub channelstring// pub/sub频道名}// NewRedisBridge 创建Redis桥接器funcNewRedisBridge(client*redis.Client,hub*Hub,channelstring)*RedisBridge{return&RedisBridge{client:client,hub:hub,channel:channel,}}// Publish 发布消息到Redis// 消息先发到Redis频道,所有节点都能收到func(r*RedisBridge)Publish(ctx context.Context,msg*Message)error{data,err:=json.Marshal(msg)iferr!=nil{returnerr}// 发布到Redis频道returnr.client.Publish(ctx,r.channel,data).Err()}// Subscribe 订阅Redis频道// 每个节点启动时调用,接收跨节点消息func(r*RedisBridge)Subscribe(ctx context.Context){// 订阅指定频道sub:=r.client.Subscribe(ctx,r.channel)defersub.Close()// 接收消息channelch:=sub.Channel()for{select{case<-ctx.Done():returncasemsg,ok:=<-ch:if!ok{return}// 收到Redis广播的消息,分发给本地Hub// 本地Hub再推送给这台机器上的在线用户r.hub.broadcast<-[]byte(msg.Payload)}}}// PersistMessage 消息持久化// 异步写入Redis List,用户可查看历史消息func(r*RedisBridge)PersistMessage(ctx context.Context,msg*Message)error{key:="chat:history:"+msg.Room data,_:=json.Marshal(msg)// 用Redis List存储历史消息,限制最多1000条pipe:=r.client.Pipeline()pipe.LPush(ctx,key,data)pipe.LTrim(ctx,key,0,999)// 保留最近1000条_,err:=pipe.Exec(ctx)returnerr}// SetOnline 设置用户在线状态// 用Redis记录在线用户,支持跨节点查询func(r*RedisBridge)SetOnline(ctx context.Context,userID,nodeIDstring)error{key:="chat:online:"+userID// 记录用户在哪个节点上线,设置30秒过期returnr.client.Set(ctx,key,nodeID,30*time.Second).Err()}// IsOnline 查询用户是否在线func(r*RedisBridge)IsOnline(ctx context.Context,userIDstring)(bool,error){key:="chat:online:"+userID result,err:=r.client.Exists(ctx,key).Result()iferr!=nil{returnfalse,err}returnresult==1,nil}

消息流程变成: 用户发消息到节点1,节点1发到Redis频道,所有节点订阅频道收到消息,各自分发给本机的在线用户。消息持久化异步写Redis List,不阻塞主流程。

三、踩坑经验:百万连接导致内存溢出

这个坑差点让我被开除。我们的聊天系统上线3个月,用户量涨到80万,高峰同时在线30万连接。某天凌晨收到告警,服务OOM重启了。

排查发现每个WebSocket连接占用内存远超预期。我们用的gorilla/websocket,每个连接除了底层TCP连接,还有读缓冲区、写缓冲区、send channel。粗算每个连接占8KB内存,30万连接就是2.4GB,加上goroutine栈内存,单机32G内存居然不够。

问题出在两个地方。第一,每个连接启动了两个goroutine(readPump和writePump),30万连接就是60万个goroutine,每个goroutine初始栈8KB,光栈内存就4.8GB。第二,send channel缓冲设了64,30万连接最坏情况缓冲区占150GB。

修复方案是控制单机连接数,加连接上限保护,同时调小缓冲区。

packagechatimport("net/http""sync/atomic")// SafeHub 带连接数保护的Hub// 防止连接数失控导致内存溢出typeSafeHubstruct{hub*Hub maxClientsint32// 最大连接数curClientsint32// 当前连接数(原子操作)}// NewSafeHub 创建带保护的Hub// maxClients: 单机最大连接数,按内存预算计算funcNewSafeHub(maxClientsint32)*SafeHub{return&SafeHub{hub:NewHub(),maxClients:maxClients,}}// GetHub 获取底层Hub,用于启动Run循环和readPumpfunc(h*SafeHub)GetHub()*Hub{returnh.hub}// HandleWebSocketSafe 带连接数检查的连接处理func(h*SafeHub)HandleWebSocketSafe(w http.ResponseWriter,r*http.Request,){// 原子检查当前连接数是否超限current:=atomic.LoadInt32(&h.curClients)ifcurrent>=h.maxClients{// 超过上限,拒绝新连接// 客户端会重试连到其他节点http.Error(w,"连接数已满",http.StatusServiceUnavailable)return}// 连接数加一atomic.AddInt32(&h.curClients,1)// 委托给Hub处理实际连接// 连接断开时在unregister中减一(需扩展Hub)h.hub.HandleWebSocket(w,r)}// Decrement 连接断开时调用,计数减一func(h*SafeHub)Decrement(){atomic.AddInt32(&h.curClients,-1)}

单机连接数上限按内存预算算。32G内存的机器,留8G给系统和进程本身,24G给连接。每个连接按2KB算(调小缓冲后),最多12万连接。设上限10万留足余量。超过上限的新连接返回503,负载均衡会重试到其他节点。

四、对比分析

广播方案延迟多节点支持消息可靠复杂度
单机Hub极低不支持丢消息风险
Redis pub/sub支持不保证
Kafka队列支持高(持久化)
直接遍历推送不支持

单机Hub延迟最低,但不能多节点。Redis pub/sub是多节点广播的主流方案,延迟低但消息不保证可靠,订阅者断开期间的消息会丢。Kafka保证消息可靠但延迟高,适合对消息完整性要求高的场景。直接遍历推送只适合几十人的小聊天室,上百人就会卡。

总结

聊天系统的核心是连接管理和消息分发。WebSocket Hub管单机连接,Redis pub/sub解决多节点广播。每连接的内存开销要算清楚,goroutine数量和缓冲区大小直接影响能承载多少连接。单机连接数必须有上限保护,超了拒绝新连接让负载均衡重试。下一篇我们聊支付系统设计,重点讲幂等性怎么保证。

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

Prompt 管理架构翻车记:当模型路由把用户订单变成诗歌时

Prompt 管理架构翻车记:当模型路由把用户订单变成诗歌时 当订单变成十四行诗:生成式AI在电商系统的实战教训与架构升级 灰度发布第三天,运营同事紧急截图丢进群聊--客户提交的电商订单在系统里变成了一首十四行诗。我盯着屏幕上的莎士比亚风格商品描述("汝之洗衣机,乃洁净…

作者头像 李华
网站建设 2026/8/19 20:07:09

2026国赛C题论文提分(十七):三线表、算法流程图与物理示意图的Visio/Python绘制规范

摘要 在数学建模竞赛中,论文的可读性与专业性往往直接影响评审专家的第一印象与最终评分。三线表、算法流程图与物理示意图作为论文中展示数据、逻辑与机理的三大可视化支柱,其绘制规范程度集中体现了参赛队伍的科学写作素养。本文系统梳理了三线表的结构要素与排版要点,深…

作者头像 李华
网站建设 2026/8/19 20:06:30

MediaBrowser 网格视图实战:快速实现照片墙与视频缩略图浏览

MediaBrowser 网格视图实战&#xff1a;快速实现照片墙与视频缩略图浏览 【免费下载链接】MediaBrowser &#x1f3de; A simple iOS photo and video browser with optional grid view, captions and selections written in Swift5.0 项目地址: https://gitcode.com/gh_mirr…

作者头像 李华
网站建设 2026/8/19 20:05:09

tQuery 游戏开发实战:TunnelGL 完整游戏源码深度解析

tQuery 游戏开发实战&#xff1a;TunnelGL 完整游戏源码深度解析 【免费下载链接】tquery extension system for three.js 项目地址: https://gitcode.com/gh_mirrors/tq/tquery 在 Web 3D 游戏开发领域&#xff0c;tQuery 是一个基于 three.js 的轻量级扩展系统&#x…

作者头像 李华