场景设计:类微信 IM 系统如何保证消息顺序一致性
本文讨论通用即时通信系统设计,不代表微信真实内部实现。
1. “保证消息顺序”首先要问:保证哪一种顺序
面试中最容易犯的错误,是直接回答“使用 Kafka 单分区”或“用消息时间戳排序”。
顺序有多个层级:
| 顺序范围 | 示例 | 是否通常需要 |
|---|---|---|
| 单客户端发送顺序 | 用户 A 连续发送 M1、M2 | 需要 |
| 单会话顺序 | A、B 在同一聊天中的所有消息 | 通常需要 |
| 单用户收件箱顺序 | 用户所有会话的消息全序 | 一般不需要严格全序 |
| 全系统消息顺序 | 所有用户所有消息统一排序 | 几乎不需要,代价极高 |
| 因果顺序 | 回复消息不能先于被回复消息出现 | 需要通过引用和补拉处理 |
类微信 IM 更合理的目标是:
1 | |
不必为了一个会话内的顺序,构造全世界唯一的中心序列器。
2. 为什么客户端时间戳不可靠
用 sendTime 排序存在以下问题:
- 用户可以修改设备时间;
- 不同设备时钟有漂移;
- 移动网络延迟差异很大;
- 消息离线缓存后延迟发送;
- 服务端重试和跨机房转发会改变到达顺序;
- 同一毫秒可产生多条消息;
- 多端同时发送无法形成稳定全序。
客户端时间可以作为展示信息或排查字段,但不应作为最终会话顺序依据。
3. 需求拆解
3.1 功能需求
- 单聊、群聊;
- 文本、图片、语音、文件、系统消息;
- 在线实时推送;
- 离线消息同步;
- 多设备登录;
- 发送重试和幂等;
- 已送达、已读状态;
- 撤回、编辑、删除;
- 历史消息分页;
- 断线重连后补齐缺失消息。
3.2 顺序相关非功能需求
- 同会话内最终顺序一致;
- 网络抖动不会产生重复消息;
- 消息暂时乱序到达时客户端可重排;
- 检测消息缺口并主动补拉;
- 服务端横向扩容后顺序不失效;
- 单个热点群不能拖垮整个集群;
- 不要求不必要的全局强一致。
4. 总体架构
flowchart LR
A[发送端 App] --> CONN[长连接网关]
CONN --> ROUTER[消息路由服务]
ROUTER --> SHARD[会话分片/Leader]
SHARD --> SEQ[会话序列分配器]
SHARD --> MSGDB[(消息存储)]
SHARD --> OUTBOX[(Outbox/事件日志)]
OUTBOX --> MQ[消息队列]
MQ --> FANOUT[收件箱/推送分发]
FANOUT --> ONLINE[在线连接路由]
FANOUT --> INBOX[(用户同步游标/收件箱)]
ONLINE --> B[接收端 App]
B --> SYNC[消息同步服务]
SYNC --> MSGDB
SYNC --> INBOX
核心原则:
1 | |
这个有序域可以是:
- 固定分片的单 Leader;
- Kafka 同一分区;
- Actor/Virtual Actor;
- 数据库分片上的串行序列分配器;
- Raft Group 中的状态机。
5. 消息标识要分层
建议至少有三类 ID。
5.1 client_msg_id
由客户端生成,用于发送幂等。
1 | |
同一条消息因超时重试时,client_msg_id 不变。
5.2 message_id
服务端全局唯一消息 ID,用于定位、存储、引用、撤回和审计。
它可以趋势递增,但不必承担严格会话顺序。
5.3 msg_seq
会话内单调递增序列:
1 | |
客户端最终按 msg_seq 排序。
示例:
1 | |
6. 会话分片
6.1 按 conversation_id 路由
1 | |
同一会话的消息进入同一分片,分片内部有一个明确 Leader 负责写入。
flowchart TD
M1[会话 C1 的消息] --> H{hash conversationId}
M2[会话 C2 的消息] --> H
M3[会话 C1 的另一条消息] --> H
H --> S1[Shard 1 Leader]
H --> S2[Shard 2 Leader]
H --> S3[Shard 3 Leader]
S1 --> Q1[有序日志]
S2 --> Q2[有序日志]
S3 --> Q3[有序日志]
C1 的所有消息始终进入同一有序域,C1 与 C2 可以并行处理。
6.2 为什么不能随机负载均衡
若 A 的 M1 请求进入节点 X,M2 进入节点 Y:
- X 和 Y 的处理速度不同;
- M2 可能先落库;
- 两个节点可能各自分配相同序列;
- 跨节点协调会变得复杂。
入口网关可以随机接收连接,但消息路由层必须按会话键确定写入分片。
7. 服务端序列号怎么生成
7.1 方案一:数据库行锁/条件更新
维护会话序列表:
1 | |
在事务中:
1 | |
然后读取新值或使用数据库 RETURNING 能力。
优点:实现直观、强一致;缺点:热点群会集中竞争同一行。
7.2 方案二:Redis INCR
1 | |
优点:快、原子;缺点:
- Redis 故障切换时如何保证序列不回退;
- 序列分配成功但消息落库失败会产生空洞;
- Redis 与消息库一致性复杂;
- 不能简单把 Redis 当不可丢失的唯一序列事实源。
序列空洞通常是可接受的,回退和重复不可接受。
7.3 方案三:分段号段
每个会话分片 Leader 一次从数据库申请一段:
1 | |
之后在内存中分配,耗尽再申请。
优点:减少中心存储访问;缺点:节点宕机会产生较大空洞。
对 IM 而言,空洞通常可以接受,只要:
1 | |
但号段如果被多个并发 Leader 同时持有,就会出问题。因此会话分片本身需要单 Leader 或 fencing token。
7.4 方案四:分片有序日志的 Offset
将同一 conversationId 固定写入同一个 Kafka Partition,使用分区 offset 作为有序位置或辅助序列。
Kafka 只保证分区内部顺序,不保证跨分区全局顺序。这个特性正好匹配“同会话有序、不同会话并行”。
但直接把 Kafka offset 当 msg_seq 也有问题:
- 同一分区包含多个会话,offset 在各会话内有间隔;
- 分区迁移和消息存储映射要处理;
- 业务查询通常仍希望有会话级序列;
- 如果生产消息前后还有数据库事务,顺序边界需统一。
更常见做法是:
1 | |
8. 推荐写入流程
sequenceDiagram
autonumber
participant App as 发送端
participant Gateway as 长连接网关
participant Router as 会话路由
participant Leader as 会话分片 Leader
participant DB as 消息数据库
participant MQ as 有序事件流
participant Receiver as 接收端
App->>Gateway: Send(clientMsgId, conversationId, content)
Gateway->>Router: 鉴权后的消息命令
Router->>Leader: 按 conversationId 路由
Leader->>DB: 查询 clientMsgId 是否已存在
alt 已处理过
DB-->>Leader: 返回原 messageId/msgSeq
Leader-->>App: 原成功结果
else 新消息
Leader->>Leader: 分配 msgSeq
Leader->>DB: 事务写 message + sender ack + outbox
DB-->>Leader: 提交成功
Leader-->>App: ACK(messageId, msgSeq)
DB->>MQ: CDC/Outbox 发布
MQ->>Receiver: 推送 message(msgSeq)
end
关键顺序:
1 | |
如果先 ACK 再落库,服务宕机会出现“发送端认为成功,但服务端没有消息”。
9. 数据模型
9.1 消息表
1 | |
主查询路径是:
1 | |
9.2 用户会话同步状态
1 | |
9.3 客户端本地状态
客户端对每个会话维护:
1 | |
10. 消息乱序到达怎么办
即使服务端写入有序,网络推送仍可能乱序:
- 消息经过不同网关;
- 重连前后的连接并存;
- 推送重试;
- 大消息上传较慢;
- 多设备同步;
- 某条消息消费延迟。
假设客户端当前连续收到到 seq=100,随后先收到 102:
1 | |
sequenceDiagram
participant Server as 服务端
participant Client as 客户端
Server-->>Client: message seq=100
Server-->>Client: message seq=102
Client->>Client: 检测到缺口 101
Client->>Server: sync(conversationId, afterSeq=100)
Server-->>Client: seq=101, seq=102
Client->>Client: 按 seq 去重并重排
客户端不能因为先收到 102 就永久把它显示在 101 前面。
11. 空洞和缺口不是同一个概念
11.1 序列分配空洞
服务端分配了 seq=101,但事务失败,最终没有 101。这是号段或先分配后持久化常见现象。
11.2 传输缺口
服务端确实存在 101,但客户端没有收到。
客户端仅凭看到 100 和 102 无法立即知道是哪一种。
解决方案有两种:
- 序列分配与消息提交在同一状态机中,尽量不产生空洞;
- 同步接口返回权威信息,例如消息列表和
committedMaxSeq,必要时返回跳号标记。
很多系统允许服务端序列存在空洞,但客户端同步协议必须能前进,不能永远等待不存在的 101。
可以提供:
1 | |
或者使用连续日志位置作为同步游标,而业务 msgSeq 只负责排序。
12. 客户端发送顺序
用户在同一设备连续发送 M1、M2,客户端可以维护 client_seq:
1 | |
如果 M1 上传图片较慢,M2 文本先到服务端,产品需要选择:
方案 A:按服务端收到顺序
简单、延迟低,但可能与用户点击顺序不同。
方案 B:同设备按 clientSeq 串行发送
客户端等 M1 获得 ACK 后再发送 M2,能保持本机点击顺序,但大文件会阻塞后续文本。
方案 C:媒体先预上传,消息提交轻量化
先把图片上传到对象存储,获得 objectKey,再发送消息元数据。这样消息提交本身很快,可以减少顺序反转。
生产系统通常将大文件上传与聊天消息提交分开,并允许产品定义“发送动作完成”的时点。
13. 多设备同时发送
同一账号在手机和电脑同时发送时,没有天然的客户端全序。
正确做法:
- 两端各自生成
client_msg_id; - 所有消息进入同一会话分片;
- 服务端按该分片接收/提交顺序分配
msg_seq; - 所有设备最终按服务端
msg_seq对齐。
不要尝试依靠手机和电脑的本地时间排序。
14. Kafka 如何使用才不会破坏顺序
14.1 Partition Key
生产事件时使用:
1 | |
同一会话进入同一 Partition。
14.2 消费并发
一个 Partition 在同一 Consumer Group 内同一时刻只由一个消费者实例处理,但消费者内部如果把消息扔进无序线程池,仍可能完成乱序。
错误示例:
1 | |
修复方式:
- 每分区串行消费;
- 按 conversationId 再次做 keyed executor;
- 异步处理时维护 per-key chain;
- 提交 offset 与处理结果严格对应。
Kafka 保证的是读取顺序,不自动保证你的业务处理完成顺序。
14.3 Producer 重试
生产者应启用幂等能力并合理配置重试。业务层仍应携带 message_id 或事件 ID,让消费者幂等。
不要把“Kafka exactly-once”理解成整个 IM 链路天然只处理一次。数据库、推送网关、客户端都需要自己的幂等边界。
15. “Exactly Once”在 IM 中通常是用户体验目标
网络系统中更实际的实现是:
1 | |
去重键:
- 服务端发送幂等:
senderId + clientMsgId; - 消息存储唯一:
messageId; - 会话顺序唯一:
conversationId + msgSeq; - 客户端展示去重:
messageId。
服务端可能重复推送,客户端必须安全忽略重复消息。
16. ACK 设计
可以区分:
1 | |
16.1 已读回执不要逐条写
对一个会话维护:
1 | |
当用户读到 108:
1 | |
这比对每条消息写一条已读记录更高效。
群聊是否展示每个人的已读状态,要根据产品和规模单独设计。
17. 撤回和编辑顺序
撤回不是删除原消息并让序列消失,而应当作状态变更事件:
1 | |
客户端先收到撤回事件但缺少原消息时:
- 先保存 tombstone;
- 后续原消息到达时直接展示“消息已撤回”;
- 同步服务按事件序列恢复最终状态。
编辑同理,应有版本号:
1 | |
不能让旧编辑事件覆盖新版本。
18. 单聊和群聊的差异
18.1 单聊
- 会话参与者少;
- 每条消息扇出目标少;
- 会话热点通常有限;
- 每用户同步状态较简单。
18.2 大群
- 单会话写热点明显;
- 一条消息可能扇出给数十万成员;
- 不能同步写每个成员 Inbox 后才 ACK;
- 需要“写一次消息日志 + 用户按游标拉取”或分层扇出;
- 群成员变更和历史可见性需要定义;
- 单会话严格序列器可能成为瓶颈。
大型群聊可采用:
1 | |
而不是为每个成员复制完整消息。
19. 热点会话
一个超大直播群或万人群会让单会话有序域变成热点。因为严格单会话顺序天然限制并行写入。
优化方向:
- 明确是否真的需要所有消息严格全序;
- 将控制消息和普通聊天消息分流;
- 对弹幕类消息允许弱顺序或分区顺序;
- 限流和合并系统通知;
- 使用专门的热点会话分片;
- 批量写日志;
- 客户端按小窗口排序。
这是不可逃避的权衡:
1 | |
20. 历史消息分页
查询接口:
1 | |
SQL:
1 | |
首次读取可以从会话最新序列开始。向前翻历史时使用 beforeSeq,不使用 OFFSET。
断线补拉:
1 | |
21. 多端同步
每台设备维护自己的设备游标:
1 | |
服务端可以维护用户级同步事件流:
1 | |
设备断线重连后按 globalSyncSeq 拉取变化,再按各会话的 msgSeq 补齐正文。
flowchart LR
D1[手机游标 500] --> SYNC[用户同步日志]
D2[电脑游标 470] --> SYNC
D3[平板游标 520] --> SYNC
SYNC --> C1[会话 C1 增量]
SYNC --> C2[会话 C2 增量]
SYNC --> STATE[已读/撤回/关系状态]
用户级同步序列和会话内消息序列解决的是不同问题,不能混为一个 ID。
22. 跨机房设计
22.1 同一会话单写地域
最简单可靠:
1 | |
其他地域收到发送请求后转发到 homeRegion。
优点:顺序清晰;缺点:跨地域延迟。
22.2 多地域同时写
如果同一会话允许多个地域同时写,要解决:
- Leader 选举;
- 冲突排序;
- 网络分区;
- 脑裂;
- 一致的 fencing;
- 用户看到的临时顺序是否允许重排。
可使用共识协议或全局排序服务,但成本高。多数普通 IM 会优先选择会话级单写主域,而不是追求无条件多主。
23. 故障场景
23.1 服务端落库成功,ACK 丢失
客户端使用相同 clientMsgId 重试,服务端返回原 messageId/msgSeq。
23.2 序列分配成功,落库失败
允许产生空洞,或把序列分配和日志提交放进同一复制状态机。绝不能复用已暴露的旧序列给另一条消息。
23.3 Kafka 重复投递
下游以 eventId/messageId 幂等,推送可重复,客户端按 messageId 去重。
23.4 客户端先收到 102 后收到 101
客户端按 msgSeq 缓冲重排,并触发缺口补拉。
23.5 分片 Leader 切换
新 Leader 必须拿到更高 epoch/fencing token,旧 Leader 即使恢复也不能继续写。序列状态从复制日志或可靠存储恢复,不能回退。
23.6 消息存储已成功但推送系统不可用
发送端仍可得到 SENT;接收端恢复连接后通过同步游标拉取,不依赖实时推送完整可靠。
24. 监控指标
- 消息发送 QPS、成功率、P99;
- 发送重试和幂等命中率;
- 同会话序列重复数;
- 序列回退检测;
- 序列空洞率;
- 客户端缺口补拉率;
- 重复推送率;
- Kafka 分区积压;
- 热点会话 QPS;
- 分片 Leader 切换次数;
- 从持久化到在线送达延迟;
- 离线同步成功率;
- 多端游标落后量。
应建立不变量报警:
1 | |
25. 常见错误回答
错误一:按消息时间戳排序
设备时间和网络延迟都不可靠。
错误二:使用全局单序列器
可以保证顺序,但很可能制造不必要的系统瓶颈。业务通常只需要会话内有序。
错误三:Kafka 能自动保证整个系统顺序
Kafka 只保证分区内读取顺序,业务线程池、数据库提交和推送仍可能乱序。
错误四:追求整个链路物理 Exactly Once
更现实的是至少一次传输与端到端幂等去重。
错误五:发现 seq 缺口就永久等待
服务端号段可能存在合法空洞,同步协议必须能确认权威进度。
错误六:实时推送丢了就丢了
IM 必须有离线同步和游标补拉,推送只是低延迟通道。
26. 面试总结话术
我不会追求全系统全局顺序,而是定义同 conversationId 内的稳定顺序。入口连接可以任意接入,但消息路由必须按 conversationId 进入同一个分片 Leader 或 Kafka 分区。客户端使用 clientMsgId 保证重试幂等,服务端为落库消息分配 messageId 和会话内单调递增 msgSeq,并在持久化成功后 ACK。Kafka 事件也以 conversationId 为 key,消费者内部按 key 串行处理。接收端即使乱序收到消息,也按 msgSeq 缓冲、去重并检测缺口,通过 afterSeq 接口补拉。多端同步再维护独立的用户/设备同步游标。系统接受序列空洞,但绝不接受重复和回退。
27. 延伸追问
- 为什么 Snowflake ID 不能直接等价于会话顺序?
- 一个十万人群如何避免单会话序列器成为瓶颈?
- Kafka 扩分区后,key 到分区映射变化怎么办?
- 消息先落库还是先写 Kafka?
- 如何实现消息撤回在多端一致?
- 服务端序列出现空洞,客户端怎样确认不是丢消息?
- 两地多活同时发送如何排序?
- 图片上传失败是否应该阻塞后续文本消息?
参考资料
- Apache Kafka Documentation: https://kafka.apache.org/documentation/
- Kafka Message Delivery Semantics: https://docs.confluent.io/kafka/design/delivery-semantics.html
- Redis INCR: https://redis.io/docs/latest/commands/incr/
- Transactional Outbox: https://microservices.io/patterns/data/transactional-outbox.html