聊天消息架构:连接落在哪、消息怎么按顺序到每台设备、人不在线时存在哪、怎么叫醒
比较 Gateways with Pub/Sub Fan-out、Per-conversation Log with Cursor Sync、Per-device Inbox Fan-out on Write 与 Presence-routed Delivery with Store-and-forward Push 四种聊天拓扑:连接落在哪个网关、消息怎么找到另一台机器上的收件人、会话的持久记录是什么、cursor 归谁、扇出发生在写入时还是读取时、离线设备怎么知道有消息在等、状态归谁,以及一次网关重启或一份过期的在线状态能波及多远、怎么收住。
一个聊天系统同时开着大量长连接,从其中一条收进一句话,然后要让它按顺序出现在这个会话每个成员的每台设备上:在线的现在就到,不在线的回来时补上。所有聊天方案都要回答四个问题:连接落在哪、消息怎么找到另一台机器上的收件人;会话的持久记录是什么、进到哪里的 cursor 归谁;扇出发生在写入时还是读取时;离线的设备怎么知道有东西在等它。这四个答案,而不是用的哪家产品,才定义了拓扑。
想亲手走一遍四种拓扑、注入故障再恢复,打开互动 Lab:/system-design-lab/chat-messaging-architectures。
有约束的设计问题
一家公司先后要做四种聊天:几千人同时在线的直播间,消息和「正在输入」要秒级到所有人,错过一两秒的人重连补一下就行;有完整历史的频道,用户在手机和电脑上都看,断线几小时后重连一条不能少、顺序不能乱;端到端加密的一对一和小群,服务器读不到内容,设备之间不共享密钥,离线的设备回来要能解开等它的消息;移动端为主的产品,任一时刻大多数收件人都不在线,消息要在他们打开应用时完整有序地出现,还不能被几十条通知轰炸。每一种该怎么搭?
四种 topology signature
| 架构 | 连接落在哪 | 持久记录 | 扇出 | 顺序 | 在线状态 | 离线投递 | State owner | 适合 | 主要代价 |
|---|---|---|---|---|---|---|---|---|---|
| Gateways with Pub/Sub Fan-out | 任一网关;注册表记用户在哪台 | 消息库,发布之前先写 | 写入时,在内存里广播给每台订阅的网关 | 库分配的会话内 seq;发布订阅按发布顺序送达 | 输入中、上下线走同样的频道 | 发布订阅不提供;客户端重连后从库补 | 消息库 | 成员大多在线的房间、输入中提示这类短命信号 | 至多一次:错过的发布要等客户端重新拉取才可见;每台网关收到它订阅的全部频道 |
| Per-conversation Log with Cursor Sync | 任一网关;读是按设备的 | 会话流本身 | 读取时:每台设备从自己的 cursor 读 | 流的 id,会话内单调递增 | 另建 | 天然:离线设备回来从 cursor 读 | 会话流 | 有历史的频道、多设备用户、可靠补齐 | 读量随设备数乘会话数增长;热会话就是一条流;cursor 要按设备持久存 |
| Per-device Inbox Fan-out on Write | 任一网关;每台设备拉自己的收件箱 | 每台设备的收件箱 | 写入时:每台收件设备一份 | 收件箱内按追加顺序 | 另建 | 天然:收件箱等到设备来拉 | 收件设备的收件箱 | 端到端加密的一对一和小群,每台设备有自己的密钥会话 | 写放大等于收件设备数;大群成倍;丢掉的设备留下一直涨的收件箱 |
| Presence-routed Delivery with Store-and-forward Push | 在线用户在网关;离线用户没有 | 离线存储,加上推送服务保留的通知 | 写入时,按在线状态逐个收件人路由 | 按收件人各自的存储 | 存在感服务就是路由依据 | 先存,再发带时限和 topic 的推送;醒来后同步 | 离线存储 | 移动端为主、多数收件人随时离线的产品 | 在线状态只和心跳一样新;推送是叫醒不是投递;平台限制时限和存储条数 |
长轮询和 SSE 是网关拓扑的另两种传输方式,比较见 长轮询、WebSocket 与 SSE;分片发布订阅(频道按 slot 哈希)是网关拓扑在集群里的扩展形态;基于花名册的状态订阅是存在感拓扑的授权模型。它们在这里解释,不单独画。后台任务的队列语义见 消息队列,WhatsApp 的整体案例见 Design WhatsApp。
1. Gateways with Pub/Sub Fan-out:连接落在网关上,先写库再广播,订阅的网关各自推给自己的连接
MDN 的 WebSocket 服务端指南把连接讲清了:握手是一次带 Upgrade: websocket 的 HTTP 请求,服务端回 101 Switching Protocols;之后「客户端或服务端都可以随时选择发送消息——这就是 WebSocket 的魔力」;「握手之后的任何时候,客户端或服务端都可以选择向对方发一个 ping。收到 ping 的一方必须尽快回一个 pong」;「要关闭连接,客户端或服务端都可以发一个带指定控制序列的控制帧来开始关闭握手」。一个用户的连接同一时刻只在一台网关上。
Redis 的 Pub/Sub 文档把广播的性质写死了:「Redis 的 Pub/Sub 表现为至多一次的消息投递语义」;「一旦消息被 Redis 服务器发出,就没有机会再发一次。如果订阅者无法处理这条消息(例如因为错误或网络断开),这条消息就永远丢失了」;「订阅者按消息发布的顺序收到消息」;「如果你的应用需要更强的投递保证,你可能想了解 Redis Streams」;从 Redis 7.0 起「引入了分片 Pub/Sub,分片频道用和 key 分配到 slot 相同的算法分配到 slot」,并「把消息的传播限制在集群的一个分片内」。Socket.IO 的 Redis adapter 展示了实际形态:发给多个客户端的每个包「发给连接到当前服务器的所有匹配客户端」并「发布到一个 Redis 频道,由集群里其他 Socket.IO 服务器收到」;「Redis adapter 用 Pub/Sub 机制在 Socket.IO 服务器之间转发数据包,所以 Redis 里不存任何 key」;「对新开发,我们推荐使用分片 adapter,它利用 Redis 7.0 引入的分片 Pub/Sub 特性」。
所以顺序是固定的:A 的消息先写进消息库拿到会话内的 seq,网关再把它发布到这个会话的频道;订阅了频道的每台网关各自推给自己手上属于这个会话的连接;A 收到带 seq 的 ack。广播只负责快,不负责不丢:客户端带着自己最后看到的 seq,重连时向库要 seq 之后的全部。
典型事故:一次滚动发布重启了 Gateway 2,它的订阅断了几秒,期间发布的消息对它的连接永远不存在——B 屏幕上少一条而 A 看到已发送;修法是先写库再发布,网关重启会断开连接,客户端重连时按 seq 从库补齐。为了简单把所有消息发到一个全局频道,每台网关收到全站消息再本地过滤,集群总线随网关数线性膨胀——频道按会话命名,网关只订阅自己有连接的会话,集群里用分片发布订阅。
2. Per-conversation Log with Cursor Sync:每个会话一条只追加的流,每台设备记自己的 cursor
Redis Streams 的文档:「Redis stream 是一种数据结构,它像一个只追加的日志,但也实现了若干操作来克服典型只追加日志的一些限制」;条目 id 是 <millisecondsTime>-<sequenceNumber>,「即使时钟向后跳,id 单调递增的性质仍然成立。序列号用于同一毫秒内创建的条目」;「Redis streams 支持按 id 的范围查询」;XREAD ... BLOCK 会等待新条目,所以实时投递就是从最后一个 id 起的一次阻塞读;消费组加上待处理条目列表和 XACK 用于协调消费;MAXLEN 给流封顶,旧条目会被淘汰。
会话本身就是日志。每条消息追加进这个会话的流拿到一个流内单调递增的 id,顺序就是流的顺序,历史就是流本身。每台设备为每个会话记着自己最后确认看到的 id,投递就是从这个 cursor 往后读:还有没读的就一条条给,追上了就阻塞等新的。实时投递和断线补齐是同一个操作:B 的手机离线一小时、期间会话里来了四十条,重连后网关查它的 cursor,从那个 id 之后读流,四十条按 id 顺序推过去,追上末尾后阻塞;B ack 最后的 id,cursor 前进。B 的笔记本有自己的 cursor,互不影响。
典型事故:B 的手机重装了系统,本地 cursor 没了,服务端没存——它要么从流的开头把几年的历史全拉一遍,要么从末尾开始把离线期间全部漏掉;cursor 按设备、按会话存在服务端的 cursor 存储里并在 ack 时更新,本地只是缓存,新设备第一次登录从最近一段历史开始。为了省内存用 MAXLEN 只保留最近一万条,B 的笔记本离线三个月,它的 cursor 指向早已淘汰的 id——截断前先归档,cursor 早于流起点的设备从归档补,截断阈值按最长可接受的离线时间定。
3. Per-device Inbox Fan-out on Write:写入时就给收件人的每台设备各投一份
Signal 的 Sesame 算法是「在异步、多设备环境下管理消息加密会话」的规范:「每台设备为它的通信对象存一组 UserRecord,按 UserID 索引。每个 UserRecord 包含一组 DeviceRecord,按 DeviceID 索引」;发送时「对 UserRecord 里每个含有活跃会话、且未过期的 DeviceRecord,发送设备用那个活跃会话加密明文」,包括发件人自己的其他设备。XMPP 的多 resource 模型写下了投递规则:用户的服务器「必须把花名册推送发给所有感兴趣的 resource」,联系人的状态要送到「用户的每一个可用 resource」,联系人的服务器也从「联系人的每一个可用 resource」向用户发当前状态;一个用户的每个已连接客户端都是一个 resource,各收自己的一份。
地址是设备,不是人。发件方知道收件人有哪些设备,为每台设备的会话各加密一份,服务器把每份追加到对应设备自己的收件箱,连发件人自己的其他设备也各得一份,否则 A 的笔记本永远看不到 A 手机发的消息。每台设备只拉自己的收件箱,处理完 ack,服务器删掉已确认的条目;离线的设备什么都不用做,它的收件箱一直等着。
典型事故:五百人的端到端加密群、每人三台设备,一条消息变成一千五百次写入,活跃的群把存储先于聊天本身压垮——端到端加密的群接受这个放大并限制群大小,发件方按设备批量提交,不需要每设备密文的大房间改用每会话一条日志。B-phone 拉取并显示了一条,ack 在网络里丢了,下一次拉取又拿到同一条——客户端按消息 id 去重,重复的直接再 ack 一次而不显示。
4. Presence-routed Delivery with Store-and-forward Push:按在线状态路由,离线先存再推一条有时限的通知
RFC 6121 关于存在感:「存在感信息只向用户已批准的其他实体披露」;用户上线时它的服务器把「联系人的每个可用 resource」的状态送到「用户的每个可用 resource」;花名册变化「推送给所有感兴趣的 resource」。RFC 8030 关于推送:「推送服务的通用模型包括三个基本角色:用户代理、推送服务和应用服务器」;「推送服务通过把推送消息保存一段时间,可以大幅提高投递的可靠性」;「应用服务器在请求推送消息投递时必须包含 TTL(Time-To-Live)头」,而「推送服务可以把推送消息保留比请求更短的时间」;「应用服务器可以在请求里包含 Urgency 头」;「带 topic 的推送消息会替换任何未送出的、topic 相同的推送消息」。
在线状态是路由的依据。存在感服务靠心跳知道谁有活着的连接,只向订阅了的联系人广播变化。路由器收到消息先问它:收件人在线就送到它所在的网关;离线就先把消息写进离线存储,再请平台的推送服务发一条小通知去唤醒设备。推送不是投递:它只是一个带时限的叫醒信号,通知里不放正文,同一个 topic 的新通知替换还没送出的旧通知,十条消息只叫醒一次;设备醒来后带着自己的 cursor 去离线存储同步,才拿到消息本身,按顺序、一条不少。A 拿到的回执先是「已存,等 B 上线」,B 同步之后才变成已送达。
典型事故:B 的手机进了电梯,连接已断但心跳还没超时,存在感服务说在线,路由器把消息送到 B 原来的网关,网关往一条死连接写帧,什么也没发生——先写离线存储再路由,网关推帧后等设备 ack,超时未 ack 就回退到离线路径并发推送,心跳超时把 B 标为离线并通知订阅的联系人。每条离线消息发一条没有 topic 也没有 TTL 的推送,B 关机一晚早上收到几十条堆在一起、内容早已过时的通知,平台还丢掉了超出上限的部分——同一会话的推送用同一个 topic,设 TTL 让过期的作废,通知只带「去同步」的提示。
四种拓扑共同的底线
- 说清状态归谁。 消息库、会话流、每台设备的收件箱,还是离线存储;广播、推送和网关都不是。
- 先存再送。 写库在发布之前,写离线存储在推送之前;先送后存的消息一旦送丢就没有第二次。
- 说清设备重连时做什么。 按 seq 从库补、从 cursor 读流、拉自己的收件箱,还是从 cursor 同步离线存储。
- 说清一次发送在哪里变成多次投递。 订阅的网关、每台设备的读、写入时的每设备一份,还是按收件人的路由。
- 把连接绑死在身份上。 升级握手时认证,一条连接对应一个用户和一台设备,帧里不允许客户端自己写发件人和会话 id;在线状态是个人数据,只给订阅的联系人。
故障与恢复
| 架构 | 故障 | 用户看到什么 | 恢复 |
|---|---|---|---|
| Gateways with Pub/Sub | Gateway 2 在发布时重启 | B 少一条,A 显示已发送 | 先写库再发布;客户端重连按 seq 从库补 |
| Gateways with Pub/Sub | B 的连接半死,帧写进去就消失 | B 什么都收不到 | ping/pong 超时关闭连接并从注册表移除;B 重连补齐 |
| Gateways with Pub/Sub | 一个全局频道让每台网关收到全部消息 | 网关和总线先于连接数撑不住 | 频道按会话命名;集群用分片发布订阅 |
| Per-conversation Log | 设备的 cursor 丢了 | 重读整段历史或漏掉离线期间 | cursor 按设备、按会话存服务端,ack 时更新;本地只是缓存 |
| Per-conversation Log | 历史被 MAXLEN 截断到 cursor 之前 | 离线太久的设备补不回中间那段 | 截断前归档;cursor 早于流起点时从归档补 |
| Per-device Inboxes | 五百人群、每人三台设备 | 所有群的消息延迟一起上升 | 限制加密群大小、批量追加;大房间改用会话日志 |
| Per-device Inboxes | ack 丢了,同一条拉到两次 | 同一条消息显示两遍 | 客户端按消息 id 去重并再 ack;服务器收到 ack 才删 |
| Presence-routed with Push | 在线状态过期,消息推进死连接 | 消息只存在于 A 的屏幕上 | 先写离线存储再路由;等 ack,超时回退离线路径并推送;心跳超时标离线 |
| Presence-routed with Push | 十条消息发十条推送,没有 topic 和 TTL | 几十条过时通知堆在一起,多的被丢 | 同一会话同一个 topic 合并;设 TTL;通知只提示去同步 |
怎么选
- 成员大多在线、信号短命:Gateways with Pub/Sub Fan-out,先写库再发布,频道按会话而不是全局,集群里分片发布订阅,客户端每次重连都从库重新拉取。
- 有历史的频道、多台设备:Per-conversation Log with Cursor Sync,用单调递增的 id 追加,每台设备每个会话一个持久 cursor,从 cursor 起阻塞读,历史归档之后才截断。
- 端到端加密、每台设备自己的会话:Per-device Inbox Fan-out on Write,给每台设备包括发件人自己的其他设备各投一份,ack 后删除,再也不回来的设备的收件箱设过期,群保持小或接受写放大。
- 移动端为主、收件人通常离线:Presence-routed Delivery with Store-and-forward Push,按在线状态路由但假设它可能过期,先存后推,每条推送带时限和 topic 让待送的合并,设备醒来从 cursor 同步而不是信推送里的内容。
- 无论哪种:写下谁拥有记录、设备重连时做什么、哪个组件把一次发送变成多次投递。
回到开头的四种聊天:直播间走 Gateways with Pub/Sub;有历史的频道走 Per-conversation Log;端到端加密走 Per-device Inboxes;移动端离线为主走 Presence-routed Delivery with Push。
面试时这样回答
- 先复述约束:成员大多在线还是离线、有没有历史、一个人几台设备、服务器能不能读内容、能不能容忍错过一两秒。
- 说消息路径:点名图上的边,例如「先写库拿 seq 再发布,订阅的网关推给自己的连接」「XADD 进会话流,查 cursor 后从它读并阻塞」「按设备加密追加到各自收件箱,拉取后 ack」「问存在感,先存离线存储再推带 TTL 的通知,醒来按 cursor 同步」。
- 说状态归谁:消息库、会话流、每台设备的收件箱,还是离线存储。
- 说代价与一个故障:例如在线状态过期把消息推进死连接,修法是先存离线存储再路由、等设备 ack、超时回退到离线路径并发推送、心跳超时标为离线。