当前位置: 首页 > news >正文

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

纲要

  • 消息存储模型:基于读扩散,一条消息只存一份,通过type字段区分私聊/群聊,receiver_id在群聊时指向群ID
  • 会话管理:用户创建群或加入群时,由im服务创建群会话,并维护用户与群的会话关系。
  • 消息推送与并发优化:利用go-zero内置的线程工具实现群消息的并发发送,避免因群成员数量大导致的延迟。
  • 消息队列处理:在taskMQ中增加群聊分支,调用社交服务获取群成员列表,完成消息扩散与落地。
  • 服务协作:社交API服务在创建群、申请进群、处理群申请等成功回调中,通过RPC调用im服务建立会话。
  • 涉及技术栈go-zerogo-zero/core/threadingWebSocketRedisMySQLRPC

消息存储与扩散模型

群聊消息采用读扩散方案:所有群成员共享同一条消息记录,避免为每个用户存储一份副本。与私聊相同,消息记录在同一张chat_log表中,通过两个字段区分场景:

  • type:消息类型,枚举值为private(私聊)和group(群聊)。
  • receiver_id:接收者 ID,私聊时为对方的用户 ID,群聊时替换为群 ID。

这样,客户端拉取群历史消息时,只需按群 ID 和消息类型查询即可获得完整的群聊记录,无需在写路径上为每个成员维护独立的收件箱。

会话的建立与管理

创建时机

会话的触发来源于两个入口:

  1. 创建群:创建者发起创建群操作后,社交服务需要同时为群本身创建者与群之间建立会话。
  2. 加入群:新成员通过申请并被批准后,社交服务需要为该用户与群建立会话。

无论在哪个入口,最终都通过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.go

packagelogicimport("context""database/sql""github.com/pkg/errors""go-zero-shop/apps/im/rpc/internal/svc""go-zero-shop/apps/im/rpc/pb""github.com/zeromicro/go-zero/core/logx")typeCreateGroupConversationLogicstruct{ctx context.Context svcCtx*svc.ServiceContext logx.Logger}funcNewCreateGroupConversationLogic(ctx context.Context,svcCtx*svc.ServiceContext)*CreateGroupConversationLogic{return&CreateGroupConversationLogic{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{return&pb.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,"添加用户会话失败")}return&pb.CreateGroupConversationResp{},nil}

说明:代码中ConversationModelUserConversationModel为 go-zero 生成的 model 层对象;constant.ChatTypeGroup是定义在常量包中的枚举值。

并发推送消息

群聊消息需要推送给所有在线成员,如果采用串行方式逐个发送,延迟会随着人数线性增长。为此,我们引入go-zero提供的线程工具进行并发控制。

并发限制与配置

im服务的Server结构体中,通过Option模式暴露并发度参数,方便运维调整。

// internal/config/config.gotypeConfigstruct{// ... 其他配置ConcurrencyLimitint`json:"ConcurrencyLimit"`}
// internal/server/option.gotypeOptionstruct{ConcurrencyLimitint}funcWithConcurrencyLimit(limitint)Option{returnfunc(s*Server){s.concurrencyLimit=limit}}

消息发送逻辑重构

推送方法原先只处理私聊,现在通过类型判定的方式分流,群聊部分使用TaskRunner并发调用私聊推送方法。

// internal/server/message.gopackageserverimport("context""fmt""go-zero-shop/apps/im/rpc/internal/constant""go-zero-shop/apps/im/rpc/internal/svc""go-zero-shop/apps/im/rpc/pb""github.com/zeromicro/go-zero/core/threading")typeMessageCenterstruct{svcCtx*svc.ServiceContext concurrencyLimitinttaskRunner*threading.TaskRunner}funcNewMessageCenter(svcCtx*svc.ServiceContext,limitint)*MessageCenter{return&MessageCenter{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("context""encoding/json""go-zero-shop/apps/im/rpc/internal/constant""go-zero-shop/apps/im/rpc/internal/svc""go-zero-shop/apps/im/rpc/pb""github.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{receivers=append(receivers,user.UserId)}}msg.Receivers=receivers// 3. 存储消息记录...// 4. 并发推送returnh.pushService.Push(ctx,msg)}

