Go 企业 IM 架构:消息可靠投递和离线推送的工程方案
Go 企业 IM 架构消息可靠投递和离线推送的工程方案一、一条消息发出去对方说没收到企业 IM 最不能出错的就是消息可靠性。但现实中消息丢失的场景远比想象的多用户手机切后台时 WebSocket 断连消息还在服务端内存里服务重启时内存中的待推送消息全部丢失跨机房网络抖动导致消息路由到错误的节点。这些问题在外表看起来很偶发但在 10 万日活的场景下每天可能出现几百次。企业 IM 的消息可靠性不是每条消息都要保证 100% 送达那不现实而是消息不丢不重丢了也要能被发现和恢复。这是 IM 系统最难也最重要的设计目标。二、消息可靠投递的架构设计核心设计包括消息持久化优先、ACK 确认机制、离线消息补偿关键设计消息在推送之前必须已落盘先写 MySQL再推送。如果推送失败用户不在线消息已经在数据库中用户上线时可以拉取。这样即使消息服务崩溃重启也不丢消息。三、Go 实现消息服务核心逻辑package im import ( context database/sql encoding/json fmt sync time github.com/go-redis/redis/v8 ) // MessageStatus 消息状态 type MessageStatus int const ( StatusSent MessageStatus iota // 已发送已落盘 StatusDelivering // 推送中 StatusDelivered // 已送达 StatusRead // 已读 StatusFailed // 推送失败 ) // Message 消息结构 type Message struct { MsgID string json:msg_id // 全局唯一ID SeqID int64 json:seq_id // 用户递增序号 FromUserID string json:from_user_id ToUserID string json:to_user_id GroupID string json:group_id,omitempty Content string json:content ContentType string json:content_type // text, image, file Timestamp time.Time json:timestamp Status MessageStatus json:status } // IMService IM 核心服务 type IMService struct { db *sql.DB redis *redis.Client // 用户在线网关映射: userID - gatewayAddr onlineMap sync.Map // 递增序列号生成器 seqGen *SequenceGenerator } // NewIMService 创建 IM 服务 func NewIMService(db *sql.DB, redis *redis.Client) *IMService { return IMService{ db: db, redis: redis, seqGen: NewSequenceGenerator(redis), } } // SendMessage 发送消息核心流程 func (s *IMService) SendMessage( ctx context.Context, msg *Message, ) error { if msg.FromUserID || (msg.ToUserID msg.GroupID ) { return fmt.Errorf(消息发送方或接收方为空) } // 第 1 步生成消息ID和序号 msg.MsgID generateMsgID() var err error msg.SeqID, err s.seqGen.Next( ctx, msg.ToUserID, ) if err ! nil { return fmt.Errorf(生成序号失败: %w, err) } msg.Timestamp time.Now() msg.Status StatusSent // 第 2 步持久化到 MySQL同步保证不丢 if err : s.persistMessage(ctx, msg); err ! nil { return fmt.Errorf(消息持久化失败: %w, err) } // 第 3 步异步推送不阻塞发送方 go s.deliverMessage(context.Background(), msg) return nil } // persistMessage 持久化消息到 MySQL func (s *IMService) persistMessage( ctx context.Context, msg *Message, ) error { query : INSERT INTO messages (msg_id, seq_id, from_user_id, to_user_id, group_id, content, content_type, timestamp, status) VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?) _, err : s.db.ExecContext(ctx, query, msg.MsgID, msg.SeqID, msg.FromUserID, msg.ToUserID, msg.GroupID, msg.Content, msg.ContentType, msg.Timestamp, msg.Status, ) if err ! nil { return fmt.Errorf(数据库写入失败: %w, err) } return nil } // deliverMessage 推送消息给接收方 func (s *IMService) deliverMessage( ctx context.Context, msg *Message, ) { // 更新状态为推送中 s.updateMessageStatus(ctx, msg.MsgID, StatusDelivering) // 检查接收方是否在线 gatewayAddr, online : s.onlineMap.Load(msg.ToUserID) if !online { // 用户离线写入离线消息队列 s.storeOfflineMessage(ctx, msg) s.pushOfflineNotification(ctx, msg.ToUserID) return } // 用户在线推送到对应网关 delivered : s.pushToGateway( ctx, gatewayAddr.(string), msg, ) if delivered { s.updateMessageStatus(ctx, msg.MsgID, StatusDelivered) } else { // 推送失败标记为失败并加入离线队列 s.updateMessageStatus(ctx, msg.MsgID, StatusFailed) s.storeOfflineMessage(ctx, msg) } } // pushToGateway 推送消息到接入网关 func (s *IMService) pushToGateway( ctx context.Context, gatewayAddr string, msg *Message, ) bool { // 实际项目中通过 RPC 或消息队列推送到接入层 data, _ : json.Marshal(msg) err : s.redis.Publish( ctx, fmt.Sprintf(gateway:%s:push, gatewayAddr), data, ).Err() return err nil } // storeOfflineMessage 存储离线消息 func (s *IMService) storeOfflineMessage( ctx context.Context, msg *Message, ) { key : fmt.Sprintf(offline_msgs:%s, msg.ToUserID) msgJSON, _ : json.Marshal(msg) // ZSet: score 为 seq_idmember 为消息 JSON s.redis.ZAdd(ctx, key, redis.Z{ Score: float64(msg.SeqID), Member: msgJSON, }) // 设置 7 天过期 s.redis.Expire(ctx, key, 7*24*time.Hour) } // PullOfflineMessages 拉取离线消息用户上线时调用 func (s *IMService) PullOfflineMessages( ctx context.Context, userID string, limit int64, ) ([]*Message, error) { key : fmt.Sprintf(offline_msgs:%s, userID) // 按 seq_id 升序拉取 members, err : s.redis.ZRange(ctx, key, 0, limit-1).Result() if err ! nil { return nil, fmt.Errorf(拉取离线消息失败: %w, err) } var messages []*Message for _, member : range members { var msg Message if json.Unmarshal([]byte(member), msg) nil { messages append(messages, msg) } } // 从队列中移除已拉取的消息 if len(messages) 0 { s.redis.ZRemRangeByRank(ctx, key, 0, int64(len(messages)-1)) } return messages, nil } // pushOfflineNotification 发送离线推送通知APNs/FCM func (s *IMService) pushOfflineNotification( ctx context.Context, userID string, ) { // 实际项目中调用推送服务APNs iOS / FCM Android fmt.Printf(发送离线通知给用户 %s\n, userID) } // updateMessageStatus 更新消息状态 func (s *IMService) updateMessageStatus( ctx context.Context, msgID string, status MessageStatus, ) { query : UPDATE messages SET status ? WHERE msg_id ? s.db.ExecContext(ctx, query, status, msgID) } // SequenceGenerator 用户消息序号生成器 type SequenceGenerator struct { redis *redis.Client } func NewSequenceGenerator(redis *redis.Client) *SequenceGenerator { return SequenceGenerator{redis: redis} } func (sg *SequenceGenerator) Next( ctx context.Context, userID string, ) (int64, error) { key : fmt.Sprintf(seq:%s, userID) return sg.redis.Incr(ctx, key).Result() } func generateMsgID() string { // 雪花算法或 UUID return fmt.Sprintf(msg_%d, time.Now().UnixNano()) }四、边界分析与 Trade-offs同步写 DB 的性能代价每条消息同步写 MySQL 会增加 5-10ms 的延迟。在高并发场景如群发 5000 人的群下如果每条消息都同步写吞吐量可能不够。解决方案单聊消息同步写数据量小可靠性要求高群聊消息异步批量写数据量大可以容忍 100ms 内的延迟。离线消息的存储策略Redis ZSet 适合做离线消息的临时队列但内存有限。离线时间超过 7 天仍未被拉取的消息应该降级到 MySQL 或对象存储中用户点查看更早消息时再拉取。消息去重 vs 投递保证移动网络下客户端可能重复发送同一消息TCP 重传。服务端通过 msg_id 做幂等处理相同 msg_id 只处理一次配合客户端生成唯一 msg_id。但不要用消息内容哈希做去重——用户可能真的想发两条相同的消息。已读回执的批量更新每条消息都发一次已读回执写 DB 的压力太大。客户端应批量提交已读回执如每 500ms 或攒够 20 条服务端用UPDATE ... WHERE msg_id IN (...)批量更新。五、总结企业 IM 的消息可靠性核心原则先落盘后推送。数据库是消息的Source of Truth内存和 Redis 是加速层。离线消息用 Redis ZSet 做临时队列按 seq_id 排序、支持批量拉取、自动过期超过 7 天不拉取的消息降级到 MySQL。推送失败不能静默——要写入失败日志并触发重试还要通知发送方消息未送达到目标不是发送失败。消息的 seq_id 服务端生成Redis INCR避免客户端排序导致的乱序问题。

