Go-Zero项目开发17: IM私聊功能实现与消息存储设计 纲要引言私聊业务流程分析消息发送与转发核心处理步骤消息模式推与拉Push模式Pull模式消息存储策略读扩散与写扩散读扩散写扩散本项目的选择读扩散 ChatLog 用户会话记录私聊功能实现数据模型设计ChatLog消息记录消息类型枚举会话 ID 生成MongoDB连接与配置消息结构体定义通用消息体Message私聊业务数据ChatData推送数据PushData错误消息构造业务逻辑ChatLogic处理私聊消息记录消息到MongoDB实时推送给接收方路由注册与Conn扩展防止重复登录与连接管理测试验证总结引言在前几篇文章中我们基于go-zero框架成功搭建了WebSocket服务实现了JWT鉴权和心跳检测。现在IM服务的基础设施已经准备就绪接下来就要进入即时通讯的核心业务——私聊。本文将首先分析私聊的实现流程介绍消息传递中的“推拉”模式以及消息存储中的“读扩散”与“写扩散”策略。随后我们将结合MongoDB在已有的WebSocket服务基础上完整实现用户与好友之间的私聊功能包括消息持久化与实时推送。私聊业务流程分析私聊是最基本的一对一通信场景。其核心流程如下图所示接收方客户端MongoDBWebSocket 服务发送方客户端接收方客户端MongoDBWebSocket 服务发送方客户端发送私聊消息 (含接收者ID、内容等)从 JWT Token 获取发送者ID将消息写入对应会话若 B 在线实时推送消息服务端在处理私聊时需要完成以下四件事解析消息类型区分文本、图片、语音等不同类型为后续扩展预留空间。确定会话标识根据发送者和接收者生成唯一的会话 ID用于消息归集。持久化消息将聊天记录存入数据库以便历史消息查询和离线消息拉取。转发消息从连接池中查找接收者的连接并将消息实时推送给对方。消息模式推与拉在即时通讯中消息的投递方式分为两种Push推服务端主动将消息发送给客户端。接收方在线时新消息由WebSocket立即推送客户端被动接收延时极低。Pull拉客户端主动向服务端请求数据。例如用户上线后拉取离线消息或向上滑动加载更早的历史记录。我们的系统采用实时推 离线拉的混合模式正常聊天通过Push实时送达当用户离线时消息暂存在数据库待其上线后通过API主动拉取。消息存储策略读扩散与写扩散消息的存储方式直接影响系统性能。常见的两种模式如下策略写入次数读取方式优点缺点读扩散每条消息写一次存于公共会话所有参与者都读同一份写入简单节省存储读取压力大尤其群聊已读/未读状态难维护写扩散每条消息为每个接收者复制一份每个用户只读自己的“收件箱”读取快个性化状态易实现写入放大严重大群尤为明显本项目的选择读扩散 ChatLog 用户会话记录在私聊场景中消息量相对可控我们采用读扩散作为基础方案。具体做法由发送方和接收方的 ID 经排序后拼接生成一个唯一的ChatId作为该对用户之间所有消息的标识。所有聊天记录存入chat_log集合每条文档都包含chat_id字段。单独维护一个用户会话列表集合记录每位用户在每个会话中的最后读取位置从而支持未读计数和离线消息的精确拉取。这种方案写操作仅一次读操作通过chat_id快速过滤性能满足私聊需求同时可方便地扩展至群聊需要额外的用户-会话状态维护。私聊功能实现数据模型设计ChatLog消息记录定义ChatLog结构体映射到MongoDB的chat_log集合// internal/model/chat_log.gopackagemodelimport(timego.mongodb.org/mongo-driver/bson/primitive)// MessageType 消息类型typeMessageTypeintconst(TextMessage MessageType1ImageMessage MessageType2VoiceMessage MessageType3)// ChatLog 聊天记录typeChatLogstruct{ID primitive.ObjectIDbson:_id,omitempty json:idChatIdstringbson:chat_id json:chat_id// 会话IDFromIdstringbson:from_id json:from_id// 发送者IDToIdstringbson:to_id json:to_id// 接收者IDType MessageTypebson:type json:type// 消息类型Contentstringbson:content json:content// 消息内容SendTime time.Timebson:send_time json:send_time// 发送时间}ChatLogModel数据库操作封装Insert方法用于写入消息// internal/model/chat_log_model.gopackagemodelimport(contexttimego.mongodb.org/mongo-driver/mongo)typeChatLogModelstruct{collection*mongo.Collection}funcNewChatLogModel(db*mongo.Database,collectionNamestring)*ChatLogModel{returnChatLogModel{collection:db.Collection(collectionName),}}func(m*ChatLogModel)Insert(ctx context.Context,log*ChatLog)error{log.SendTimetime.Now()_,err:m.collection.InsertOne(ctx,log)returnerr}会话 ID 生成私聊会话的ChatId由双方的 ID 按字典序拼接而成保证无论哪一方发起都使用同一个会话 ID。// internal/logic/chat_utils.gopackagelogicfuncGenerateChatId(uid1,uid2string)string{ifuid1uid2{returnuid1_uid2}returnuid2_uid1}MongoDB 连接与配置在internal/config/config.go中添加MongoDB配置typeConfigstruct{rest.RestConf JwtAuth JwtAuthConf MaxIdleTime time.Duration MongoDB MongoDBConf}typeMongoDBConfstruct{URIstringDatabasestring}etc/im-ws.yaml配置示例Name:im-wsHost:0.0.0.0Port:8888MaxIdleTime:30sJwtAuth:Secret:your-secret-keyMongoDB:URI:mongodb://root:password192.168.1.100:27017Database:im_chat在ServiceContext中初始化MongoDB客户端并创建ChatLogModel// internal/svc/service_context.gopackagesvcimport(contexttimestrconvim-ws/internal/configim-ws/internal/logicim-ws/internal/modelgithub.com/zeromicro/go-zero/core/logxgo.mongodb.org/mongo-driver/mongogo.mongodb.org/mongo-driver/mongo/options)typeServiceContextstruct{Config config.Config WsServer*logic.Server MongoClient*mongo.Client ChatLogModel*model.ChatLogModel}funcNewServiceContext(c config.Config,auth logic.Auth)*ServiceContext{clientOptions:options.Client().ApplyURI(c.MongoDB.URI)client,err:mongo.Connect(context.Background(),clientOptions)iferr!nil{logx.Must(err)}iferrclient.Ping(context.Background(),nil);err!nil{logx.Must(err)}db:client.Database(c.MongoDB.Database)chatLogModel:model.NewChatLogModel(db,chat_log)server:logic.NewServer(c.Host:strconv.Itoa(c.Port),logic.WithAuth(auth),logic.WithMaxIdleTime(c.MaxIdleTime),)returnServiceContext{Config:c,WsServer:server,MongoClient:client,ChatLogModel:chatLogModel,}}消息结构体定义我们沿用并扩展之前的Message结构并新定义了ChatData和PushData。// internal/types/message.gopackagetypes// Message 通用消息结构客户端与服务端通信typeMessagestruct{Methodstringjson:methodFromIDstringjson:from_idDatainterface{}json:data// 业务数据如 ChatData}// ChatData 私聊/群聊消息的具体业务数据typeChatDatastruct{ChatTypeintjson:chat_type// 1:私聊, 2:群聊ToIdstringjson:to_idContentstringjson:contentMsgTypeintjson:msg_type// 1:text, 2:image, 3:voice}// PushData 推送给接收方的消息格式typePushDatastruct{FromIdstringjson:from_idToIdstringjson:to_idContentstringjson:contentMsgTypeintjson:msg_typeSendTimeint64json:send_time}错误消息构造为方便向客户端返回统一格式的错误我们在Server上增加SendError方法// internal/logic/server.go 补充func(s*Server)SendError(conn*Conn,method,errMsgstring){errResp:types.Message{Method:method_err,FromID:server,Data:errMsg,}data,_:json.Marshal(errResp)conn.WriteMessage(websocket.TextMessage,data)}业务逻辑ChatLogicChatLogic负责处理聊天消息的存储与转发。// internal/logic/chat_logic.gopackagelogicimport(contextencoding/jsontimeim-ws/internal/modelim-ws/internal/typesgithub.com/gorilla/websocketgithub.com/zeromicro/go-zero/core/logx)typeChatLogicstruct{server*Server chatLogModel*model.ChatLogModel}funcNewChatLogic(s*Server,chatLogModel*model.ChatLogModel)*ChatLogic{returnChatLogic{server:s,chatLogModel:chatLogModel,}}// HandleChat 统一处理聊天消息私聊或群聊func(l*ChatLogic)HandleChat(conn*Conn,msg*types.Message){// 将 interface{} 转为具体结构dataBytes,err:json.Marshal(msg.Data)iferr!nil{l.server.SendError(conn,msg.Method,invalid data format)return}varchatData types.ChatDataiferrjson.Unmarshal(dataBytes,chatData);err!nil{l.server.SendError(conn,msg.Method,invalid chat data)return}// 获取发送者ID已在鉴权时通过 context 传递给连接fromID:conn.UserID()toID:chatData.ToIdiffromID||toID{l.server.SendError(conn,msg.Method,missing user id)return}// 根据聊天类型分别处理ifchatData.ChatType1{l.handlePrivateChat(fromID,toID,chatData,conn)}elseifchatData.ChatType2{// 群聊逻辑后续补充l.server.SendError(conn,msg.Method,group chat not implemented yet)}else{l.server.SendError(conn,msg.Method,unsupported chat type)}}func(l*ChatLogic)handlePrivateChat(fromID,toIDstring,chatData types.ChatData,conn*Conn){// 生成会话IDchatID:GenerateChatId(fromID,toID)// 构建 ChatLog 并写入 MongoDBchatLog:model.ChatLog{ChatId:chatID,FromId:fromID,ToId:toID,Type:model.MessageType(chatData.MsgType),Content:chatData.Content,}iferr:l.chatLogModel.Insert(context.Background(),chatLog);err!nil{logx.Errorf(Failed to insert chat log: %v,err)l.server.SendError(conn,chat,message store failed)return}// 构建推送给接收方的消息pushData:types.PushData{FromId:fromID,ToId:toID,Content:chatData.Content,MsgType:chatData.MsgType,SendTime:time.Now().Unix(),}pushMsg:types.Message{Method:push_message,FromID:fromID,Data:pushData,}// 实时推送查找接收方连接并发送iftargetConn,ok:l.server.connPool.GetConnByUserID(toID);ok{payload,_:json.Marshal(pushMsg)iferr:targetConn.WriteMessage(websocket.TextMessage,payload);err!nil{logx.Errorf(Push to user %s failed: %v,toID,err)}}else{// 用户不在线消息已存入数据库待上线后拉取logx.Infof(User %s offline, message stored,toID)}}路由注册与Conn扩展为方便业务逻辑中获取当前连接的用户 ID我们为Conn结构体增加userID字段并在鉴权成功后进行设置。首先在Conn结构体中添加字段和方法// internal/logic/connection.go 补充typeConnstruct{*websocket.Conn server*Server lastRead time.Time maxIdle time.Duration closeChchanstruct{}mu sync.Mutex once sync.Once userIDstring// 新增绑定的用户ID}// SetUserID 设置用户IDfunc(c*Conn)SetUserID(uidstring){c.userIDuid}// UserID 获取用户IDfunc(c*Conn)UserID()string{returnc.userID}然后在ServeWS方法中鉴权通过后立即设置用户 ID。下面是ServeWS的关键片段// 鉴权成功后获取用户ID并设置到连接上userID:s.auth.UserID(r)conn.SetUserID(userID)// 防止重复登录踢掉旧连接ifoldConn,ok:s.connPool.GetConnByUserID(userID);ok{oldConn.Close()}s.connPool.Add(userID,conn)defers.connPool.Remove(userID,conn)这样后续的所有业务处理方法如ChatLogic.HandleChat都可以通过conn.UserID()安全地获取当前操作的用户身份无需再从请求中重复解析。相应的路由注册保持不变只需在初始化时创建ChatLogic并绑定路由// internal/handler/routes.gopackagehandlerimport(im-ws/internal/logicim-ws/internal/svcim-ws/internal/types)funcRegisterRoutes(ctx*svc.ServiceContext){server:ctx.WsServer chatLogic:logic.NewChatLogic(server,ctx.ChatLogModel)server.AddRoutes([]logic.Route{{Method:chat,Handler:func(s*logic.Server,conn*logic.Conn,msg*types.Message){chatLogic.HandleChat(conn,msg)},},// ... 其他路由})}注意路由处理函数的第二个参数类型已由*websocket.Conn改为*logic.Conn这与我们前文对Server路由的改造保持一致确保类型统一。防止重复登录与连接管理在ServeWS中我们已经加入了踢掉旧连接的逻辑。当同一个用户再次连接时旧的连接会被关闭并从连接池中移除然后添加新连接。这一机制保证了每个用户同时只保留一个有效连接避免消息重复投递或投递到已失效的连接。测试验证分别用两个用户登录获取JWT Token。使用websocat或ApiPost连接到WebSocket服务并携带Token。用户 A 发送私聊消息{method:chat,from_id:userA,data:{chat_type:1,to_id:userB,content:Hello, B!,msg_type:1}}用户 B 实时收到推送消息{method:push_message,from_id:userA,data:{from_id:userA,to_id:userB,content:Hello, B!,msg_type:1,send_time:1698765432}}断开 B 的连接A 再次发送消息服务端日志提示User userB offline, message stored消息已成功写入MongoDB的chat_log集合。B 重新上线后可通过API拉取离线消息。总结本文从理论到实践完整阐述了私聊功能的实现过程。我们首先分析了私聊的业务流程回顾了实时通讯中“推”与“拉”两种模式进而对比了“读扩散”与“写扩散”两种存储方案并选择了适合当前项目规模的读扩散策略。在代码实现层面我们设计了ChatLog数据模型集成了MongoDB定义了清晰的消息结构体最终在WebSocket服务中通过ChatLogic实现了消息的存储和实时转发。同时通过对Conn的扩展使连接对象携带用户身份优化了业务逻辑的调用方式。至此IM 服务的私聊功能已具备基本可用性后续可在此基础上继续扩展群聊、已读回执、离线消息拉取等高级功能。