Go-Zero项目开发14: IM服务业务分析及WebSocket服务工程实现
纲要IM 核心业务与数据建模聊天记录与聊天会话存储方案对比客户端存储、服务端存储、混合存储数据库选型MySQLvsMongoDB聊天记录ID设计策略WebSocket服务工程化基于go-zero风格的项目目录结构服务对象封装Server结构体、Upgrader、启停方法配置加载与ServiceContext消息体定义与路由分发连接池管理map 读写锁鉴权接口抽象 (Auth) 与默认实现消息发送按用户列表或按连接对象灵活配置ServerOptions模式服务启动与测试验证总结IM 核心业务与数据建模在即时通讯系统中聊天记录和聊天会话是两块核心数据。聊天记录保存用户之间交换的每一条消息而聊天会话则记录用户当前正在与哪些对象个人或群组进行对话并维护未读消息数量。每次用户打开应用都会通过会话列表拉取最近的聊天记录。存储方案对比方案描述代表产品纯客户端存储聊天记录仅保存在本地服务端只暂存离线消息上线后推送并删除微信纯服务端存储所有记录存储在服务端客户端按需拉取钉钉混合存储客户端与服务端均保存客户端优先读取本地缺失时从服务端拉取服务端承担历史归档和离线推送本课程采用混合存储兼具了本地缓存的快速访问和服务端数据的可靠性是较为折中的选择。离线消息可以在服务端暂存待用户上线后推送并清理。数据库选型针对聊天记录和会话数据关系型数据库MySQL虽然擅长结构化数据和事务处理但在面对海量消息写入、字段结构多变的场景时容易遇到性能瓶颈往往需要分库分表增加系统复杂度。而文档型数据库MongoDB天然适合存储半结构化、高写入吞吐的数据具备更好的水平扩展能力。在后续实现中我们选择MongoDB存储聊天记录和用户会话因为它能灵活支持不同类型的消息体且写入性能优异无需事先定义严格的表结构。聊天记录 ID 设计消息记录的唯一标识通常有两种方案基于双方用户 ID 组合生成例如将发送者 ID 与接收者 ID 按固定规则拼接形成会话范围内的消息 ID。关联表方案通过一张关联表将消息 ID 与参与用户 ID 一一对应。方案一更简洁因为用户查询消息时总是先选定一个会话由双方 ID 确定再获取该会话下的消息列表。因此我们采用第一种方式将发送方和接收方 ID 组合作为消息记录的关联标识简化查询逻辑。点击会话双方ID消息ID用户 A聊天会话消息列表MongoDBWebSocket 服务工程化根据前文架构设计WebSocket 服务负责实时消息收发。在微服务体系内它是一个独立的功能模块但可以和 API 服务共用端口通过路由区分。项目目录结构我们继续沿用go-zero推荐的工程目录布局im-ws/ ├── etc/ │ └── im-ws.yaml ├── internal/ │ ├── config/ │ │ └── config.go │ ├── handler/ │ │ └── routes.go │ ├── logic/ │ │ └── server.go │ ├── svc/ │ │ └── service_context.go │ └── types/ │ └── message.go ├── main.go └── go.mod服务对象封装gorilla/websocket提供了基础的Upgrader我们需要在其之上构建一个高层服务结构Server将连接生命周期、鉴权、路由、连接池等统一管理。// internal/logic/server.gopackagelogicimport(net/httpsyncgithub.com/gorilla/websocketgithub.com/zeromicro/go-zero/core/logx)// Server 封装 WebSocket 服务typeServerstruct{addrstringupgrader websocket.Upgrader routesmap[string]HandlerFunc// 方法路由表connPool*ConnectionPool auth Auth// 鉴权接口opts*ServerOptions mu sync.RWMutex}// HandlerFunc 路由处理函数参数server当前连接消息体typeHandlerFuncfunc(server*Server,conn*websocket.Conn,msg*Message)// NewServer 创建服务并应用选项funcNewServer(addrstring,opts...ServerOption)*Server{s:Server{addr:addr,upgrader:websocket.Upgrader{CheckOrigin:func(r*http.Request)bool{returntrue// 生产环境应校验 Origin},},routes:make(map[string]HandlerFunc),connPool:NewConnectionPool(),}// 默认选项defaultOpts:ServerOptions{Auth:defaultAuth{},}for_,o:rangeopts{o(defaultOpts)}s.optsdefaultOpts s.authdefaultOpts.Authreturns}// Start 启动服务func(s*Server)Start()error{http.HandleFunc(/ws,s.ServeWS)logx.Infof(WebSocket server starting on %s,s.addr)returnhttp.ListenAndServe(s.addr,nil)}// Stop 停止服务预留资源清理func(s*Server)Stop(){logx.Info(WebSocket server stopped)}// AddRoutes 批量注册路由func(s*Server)AddRoutes(routes[]Route){for_,r:rangeroutes{s.routes[r.Method]r.Handler}}// Route 路由条目typeRoutestruct{MethodstringHandler HandlerFunc}配置加载与 ServiceContextgo-zero支持通过conf包从yaml文件加载配置我们定义Config结构体并嵌入rest.RestConf以获得日志、监听等内置能力。Name: im-ws Host: 0.0.0.0 Port: 8888// internal/config/config.gopackageconfigimportgithub.com/zeromicro/go-zero/resttypeConfigstruct{rest.RestConf// 可扩展自定义字段如 MongoDB 连接、RPC 地址等}// internal/svc/service_context.gopackagesvcimport(im-ws/internal/configim-ws/internal/logic)typeServiceContextstruct{Config config.Config WsServer*logic.Server}funcNewServiceContext(c config.Config)*ServiceContext{server:logic.NewServer(c.Host:strconv.Itoa(c.Port))returnServiceContext{Config:c,WsServer:server,}}启动入口// main.gopackagemainimport(flagstrconvim-ws/internal/configim-ws/internal/handlerim-ws/internal/svcgithub.com/zeromicro/go-zero/core/confgithub.com/zeromicro/go-zero/core/logx)varconfigFileflag.String(f,etc/im-ws.yaml,the config file)funcmain(){flag.Parse()varc config.Config conf.MustLoad(*configFile,c)logx.MustSetup(c.Log)ctx:svc.NewServiceContext(c)// 注册路由handler.RegisterRoutes(ctx)logx.Infof(Starting WebSocket service at %s:%d,c.Host,c.Port)ctx.WsServer.Start()}消息体与路由分发定义通用的消息结构体客户端与服务器均按照此格式通信// internal/types/message.gopackagetypesimportencoding/json// Message 通用消息结构typeMessagestruct{Methodstringjson:method// 方法名FromIDstringjson:from_idDatainterface{}json:data}// NewMessage 快速创建消息funcNewMessage(method,fromIDstring,datainterface{})*Message{returnMessage{Method:method,FromID:fromID,Data:data,}}// Marshal 序列化func(m*Message)Marshal()([]byte,error){returnjson.Marshal(m)}Server需要实现一个循环读取协程对每条消息进行路由分发。// internal/logic/server.go 补充方法// ServeWS WebSocket 连接处理入口func(s*Server)ServeWS(w http.ResponseWriter,r*http.Request){conn,err:s.upgrader.Upgrade(w,r,nil)iferr!nil{logx.Errorf(upgrade error: %v,err)return}deferconn.Close()// 鉴权if!s.auth.Authenticate(w,r){conn.WriteMessage(websocket.TextMessage,[]byte({error:auth failed}))return}userID:s.auth.UserID(r)// 将连接加入连接池s.connPool.Add(userID,conn)defers.connPool.Remove(userID,conn)// 消息处理循环for{msgType,payload,err:conn.ReadMessage()iferr!nil{ifwebsocket.IsUnexpectedCloseError(err,websocket.CloseGoingAway,websocket.CloseNormalClosure){logx.Errorf(read error: %v,err)}break}ifmsgType!websocket.TextMessage{continue}varmsg Messageiferr:json.Unmarshal(payload,msg);err!nil{logx.Errorf(unmarshal error: %v, payload: %s,err,string(payload))continue}// 路由查找handler,ok:s.routes[msg.Method]if!ok{resp:NewMessage(error,,method not found)data,_:resp.Marshal()conn.WriteMessage(websocket.TextMessage,data)continue}// 调用具体逻辑handler(s,conn,msg)}}连接池管理连接池使用两个map交叉存储用户 ID 与连接对象并配合读写锁保证并发安全。// internal/logic/connection.gopackagelogicimport(syncgithub.com/gorilla/websocket)typeConnectionPoolstruct{mu sync.RWMutex userConnmap[string]*websocket.Conn// userID - connconnUsermap[*websocket.Conn]string// conn - userID}funcNewConnectionPool()*ConnectionPool{returnConnectionPool{userConn:make(map[string]*websocket.Conn),connUser:make(map[*websocket.Conn]string),}}func(p*ConnectionPool)Add(userIDstring,conn*websocket.Conn){p.mu.Lock()deferp.mu.Unlock()p.userConn[userID]conn p.connUser[conn]userID}func(p*ConnectionPool)Remove(userIDstring,conn*websocket.Conn){p.mu.Lock()deferp.mu.Unlock()delete(p.userConn,userID)delete(p.connUser,conn)}func(p*ConnectionPool)GetConnByUserID(userIDstring)(*websocket.Conn,bool){p.mu.RLock()deferp.mu.RUnlock()conn,ok:p.userConn[userID]returnconn,ok}func(p*ConnectionPool)GetUserIDByConn(conn*websocket.Conn)(string,bool){p.mu.RLock()deferp.mu.RUnlock()id,ok:p.connUser[conn]returnid,ok}// GetAllUserIDs 获取所有在线用户 IDfunc(p*ConnectionPool)GetAllUserIDs()[]string{p.mu.RLock()deferp.mu.RUnlock()ids:make([]string,0,len(p.userConn))forid:rangep.userConn{idsappend(ids,id)}returnids}鉴权抽象鉴权逻辑因项目而异我们抽取接口Auth并提供基于请求头User-ID的默认实现。// internal/logic/auth.gopackagelogicimportnet/http// Auth 鉴权接口typeAuthinterface{Authenticate(w http.ResponseWriter,r*http.Request)boolUserID(r*http.Request)string}// defaultAuth 默认鉴权直接从请求头或参数获取用户 IDtypedefaultAuthstruct{}func(a*defaultAuth)Authenticate(w http.ResponseWriter,r*http.Request)bool{// 实际项目应校验 token此处简化为存在即合法userID:a.UserID(r)returnuserID!}func(a*defaultAuth)UserID(r*http.Request)string{// 优先从 URL 参数获取其次从 Header 获取ifuid:r.URL.Query().Get(user_id);uid!{returnuid}returnr.Header.Get(User-ID)}消息发送提供两种发送方式根据用户 ID 列表发送或直接对某个连接发送。// internal/logic/sender.gopackagelogicimport(github.com/gorilla/websocketgithub.com/zeromicro/go-zero/core/logx)// SendToUsers 向指定用户列表发送消息func(s*Server)SendToUsers(msg*Message,userIDs...string){for_,uid:rangeuserIDs{conn,ok:s.connPool.GetConnByUserID(uid)if!ok{logx.Infof(user %s offline, skip send,uid)continue}s.sendToConn(conn,msg)}}// SendToConn 向单个连接发送消息func(s*Server)SendToConn(conn*websocket.Conn,msg*Message){s.sendToConn(conn,msg)}func(s*Server)sendToConn(conn*websocket.Conn,msg*Message){data,err:msg.Marshal()iferr!nil{logx.Errorf(marshal message error: %v,err)return}iferr:conn.WriteMessage(websocket.TextMessage,data);err!nil{logx.Errorf(write message error: %v,err)}}灵活配置ServerOptions为了增强扩展性我们使用选项模式Functional Options来配置鉴权逻辑等可选参数。// internal/logic/server_options.gopackagelogictypeServerOptionsstruct{Auth Auth}typeServerOptionfunc(*ServerOptions)funcWithAuth(auth Auth)ServerOption{returnfunc(opts*ServerOptions){opts.Authauth}}在NewServer中应用这些选项后开发者可以通过logic.NewServer(addr, logic.WithAuth(myAuth))灵活替换鉴权实现。路由注册示例获取在线用户我们注册一个getOnlineUsers方法用于返回当前所有在线用户 ID 列表。该示例演示路由定义、消息构造与发送的完整流程。// internal/handler/routes.gopackagehandlerimport(im-ws/internal/logicim-ws/internal/svcim-ws/internal/types)funcRegisterRoutes(ctx*svc.ServiceContext){server:ctx.WsServer server.AddRoutes([]logic.Route{{Method:getOnlineUsers,Handler:func(s*logic.Server,conn*websocket.Conn,msg*types.Message){// 获取所有在线用户 IDusers:s.connPool.GetAllUserIDs()// 构造响应消息resp:types.NewMessage(onlineUsers,server,users)s.SendToConn(conn,resp)},},// 更多路由在此注册...})}客户端只需发送{method:getOnlineUsers,from_id:user123,data:null}即可得到在线用户列表。测试验证启动服务后使用ApiPost或websocat等工具连接ws://127.0.0.1:8888/ws?user_idalice发送上述请求即可收到服务器返回的在线用户列表。若再连接另一个客户端如user_idbob再次调用getOnlineUsers将会看到两人均在线。$ websocat ws://127.0.0.1:8888/ws?user_idalice{method:getOnlineUsers,from_id:alice,data:null}{method:onlineUsers,from_id:server,data:[alice]}通过这种方式我们验证了WebSocket服务的连接管理、鉴权、路由分发、消息收发全部链路畅通。总结本文从 IM 业务的聊天记录与会话数据设计出发对比了不同存储方案和数据库选型最终确定了混合存储 MongoDB的技术路线。接着我们基于go-zero工程规范完整实现了WebSocket服务的核心组件服务对象封装、配置加载、消息路由、连接池、鉴权抽象以及消息发送。通过选项模式各模块间保持松耦合便于后续集成真实的Token校验、消息持久化等功能。至此IM 服务的实时通信骨架已经搭建完毕为后续实现私聊、群聊、已读回执等高级业务打下了坚实基础。