Go-Zero 项目开发22:用户群聊功能的实现与完善

发布时间:2026/7/26 13:41:56
Go-Zero 项目开发22:用户群聊功能的实现与完善 纲要消息存储模型基于读扩散一条消息只存一份通过type字段区分私聊/群聊receiver_id在群聊时指向群ID。会话管理用户创建群或加入群时由im服务创建群会话并维护用户与群的会话关系。消息推送与并发优化利用go-zero内置的线程工具实现群消息的并发发送避免因群成员数量大导致的延迟。消息队列处理在taskMQ中增加群聊分支调用社交服务获取群成员列表完成消息扩散与落地。服务协作社交API服务在创建群、申请进群、处理群申请等成功回调中通过RPC调用im服务建立会话。涉及技术栈go-zero、go-zero/core/threading、WebSocket、Redis、MySQL、RPC。消息存储与扩散模型群聊消息采用读扩散方案所有群成员共享同一条消息记录避免为每个用户存储一份副本。与私聊相同消息记录在同一张chat_log表中通过两个字段区分场景type消息类型枚举值为private私聊和group群聊。receiver_id接收者 ID私聊时为对方的用户 ID群聊时替换为群 ID。这样客户端拉取群历史消息时只需按群 ID 和消息类型查询即可获得完整的群聊记录无需在写路径上为每个成员维护独立的收件箱。会话的建立与管理创建时机会话的触发来源于两个入口创建群创建者发起创建群操作后社交服务需要同时为群本身和创建者与群之间建立会话。加入群新成员通过申请并被批准后社交服务需要为该用户与群建立会话。无论在哪个入口最终都通过im服务提供的RPC接口完成会话的初始化。时序梳理数据库IM RPC社交 RPC社交 API客户端数据库IM RPC社交 RPC社交 API客户端alt[会话不存在][会话已存在]创建群/审批加入执行群业务逻辑返回群 IDCreateGroupConversation(groupId, userId)查询群会话是否已存在插入群会话记录为用户插入群会话关系成功直接返回操作完成项目结构速览apps/ ├─ social/ │ ├─ api/ # 社交 API 服务 │ │ ├─ internal/ │ │ │ ├─ config/ │ │ │ ├─ logic/ # 创建群、申请群、处理申请等逻辑 │ │ │ └─ svc/ │ │ └─ social.api │ └─ rpc/ # 社交 RPC 服务 │ ├─ internal/ │ │ ├─ logic/ # GetGroupUserList 等 │ │ └─ svc/ │ └─ social.proto └─ im/ └─ rpc/ # IM RPC 服务 ├─ internal/ │ ├─ config/ │ ├─ logic/ # CreateGroupConversation 等 │ ├─ mq/ # taskMQ 消费者 │ ├─ server/ # WebSocket 连接管理、并发推送 │ └─ svc/ ├─ model/ # 会话、用户会话模型 └─ im.proto代码实现IM 服务中的会话逻辑以下代码位于 im 的 RPC 服务中负责创建群会话并关联用户会话列表。文件internal/logic/creategroupconversationlogic.gopackagelogicimport(contextdatabase/sqlgithub.com/pkg/errorsgo-zero-shop/apps/im/rpc/internal/svcgo-zero-shop/apps/im/rpc/pbgithub.com/zeromicro/go-zero/core/logx)typeCreateGroupConversationLogicstruct{ctx context.Context svcCtx*svc.ServiceContext logx.Logger}funcNewCreateGroupConversationLogic(ctx context.Context,svcCtx*svc.ServiceContext)*CreateGroupConversationLogic{returnCreateGroupConversationLogic{ctx:ctx,svcCtx:svcCtx,Logger:logx.WithContext(ctx),}}// CreateGroupConversation 创建群会话func(l*CreateGroupConversationLogic)CreateGroupConversation(in*pb.CreateGroupConversationReq)(*pb.CreateGroupConversationResp,error){// 1. 检查群会话是否已存在existing,err:l.svcCtx.ConversationModel.FindOneByConversationId(l.ctx,in.GroupId)iferr!nil!errors.Is(err,sql.ErrNoRows){l.Logger.Errorf(查询群会话失败: %v,err)returnnil,errors.Wrap(err,查询会话失败)}ifexisting!nil{returnpb.CreateGroupConversationResp{},nil}// 2. 创建群会话groupConv:model.Conversation{ConversationId:in.GroupId,Type:constant.ChatTypeGroup,}if_,err:l.svcCtx.ConversationModel.Insert(l.ctx,groupConv);err!nil{l.Logger.Errorf(创建群会话失败: %v,err)returnnil,errors.Wrap(err,创建会话失败)}// 3. 为创建者添加群会话关系userConv:model.UserConversation{UserId:in.CreatorId,ConversationId:in.GroupId,Type:constant.ChatTypeGroup,}if_,err:l.svcCtx.UserConversationModel.Insert(l.ctx,userConv);err!nil{l.Logger.Errorf(为用户添加群会话失败: %v,err)returnnil,errors.Wrap(err,添加用户会话失败)}returnpb.CreateGroupConversationResp{},nil}说明代码中ConversationModel和UserConversationModel为 go-zero 生成的 model 层对象constant.ChatTypeGroup是定义在常量包中的枚举值。并发推送消息群聊消息需要推送给所有在线成员如果采用串行方式逐个发送延迟会随着人数线性增长。为此我们引入go-zero提供的线程工具进行并发控制。并发限制与配置在im服务的Server结构体中通过Option模式暴露并发度参数方便运维调整。// internal/config/config.gotypeConfigstruct{// ... 其他配置ConcurrencyLimitintjson:ConcurrencyLimit}// internal/server/option.gotypeOptionstruct{ConcurrencyLimitint}funcWithConcurrencyLimit(limitint)Option{returnfunc(s*Server){s.concurrencyLimitlimit}}消息发送逻辑重构推送方法原先只处理私聊现在通过类型判定的方式分流群聊部分使用TaskRunner并发调用私聊推送方法。// internal/server/message.gopackageserverimport(contextfmtgo-zero-shop/apps/im/rpc/internal/constantgo-zero-shop/apps/im/rpc/internal/svcgo-zero-shop/apps/im/rpc/pbgithub.com/zeromicro/go-zero/core/threading)typeMessageCenterstruct{svcCtx*svc.ServiceContext concurrencyLimitinttaskRunner*threading.TaskRunner}funcNewMessageCenter(svcCtx*svc.ServiceContext,limitint)*MessageCenter{returnMessageCenter{svcCtx:svcCtx,concurrencyLimit:limit,taskRunner:threading.NewTaskRunner(limit),}}// Push 消息推送入口func(m*MessageCenter)Push(ctx context.Context,msg*pb.ChatMessage)error{switchmsg.Type{caseconstant.ChatTypePrivate:returnm.pushPrivate(ctx,msg,msg.ReceiverId)caseconstant.ChatTypeGroup:returnm.pushGroup(ctx,msg)default:returnfmt.Errorf(不支持的消息类型: %d,msg.Type)}}// pushPrivate 私聊推送func(m*MessageCenter)pushPrivate(ctx context.Context,msg*pb.ChatMessage,receiverIdstring)error{conn,err:m.svcCtx.ConnectionManager.Get(receiverId)iferr!nil{// 用户离线可记录日志或丢弃returnnil}// 假设存在 packResponse 将消息序列化为 WebSocket 帧data,err:packResponse(msg)iferr!nil{returnerr}returnconn.WriteMessage(data)}// pushGroup 群聊推送func(m*MessageCenter)pushGroup(ctx context.Context,msg*pb.ChatMessage)error{// msg.Receivers 由上游填充包含剔除发送者后的所有成员 IDfor_,uid:rangemsg.Receivers{uid:uid// 防止闭包引用问题m.taskRunner.Schedule(func(){iferr:m.pushPrivate(ctx,msg,uid);err!nil{logx.WithContext(ctx).Errorf(群聊推送失败, receiver%s, err%v,uid,err)}})}returnnil}注释ConnectionManager是我们实现的局部连接管理组件负责根据用户 ID 查找对应的WebSocket连接。TaskRunner.Schedule使用channel控制并发数当队列满时调用方会被阻塞从而实现反压。消息队列的群聊支持为了提高可靠性消息先被投递到消息队列由taskMQ异步消费并完成持久化与推送。需要在消费端增加群聊类型的处理并通过社交RPC服务获取群成员列表。消费端骨架// internal/mq/task.gopackagemqimport(contextencoding/jsongo-zero-shop/apps/im/rpc/internal/constantgo-zero-shop/apps/im/rpc/internal/svcgo-zero-shop/apps/im/rpc/pbgithub.com/zeromicro/go-zero/core/logx)typeTaskHandlerstruct{svcCtx*svc.ServiceContext pushService*server.MessageCenter}func(h*TaskHandler)Handle(ctx context.Context,raw[]byte)error{varmsg pb.ChatMessageiferr:json.Unmarshal(raw,msg);err!nil{returnerr}switchmsg.Type{caseconstant.ChatTypePrivate:returnh.handlePrivate(ctx,msg)caseconstant.ChatTypeGroup:returnh.handleGroup(ctx,msg)default:returnnil}}func(h*TaskHandler)handlePrivate(ctx context.Context,msg*pb.ChatMessage)error{// 存储消息记录...returnh.pushService.Push(ctx,msg)}func(h*TaskHandler)handleGroup(ctx context.Context,msg*pb.ChatMessage)error{// 1. 获取群成员rpcResp,err:h.svcCtx.SocialRpc.GroupUserList(ctx,social_pb.GroupUserListReq{GroupId:msg.ReceiverId,})iferr!nil{logx.WithContext(ctx).Errorf(获取群成员失败: %v,err)returnerr}// 2. 过滤发送者构建接收列表varreceivers[]stringfor_,user:rangerpcResp.Users{ifuser.UserId!msg.SenderId{receiversappend(receivers,user.UserId)}}msg.Receiversreceivers// 3. 存储消息记录...// 4. 并发推送returnh.pushService.Push(ctx,msg)}配置社交 RPC 客户端在im的config和service context中引入社交RPC客户端。// internal/config/config.gotypeConfigstruct{// ...SocialRpc zrpc.RpcClientConf}// internal/svc/servicecontext.gotypeServiceContextstruct{Config config.Config SocialRpc socialpb.SocialClient// ...其他依赖}funcNewServiceContext(c config.Config)*ServiceContext{returnServiceContext{Config:c,SocialRpc:socialpb.NewSocialClient(zrpc.MustNewClient(c.SocialRpc).Conn()),}}社交服务触发会话建立im服务的会话创建接口需要通过具体业务行为触发。在社交API服务中当创建群、申请入群、处理入群申请成功后应异步回调im RPC建立会话。社交 API 中的调用逻辑以创建群为例其余两个场景类似。// internal/logic/creategrouplogic.go (社交 API)func(l*CreateGroupLogic)CreateGroup(req*types.CreateGroupReq)(*types.CreateGroupResp,error){// ... 创建群业务逻辑获得 groupIdgroupId:xxx// 调用 IM RPC 创建群会话_,err:l.svcCtx.ImRpc.CreateGroupConversation(l.ctx,im_pb.CreateGroupConversationReq{GroupId:groupId,CreatorId:req.CreatorId,})iferr!nil{l.Logger.Errorf(创建群会话失败, groupId%s, err%v,groupId,err)// 通常这里可容忍失败通过定时任务补偿}returntypes.CreateGroupResp{GroupId:groupId},nil}社交服务的 IM RPC 配置// internal/config/config.go (社交 API)typeConfigstruct{// ...ImRpc zrpc.RpcClientConf}// internal/svc/servicecontext.go (社交 API)typeServiceContextstruct{Config config.Config ImRpc impb.ImClient// ...}funcNewServiceContext(c config.Config)*ServiceContext{returnServiceContext{Config:c,ImRpc:impb.NewImClient(zrpc.MustNewClient(c.ImRpc).Conn()),}}总结群聊功能的实现本质上复用了私聊的存储与推送链路核心差异体现在三处会话建模在群创建/加入时通过 im 服务统一管理群会话与用户‑会话关系。消息扩散服务端根据群 ID 查询成员列表借助go-zero的并发工具高效推送。异步处理消息队列消费端区分消息类型调用社交服务获取最新成员列表保证成员变动的实时性。整套方案在保持代码简洁的同时充分利用了go-zero框架的微服务能力RPC调用、线程池、消息队列可以平稳支撑较大规模的群组聊天场景。