相关新闻

独立开发者的API monetization实战:从免费工具到付费API服务

独立开发者的API monetization实战:从免费工具到付费API服务

独立开发者的API monetization实战:从免费工具到付费API服务 API monetization的三个阶段 独立开发者的产品,如果有"数据"或"计算能力",就可以把这部分能力开放成API,变成"第二收入来源"。 阶段一&…

2026/7/24 16:45:55阅读更多 →
Python毕设项目: 基于Python的香港历史资料整理与科普传播系统实现 校园香港历史科普学习平台设计与实现(源码+文档,讲解、调试运行,定制等)

Python毕设项目: 基于Python的香港历史资料整理与科普传播系统实现 校园香港历史科普学习平台设计与实现(源码+文档,讲解、调试运行,定制等)

博主介绍:✌️码农一枚 ,专注于大学生项目实战开发、讲解和毕业🚢文撰写修改等。全栈领域优质创作者,博客之星、掘金/华为云/阿里云/InfoQ等平台优质作者、专注于Java、小程序技术领域和毕业项目实战 ✌️技术范围:&am…

2026/7/24 16:43:49阅读更多 →
GG3M技术:学术理论与产业实践的差异分析

GG3M技术:学术理论与产业实践的差异分析

1. 关于GG3M的学术与产业视角差异解析GG3M作为一项新兴技术范式,近年来在学术界和产业界引发了截然不同的讨论声浪。上周参加行业技术峰会时,我与几位高校教授和科技公司CTO的交流就充分体现了这种认知鸿沟——教授们更关注理论突破的可能性,…

