OpenIM 源码拆解 第三篇
前两篇我们跟着一条消息走完了完整旅程——从发送端到接收端,经过网关、rpc、Kafka、msg-transfer、push 服务,七八站下来。链路篇里我提到过一句:消息最终写进 MongoDB 时,用的是"分块存储"——100 条消息塞一个文档。当时没展开,因为翻到这段源码的时候,我自己都没想明白为什么。
第一反应是:这不是给自己找麻烦吗?读一条消息还得先把整个文档捞出来?
带着这个疑问啃完源码,我才明白:这不是拍脑袋的设计,而是针对 IM 消息访问模式做的一次教科书级优化。这篇就把它讲透。
一、先看它到底怎么存的:三层嵌套结构
OpenIM 的消息文档模型是三层嵌套的(pkg/common/storage/model/msg.go):
// 第一层:一个文档 = 一个消息块(bucket)
type MsgDocModel struct {
DocID string bson:"doc_id" // 文档ID
Msg []*MsgInfoModel bson:"msgs" // 定长数组,装 100 条
}
// 第二层:每个数组元素,是对"一条消息槽位"的封装
type MsgInfoModel struct {
Msg *MsgDataModel bson:"msg" // 消息本体
Revoke *RevokeModel bson:"revoke" // 撤回信息
DelList []string bson:"del_list" // 对哪些人删除了
IsRead bool bson:"is_read"
}
// 第三层:真正的消息内容
type MsgDataModel struct {
SendID string bson:"send_id"
RecvID string bson:"recv_id"
ContentType int32 bson:"content_type"
Content string bson:"content"
Seq int64 bson:"seq" // 会话内递增序号,全场的主角
SendTime int64 bson:"send_time"
// ……
}
三层各司其职,别嫌绕,它的分工很清晰:
MSG_DOC(块):存储和读取的最小单位,一个文档打包 100 条消息。索引只加在它身上。
MSG_INFO(槽位):一条消息的"元信息容器"。为什么不直接放消息本体?因为一条消息除了内容,还有撤回、删除、已读这些会变化的状态。把它们和消息本体拆开,撤回一条消息时只改 revoke 字段,不动 msg。
MSG_DATA(本体):消息内容,写进去基本不变。
一句话:MSG_DOC 管"存哪",MSG_INFO 管"状态",MSG_DATA 管"内容"。
看到这个结构的时候,我反应了一会儿才理解为什么不把消息本体直接塞进数组——后来意识到,一条消息除了内容,还有撤回、删除、已读这些会变的状态。如果把状态和内容混在一起,每次撤回消息都要改整个文档,还可能影响其他消息。拆开之后,撤回一条消息只改 revoke 字段,不动 msg,干净利落。
你可能会问:100 条消息塞一个文档,读一条不是得把整个文档捞出来吗?这个疑问先放着,等讲完读取逻辑你就明白。
二、doc_id 怎么算?100 条一个文档是怎么落位的
这是整个设计的数学核心,就三个公式(model/msg.go):
const singleGocMsgNum = 100 // 每个文档装 100 条
// 1. 这条消息该进哪个文档?
func (m *MsgDocModel) GetDocID(conversationID string, seq int64) string {
seqSuffix := (seq - 1) / singleGocMsgNum // 块号 = (seq-1)/100
return conversationID + ":" + strconv.FormatInt(seqSuffix, 10)
}
// 2. 这条消息在文档数组的第几个槽位?
func (m *MsgDocModel) GetMsgIndex(seq int64) int64 {
return (seq - 1) % singleGocMsgNum // 槽位 = (seq-1)%100
}
拿会话 si_A_B 举个例子:
• seq=1 → docID = si_A_B:0,槽位 0
• seq=100 → docID = si_A_B:0,槽位 99
• seq=101 → docID = si_A_B:1,槽位 0
• seq=250 → docID = si_A_B:2,槽位 49
看出规律了吗?同一个会话,每 100 个连续 seq 落在同一个文档里,seq 直接决定了它在数组的第几格。这里的 seq 是会话内单调递增的序号(在消息入库前由 Redis 分配),它既是排序依据,也是寻址坐标——记住这点,后面读取全靠它。
三、哪些消息才会落库?不是一个开关,是"类别 × 留存"的组合判定
很多人(包括我最初)以为落库就靠一个 IsHistory 开关,翻源码才发现是两个维度在决策(categorizeMessageLists):
第一维 IsNotNotification——先分两条线:普通聊天消息走一条线,系统通知(入群、撤回等)走另一条线,分别由 handleMsg 和 handleNotification 处理。
第二维 IsHistory——每条线内部再判是否留存:true 进"待落库"列表,false 只走 Redis + 推送、绝不进 Mongo。
两维交叉出四象限,另外还有个 IsSendMsg:某些通知既要作为通知下发,又要在会话里当一条消息展示,于是会额外克隆一份进消息存储线。
• 落库:普通文字/图片/语音/文件(普通消息 + History),以及要留痕的通知(通知 + History)
• 不落库:正在输入、音视频信令、已读回执(History=false 的瞬时消息)
所以准确地说:IsHistory 是"要不要留存"的直接开关,但"存到哪条线、要不要额外存一份"是 IsNotNotification 和 IsSendMsg 共同决定的。这套组合拳让不同性质的消息各走各的路——这才是它设计精细的地方。
四、写入的精妙之处:先 update,再 insert
100 条共用一个文档,那写第 50 条消息时,文档可能已经存在(前 49 条建的),也可能不存在(本块第一条)。我一开始以为逻辑很简单——先 update 试试,不行再 insert 嘛。翻到 BatchInsertBlock 的源码才发现,它比我想的要精细:
tryUpdate := true
for i := 0; i < len(fields); i++ {
seq := firstSeq + int64(i)
if tryUpdate {
matched, _ := updateMsgModel(seq, i) // 先尝试 update 对应槽位
if matched {
continue // 文档已存在,填好槽位,跳过
}
}
// 走到这,说明文档不存在 → 新建,并预分配 100 个空槽
doc := model.MsgDocModel{
DocID: db.msgTable.GetDocID(conversationID, seq),
Msg: make([]*model.MsgInfoModel, num), // num = 100
}
// ……把当前块的消息填进对应槽位……
if err := db.msgDocDatabase.Create(ctx, &doc); err != nil {
if mongo.IsDuplicateKeyError(err) { // 撞车了?说明别人刚建好
i--
tryUpdate = true // 退回 update 模式
continue
}
}
tryUpdate = false // 新块建成,后续优先 insert
}
拆开看这套"先 update 后 insert"的组合拳:
1. 优先 update:假设文档已存在,直接把消息填进 msgs 数组的对应槽位。这是最常见的情况(一个块要填 100 次,99 次都是 update)。
2. update 没命中 → insert:说明是这个块的第一条消息,新建文档,并一次性预分配 100 个空槽(make([]*MsgInfoModel, 100))。
3. insert 撞唯一键 → 退回 update:并发场景下可能有人抢先建好了文档,靠 duplicate key error 兜底,退回去做 update。
这里有个容易忽略的细节:新建文档时就把 100 个槽位占好了。所以一个块从第 1 条到第 100 条,物理上只有 1 次 insert + 99 次 update,不会反复创建文档。
五、读取历史消息:按 seq 定位块,再定位槽位
写得巧,读起来才爽。客户端翻历史消息时,传的是一批 seq,OpenIM 的读取逻辑(msg.go 的 GetDocIDSeqsMap):
// 把要读的一批 seq,按所属文档分组
func (m *MsgDocModel) GetDocIDSeqsMap(conversationID string, seqs []int64) map[string][]int64 {
t := make(map[string][]int64)
for _, seq := range seqs {
docID := m.GetDocID(conversationID, seq) // 同一个块的 seq 归到一起
t[docID] = append(t[docID], seq)
}
return t
}
于是读取变成:
• 客户端要读 seq 1~100 → 全落在 si_A_B:0 一个文档里 → 一次查询就够了。
• 要读 seq 90~110 → 分到 si_A_B:0 和 si_A_B:1 两个文档 → 只需 2 次查询。
拿到文档后,用 GetMsgIndex(seq) 直接算出槽位下标,从数组里取出来,不用遍历。
对比一下"一条一个文档"的方案:读 100 条消息 = 100 次文档查询(或一个 $in 100 个 id 的大查询 + 100 次索引命中)。而分块方案:读 100 条 = 1 次文档查询。
同样是读这 21 条消息,一条一个文档要查 21 次,分块只查 2 次——这就是按块分组的威力。
六、核心观点:这不是"省空间",是对访问模式的定向优化
聊到这,很多人会觉得"哦,就是为了省空间嘛,合并存储减少文档数"。这只是表象。分块存储真正解决的核心问题只有一个:让存储布局对齐访问模式。具体来说,体现在三个层面:
1. 大幅降低索引压力
一条一个文档,1 亿条消息就是 1 亿个文档、1 亿条索引项。分块后除以 100,文档数和索引项都降到原来的 1/100——降了两个数量级。索引小了,B+ 树层级浅,查询和写入都更快,内存也扛得住。
2. 批量读取效率碾压
IM 读消息从来不是"读某一条",而是"读某个会话最近 20 条 / 翻一页历史"。这是连续 seq 的范围读。分块存储让连续的消息物理上聚集在同一个文档,一次 IO 捞出一大片,完美契合这个模式。
3. 顺序写 + 顺序读的天然契合
IM 消息的访问模式非常鲜明:按会话顺序写(seq 递增),按会话顺序读(翻页)。分块存储用 conversationID:块号 做 docID,让同一会话的消息在存储上就是按块顺序排列的——写是往当前块尾部追加,读是按块区间扫描。存储布局和访问模式对齐了,性能自然就上来了。
一句话总结这个设计哲学:不要用通用的"一行一记录"去硬套 IM 场景,而要顺着"数据怎么被访问"去设计"数据怎么存"。100 条塞一个文档,塞的不是数据,是对 IM 访问模式的深刻理解。
我翻完这段源码最大的感受是:很多"反直觉"的设计,不是作者在炫技,而是他比你更早看到了数据的访问模式。你以为是麻烦,其实是优化。
写在最后:你能拿走什么
OpenIM 用的这种分块存储方案,思想能迁移到很多"按某个维度顺序写入、又按同一维度范围读取"的场景:
• 物联网设备的时序数据(按设备+时间分块)
• 日志/审计流水(按来源+时间分块)
• Feed 流、评论列表(按主题+序号分块)
下次你要设计一张"会无限增长、且总是范围查询"的表时,别急着一行一记录。先问自己一句:我的数据是怎么被读的?想清楚这个,再决定怎么存——这才是 OpenIM 这段源码真正想教给我们的东西。
这篇聊的是"消息在 MongoDB 里怎么存",但你可能已经注意到了,文中反复出现的 seq,它的故事远不止"寻址坐标"这一个角色。
seq 怎么驱动多端同步?未读数为什么不存计数器而是用 seq 差值?Redis 缓存里的消息和 MongoDB 里的消息是什么关系?
这些才是消息存储设计里更深的水。下篇继续挖。
关注我,跟着源码学 IM。也欢迎在评论区聊聊——你做过的 IM 系统,消息存储是怎么设计的?
夜雨聆风