配置社交 RPC 客户端

imconfigservice context中引入社交RPC客户端。

// internal/config/config.gotypeConfigstruct{// ...SocialRpc zrpc.RpcClientConf}
// internal/svc/servicecontext.gotypeServiceContextstruct{Config config.Config SocialRpc socialpb.SocialClient// ...其他依赖}funcNewServiceContext(c config.Config)*ServiceContext{return&ServiceContext{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)// 通常这里可容忍失败,通过定时任务补偿}return&types.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{return&ServiceContext{Config:c,ImRpc:impb.NewImClient(zrpc.MustNewClient(c.ImRpc).Conn()),}}

总结

群聊功能的实现本质上复用了私聊的存储与推送链路,核心差异体现在三处:

  • 会话建模:在群创建/加入时通过 im 服务统一管理群会话与用户‑会话关系。
  • 消息扩散:服务端根据群 ID 查询成员列表,借助go-zero的并发工具高效推送。
  • 异步处理:消息队列消费端区分消息类型,调用社交服务获取最新成员列表,保证成员变动的实时性。

整套方案在保持代码简洁的同时,充分利用了go-zero框架的微服务能力(RPC调用、线程池、消息队列),可以平稳支撑较大规模的群组聊天场景。

http://www.jsqmd.com/news/1267753/

相关文章:

  • 终极HTML转Figma指南:3步将任何网页变成可编辑设计稿
  • 机器学习核心算法实践指南:从回归到神经网络完整学习路径
  • 从源码到JAR:PDFFigures2本地化部署与依赖库配置完全手册
  • SCMP供应链管理证书考取流程 - 众智商学院官方
  • 2026上海GEO优化服务|企业做GEO服务商怎么选?本地靠谱选型指南与五家机构深度测评 - 企业新闻快传
  • 3步极速指南:用SlopeCraft将任意图片变成立体Minecraft地图艺术
  • QuickRecorder自动化录屏终极指南:5步打造你的专属录屏工作流
  • 基于Dify和RAG技术构建智能数据治理知识库
  • 终极掌控:GTA圣安地列斯存档编辑器完全指南
  • 5大创新:开源眼动追踪系统如何重新定义人机交互
  • Keyboard Chatter Blocker:精准修复机械键盘连击的终极免费方案
  • Nucleoid runtime工作原理解析:从语句追踪到知识图谱生成的全过程
  • TMS320C6670热阻解析与FCBGA封装散热设计实战
  • 无线MCU命令调度机制:射频时序的精准掌控与低功耗设计
  • 2026西安除甲醛公司实测调研:多维度筛选靠谱室内空气治理机构 - 西安治泉环保
  • LibreSpeed-cli vs 其他测速工具:为什么选择这款开源命令行神器?
  • MiniMax多模态AI技术与企业级应用解析
  • 终极指南:3步快速安装CZSC缠论可视化分析插件
  • 2026.7月衢州房屋漏水维修实用指南|厨卫/阳台/外墙/屋面/地下室一站式防水修缮参考 - 超人防水
  • 上海快消消费品牌企业做GEO服务商怎么选?2026五家代表性服务商深度测评与靠谱选型指南 - 子柔传媒
  • Python评分卡开发终极指南:从业务痛点到专业风控模型的完整路线图
  • M9A:重返未来1999智能自动化助手完整指南,彻底解放你的游戏时间
  • C++模板链接错误:undefined reference to的根源与解决方案
  • HarmonyOS ArkTS 实战:实现一个标准计算器
  • Streetmix安全机制详解:内容安全策略与用户数据保护措施
  • 5分钟解决Calibre中文路径乱码:让“科幻小说“不再变成“Ke_Huan_Xiao_Shuo“
  • Calibre中文路径保护插件:告别拼音乱码,重拾清晰文件管理
  • 溧阳装修公司对比:哪家好、口碑好、靠谱又高性价比?2026 - 装企精灵GEO
  • 终极Linux打印机驱动解决方案:foo2zjs让100+款打印机完美工作
  • 2026最新菏泽防水补漏本地人必选的正规靠谱公司推荐-房屋漏水检测维修师傅上门-卫生间厨房阳台房顶外墙漏水检测精准测漏 - 吉林同城获客