2026/7/24 16:43:49阅读更多 →
HarmonyOS开发实战:小分享-HAR 与 HSP 共享包的发布与依赖

HarmonyOS开发实战:小分享-HAR 与 HSP 共享包的发布与依赖

HarmonyOS开发实战:小分享-HAR 与 HSP 共享包的发布与依赖 前言 HAR(静态共享包) 和 HSP(动态共享包) 是 HarmonyOS 模块化开发的核心技术。HAR 在编译时打包到 HAP 中,HSP 在运行时动态加载。小分享 App …

2026/7/24 18:18:14阅读更多 →
2026最新:职场人怎么做职场视频总结?这5个实用技巧亲测好用

2026最新:职场人怎么做职场视频总结?这5个实用技巧亲测好用

2026年职场人做职场视频总结,不需要全程手动剪辑逐字整理,核心逻辑是先用AI提取视频核心信息,再结构化整理成可复用的总结内容。我长期测试各类AI效率工具,亲测5个实用技巧覆盖从素材提取到协作导出全流程,对自媒体人做…

2026/7/24 18:18:14阅读更多 →
准备面试的研究生:2026年3款自己录音背书的app哪款更实用?

准备面试的研究生:2026年3款自己录音背书的app哪款更实用?

先按场景给答案 针对准备面试的研究生整理背书录音、内容创作者处理自主录音素材的需求,本次评测2026年主流的3款自己录音背书的app:网易见外工作台、听脑AI、迅捷录音转文字,无绝对排名,仅按场景匹配:仅需要免费基础…

2026/7/24 18:18:14阅读更多 →
AI技术如何重塑CMS建站行业格局与应用实践

AI技术如何重塑CMS建站行业格局与应用实践

1. AI技术如何重塑CMS建站行业格局过去三年里,我亲眼见证了AI技术从内容生产环节渗透到CMS系统的每个毛细血管。最直观的变化发生在内容创作模块——传统CMS需要手动填写的标题、摘要、标签,现在通过GPT-3.5级别的模型就能自动生成符合SEO规范的内容框架…

2026/7/24 18:18:14阅读更多 →
深度强化学习在工业控制中的多环境自适应实践

深度强化学习在工业控制中的多环境自适应实践

1. 项目背景与核心挑战在工业控制系统和机器人领域,我们经常遇到这样的困境:实际部署环境与训练环境存在差异,传感器数据可能缺失或部分不可观测。就像驾驶员在雾天行车时,视线受阻但依然需要安全驾驶。这种"部分可观测多环境…

2026/7/24 18:18:14阅读更多 →
2026年SEO优化公司选型指南:行业标准、适配场景与主流服务商分析

2026年SEO优化公司选型指南:行业标准、适配场景与主流服务商分析

近年来,国内企业线上获客体系不断成熟,自然搜索流量的长期价值逐步得到行业认可。根据中国广告协会发布的《2025年中国搜索营销发展白皮书》数据显示,国内企业自然搜索流量贡献的转化占比已提升至42%,超过付费搜索流量&#xff0c…

