乐于分享
好东西不私藏

OpenIM 的消息为什么不一条存一个文档?我翻了源码才懂这个设计

OpenIM 的消息为什么不一条存一个文档?我翻了源码才懂这个设计

 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 条消息塞一个文档,读一条不是得把整个文档捞出来吗?这个疑问先放着,等讲完读取逻辑你就明白。

第一层 · 块 bucket+string DocID  "conversationID:块号"+MsgInfoModel[] Msg  "定长数组 · 100 个槽位"职责:存储/读取的最小单位1 *-- 100  内嵌数组 msgs[]第二层 · 槽位 slot+MsgDataModel Msg  "消息本体"+RevokeModel Revoke  "撤回信息"+string[] DelList  "对哪些人已删除"+bool IsRead职责:管理可变状态(撤回/删除/已读)1 *-- 1  内嵌 msg第三层 · 本体 data+int64 Seq  "会话内递增序号"+string SendID, RecvID, Content+int32 ContentType  +int64 SendTime职责:消息内容(写入后基本不变)RevokeModel+int32 Role+string UserID+int64 Time0..1 内嵌 revoke

二、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 分配),它既是排序依据,也是寻址坐标——记住这点,后面读取全靠它。

消息 seq = 250会话 = si_A_B寻址计算块号 = (seq-1) / 100 = 249 / 100 = 2docID = si_A_B:2槽位 = (seq-1) % 100 = 249 % 100 = 49文档 si_A_B:0seq 1~100Msg[0..99]文档 si_A_B:1seq 101~200Msg[0..99]文档 si_A_B:2seq 201~300Msg[0..99]命中Msg[49] = 这条消息第49格MongoDB 存储布局 (同一会话)

三、哪些消息才会落库?不是一个开关,是"类别 × 留存"的组合判定

很多人(包括我最初)以为落库就靠一个 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,不会反复创建文档。

开始批量写入 tryUpdate = true还有待写消息?写入完成取下一条 seq = firstSeq + itryUpdate == true?否·上一块刚新建尝试 UpdateMsg定位 docID + 槽位index,更新已有文档MatchedCount > 0?文档已存在?是·命中填好槽位 continue否·文档不存在新建文档 MsgDocModel预分配 100 个空槽 make Msg, 100Create 插入新文档DuplicateKeyError?并发下被人抢先建了?是·撞键i-- 回退tryUpdate=true成功tryUpdate=false 后续优先insert其他错误

五、读取历史消息:按 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 次文档查询。

客户端请求seq 90 ~ 110 共 21 条GetDocIDSeqsMap按 docID 分组seq 90~100seq 101~110查文档 si_A_B:0取 Msg[89..99]查文档 si_A_B:1取 Msg[0..9]合并 · 按 seq 排序

同样是读这 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 系统,消息存储是怎么设计的?