2026/7/24 18:16:14阅读更多 →
Go语言静态资源打包方案对比与实践指南

Go语言静态资源打包方案对比与实践指南

1. 项目背景与核心需求在Go语言开发中,我们经常需要处理静态资源文件的打包问题。无论是Web应用的模板文件、前端资源,还是配置文件、证书等,都需要随程序一起分发。传统做法是将这些文件与编译后的二进制文件放在同一目录下,但这…

2026/7/24 0:58:53阅读更多 →
Go语言实现高性能LDAP认证服务的架构与实践

Go语言实现高性能LDAP认证服务的架构与实践

1. 项目背景与核心价值LDAP(轻量级目录访问协议)作为企业级身份认证的黄金标准,已经服务了超过80%的财富500强公司。我在金融科技领域实施统一认证体系时,发现传统Java方案存在启动慢、内存占用高等痛点。而Go语言凭借其协程并发模…

2026/7/24 0:58:53阅读更多 →
【AI面试官实战指南】:用ChatGPT模拟10类高频技术岗面试,3天提升应答精准度92%

【AI面试官实战指南】:用ChatGPT模拟10类高频技术岗面试,3天提升应答精准度92%

更多请点击: https://intelliparadigm.com 第一章:AI面试官实战指南的核心价值与适用场景 AI面试官并非替代人类HR的“黑箱工具”,而是以可解释、可审计、可迭代的方式,赋能招聘全链路的关键基础设施。其核心价值在于将主观经验沉…

2026/7/24 0:58:53阅读更多 →
我的编程之路:第一篇博客

我的编程之路:第一篇博客

大家好,我是一名编程初学者,同时这也是我编程学习之路上的第一篇博客。在这里,我想要向大家介绍我的一些想法和规划。a.自我介绍我是一个刚刚接触编程的新手,目前在学习c语言,我对编程世界充满了强烈的好奇。当然&…

2026/7/24 0:00:06阅读更多 →
【LeetCode 54】螺旋矩阵

【LeetCode 54】螺旋矩阵

问题描述: 解法: 1、模拟(参考自【LeetCode 54】螺旋矩阵-CSDN博客) int *spiralOrder(int **matrix, int matrixSize, int *matrixColSize, int *returnSize) {static const int dirs[4][2] {{0, 1}, {1, 0}, {0, -1}, {-1, …

2026/7/24 0:00:06阅读更多 →
2026 WAIC:模型隐身、智能体疯野,厂商竞赛聚焦办公场景与商业闭环

2026 WAIC:模型隐身、智能体疯野,厂商竞赛聚焦办公场景与商业闭环

知春路不相信模型领先今年WAIC大会,昔日AI六小龙来了五家,分别是Kimi、阶跃星辰、Minimax、百川智能、零一万物。连放弃基模的百川和零一万物都来了,唯一缺席的竟是近几个月来风光无限的智谱。(DeepSeek一直不参加)WAI…

2026/7/24 0:00:06阅读更多 →
YOLOv8推理性能优化:从1.2FPS到35FPS的全链路加速实践

YOLOv8推理性能优化:从1.2FPS到35FPS的全链路加速实践

如果你在部署 YOLOv8 时,发现推理速度只有可怜的 1-2 FPS,而别人的演示视频却能跑到 30 FPS 以上,那么问题很可能不在模型本身,而在于你的整个处理链路。很多开发者拿到一个训练好的 YOLOv8 模型后,会直接使用官方示例…

2026/7/23 22:58:43阅读更多 →
Coze与Dify对比指南:低代码AI应用开发从入门到实战

Coze与Dify对比指南:低代码AI应用开发从入门到实战

1. 从零到一:为什么你需要了解 Coze 和 Dify?如果你对 AI 应用开发感兴趣,但一看到“大模型”、“智能体”、“工作流”这些词就头疼,觉得门槛太高,那这篇文章就是为你准备的。很多开发者,包括我自己&#…

2026/7/23 18:58:18阅读更多 →
AI生图工具怎么选?2026年6月版实测对比

AI生图工具怎么选?2026年6月版实测对比

做自媒体的朋友应该都有体会:配图一直是个让人头疼的问题。2026年,AI生图工具已经非常成熟了,但工具太多反而不知道怎么选。以下是截至2026年6月我对主流AI生图工具的实测对比。Midjourney V8.1:速度之王2026年6月11日&#xff0c…

2026/7/23 18:58:18阅读更多 →