T18 · Realtime System
怎样把链上的变化实时推给用户?
- 练习的能力
- Builder
- 动手
- 实现一个订阅接口,断开 30 秒后重连,验证中间的事件都能补齐。
- AI Lab
- 让 AI 设计断线重连协议,自己在弱网下实测它的丢包表现。
一个现实问题
你的索引器已经很可靠了(T17 那一章的成果),现在要把变化推给用户。你加了一条长连接,索引器每处理完一个区块就把事件推出去。本地测一遍,秒级到达,非常爽。
上线两周后,客诉里开始出现同一类描述,而且很难复现:
- 「我的活动列表少了几笔,刷新一下又有了。」
- 「有一笔一直显示处理中,其实早就成功了。」
- 「界面上的余额和我在浏览器上看到的对不上。」
你去查服务端日志,干干净净:消息都发出去了。 发送函数返回成功,没有异常,没有堆积告警。
然后你注意到这几个用户有一个共同点:他们都在移动网络下,都经历过短暂断线。
这时候问题就清楚了。你的推送是「发出去就忘」的:那一刻谁连着,谁就收到;谁没连着,那条消息就永远消失了。它不是丢在网络上,是丢在你的设计里——因为你从来没有为「他不在」这件事准备过任何东西。
于是真正的问题是:用户掉线 30 秒再回来,中间发生的事,你从哪里补给他?
思想实验
三种把消息送出去的方式。
第一种:电台广播。 你在演播室念,谁开着收音机谁听见。你不在收音机前,那段节目就没了,而且没有任何人知道你错过了什么。
这种方式的特点是:服务端不需要记住任何关于听众的事,成本极低。代价是「送达」完全取决于运气。
第二种:留言板加编号。 每条消息贴在墙上,带一个连续递增的编号。你回来时对管理员说「我看到第 812 号了」,管理员把 813 之后的念给你听。
现在「掉线」这件事从一个不可恢复的事故,变成了一次可以补齐的追赶。代价是有人得保存这面墙,而墙的容量是有限的——旧消息总会被擦掉。
第三种:只喊「有变化了,自己去看」。 不传内容,只传一个信号。听到的人自己去查最新状态。
这一种有一个被低估的优点:丢掉一个信号几乎没有后果,因为下一个信号来的时候会一起刷新。它天然幂等、天然抗丢。它答不了的问题是「中间发生了什么」——只能给你「现在是什么样」。
三种方式指向同一个判断标准:
实时推送要么可重放,要么可收敛。两者都不是,就是在赌网络。
第一种既不可重放也不可收敛,所以它必然会丢。开头那三条客诉,就是这一件事的三种表现。
你来决定
给「链上一有变化就更新界面」这个需求选一种做法。
观察结果
把四种做法按两个不同的指标排一次,你会看到一件很多人没分清的事:
| 延迟 | 掉线会不会丢 | 服务端状态成本 | 能不能回答「中间发生了什么」 | |
|---|---|---|---|---|
| 定期拉取 | 高 | 不会丢状态,会丢过程 | 无 | 不能 |
| 推完就忘 | 最低 | 会永久丢 | 无 | 能,但不可靠 |
| 带编号加补齐 | 最低 | 不会 | 高(要存窗口) | 能 |
| 只推变化信号 | 低 | 不会丢状态,会丢过程 | 极低 | 不能 |
第一列和第二列没有任何关系。这就是这一章最该带走的一句话:
「实时」是一个延迟指标,「不丢」是一个正确性指标,两者由完全不同的机制保证。
绝大多数「实时推送」只做了第一件事。它上线时看起来很好,因为开发和测试都在稳定网络下;它的问题只会出现在用户那里,而且以「偶尔少几条」这种最难排查的形式出现。
还有第二条结论,它把这一章和上一章缝在了一起:
链头的事件是可能被撤回的。已经推给用户的事件,不会自己收回来。
所以推送流里必须有一类「撤销」消息,或者干脆只推已经足够深的事件——也就是 T17 那张待发表里的投递条件。「不丢」和「不假」是两个独立的要求,都要做。
建立模型
三个要素
一个不丢事件的推送系统,只需要三样东西:
- 单调游标
- 可重放窗口
- 重连时的补齐协议
游标从哪来:三个候选,两个是错的
时间戳——错。 它不严格单调(同一毫秒内多条)、受时钟调整影响、多实例之间还会漂移。用时间戳做游标,补齐时要么重复要么遗漏。
数据库自增序列——大多数情况下也是错的。 原因很隐蔽:序列号是在插入时取的,提交却不同步。事务 A 拿到 102,事务 B 拿到 103,B 先提交。这时候一个读取者看到了 103,把游标推进到 103——然后 A 才提交,102 就永远不会被任何人读到了。
这个空洞不会报错,只会让某些事件消失。要用它,必须配合「只读取足够旧的记录」或者「单一写入者串行分配」。
链上坐标——对。 (chain_id, block_number, log_index) 这个组合本身就是一个严格单调、全局唯一、跨实例一致、而且重启后还能对得上的游标。它还有一个额外好处:客户端拿着它,可以直接去链上核对。
把三段编码成一个可比较的字符串或者一个大整数,就是你的游标:
cursor = chainId : blockNumber : logIndex
例: 1 : 00000000020134 : 0007消息协议:五种消息,缺一不可
不依赖任何具体的库,先把协议本身写清楚:
客户端 → 服务端 subscribe { topics, cursor } 订阅,带上我看到的位置
服务端 → 客户端 snapshot { cursor, state } 当前状态快照,可选
服务端 → 客户端 event { cursor, type, payload }
服务端 → 客户端 revoke { fromCursor } 这个位置之后的事件作废
服务端 → 客户端 reset { reason } 你的游标太旧,请重新订阅
双向 ping / pong 心跳后面三种是新手最容易漏的:
reset 必须有。 可重放窗口是有限的,客户端离线一小时回来,它的游标可能已经被清理掉了。没有 reset,服务端会静默地不补任何东西,客户端会以为自己已经是最新的——这是一种比掉线更糟的状态,因为它看起来是正常的。
revoke 是重组的出口。 如果你选择推送未确认的事件,就必须有能力告诉客户端「刚才那几条作废」。不想实现它,就只推足够深的事件,把问题交给确认深度。
心跳必须是双向的。 只有客户端发心跳,服务端可以正确判死;只有服务端发,客户端可以正确判死。两边都要有超时判断,否则会出现一方以为连着、另一方早就走了的半开连接。
快照与增量的接缝:这一节是最容易写错的地方
「先取当前状态,再订阅后续变化」听起来天经地义,但它是错的:
错误顺序:
t0 查询快照(拿到截至 cursor=100 的状态)
t1 ← 这一瞬间发生了 cursor=101 的事件,没有人在监听
t2 建立订阅,从此收到 102 及之后
结果:101 永久丢失,而且没有任何迹象正确顺序是先订阅、缓冲、再取快照、最后放行:
1. 建立连接并订阅,服务端开始往这个连接的队列里放事件
2. 客户端(或服务端)此时不投递,先把事件缓冲起来
3. 取快照,得到快照对应的 cursor
4. 把缓冲里 cursor 大于快照 cursor 的事件依次放出去,丢弃更早的
5. 之后转入正常投递这一段顺序写反,是「重连后少几条」这类问题最常见的根因,而且它在本地测试中几乎不可能复现——你需要恰好在那一瞬间有一个事件。
背压:出口永远是降级,不是等待
一个客户端读得比你发得慢(弱网、后台标签页、一个写得不好的客户端)。三种处理方式:
| 做法 | 后果 |
|---|---|
| 阻塞等它读完 | 一个慢客户端拖垮整个服务,所有人一起卡住 |
| 无限缓冲 | 内存一路涨到进程被杀,而且是在流量高峰时 |
| 有界队列,满了就降级 | 这个客户端收到 reset,自己去拉快照,其他人不受影响 |
第三种是唯一正确的答案。把它记成一句话:
背压的出口永远是降级,不是等待。
具体做法:每个连接一个有界队列(比如一千条)。队列满时,丢掉这个连接的全部排队事件,给它发一个 reset,让它重新走一遍「订阅、快照、补齐」的流程。
注意不要只丢最老的那几条——那会在客户端造成一个它察觉不到的空洞。要丢就全丢,然后明确告诉它「你需要重来」。
重连:退避要带抖动
客户端断线后立刻重连,失败再立刻重连,会在服务端故障时把它彻底压死。标准做法是指数退避:
第 n 次重试的等待时间 = min(基础间隔 × 2 的 n 次方, 上限) × 随机系数
随机系数取 0.5 到 1.5 之间那个随机系数不是可有可无的。没有它,一次服务端重启会让所有客户端在同一秒回来,第二次重启又会让它们在同一秒回来——你会得到一个周期性的自我攻击。
交付语义:不要追求服务端的恰好一次
网络层面做不到「恰好一次」。可行的组合是:
服务端:至少一次投递(重连补齐可能重发几条)
客户端:按 cursor 去重(小于等于已见 cursor 的直接丢弃)
合起来:效果上的恰好一次这和 T15 的幂等、T16 的主键幂等、T17 的断点幂等是同一个套路的第四次出现:不追求「只做一次」,而是保证「做多次和做一次结果相同」。
扇出与多实例
每个连接一个有界发送队列,一张从主题到连接的索引表,事件来了查表投递。这在单实例下很直接。
多实例时会冒出一个诱惑:把事件塞进一个消息中间件,各实例订阅。可以做,但要守住 T16 那条规则:事实只有一份。 中间件是传输通道,不是事实来源——补齐时必须回数据库按游标查,而不是去问中间件要历史。
理由很简单:中间件的保留策略、分区顺序、消费位点,全都是另一套可能出错的状态。让它只负责「快」,让数据库负责「对」。
它叫什么
浏览器与服务端之间的一条双向长连接,两边都能主动发消息。
它解决的只是通道问题。「掉线不丢事件」这件事它一点忙都帮不上——那是你的协议设计要回答的。把长连接等同于可靠推送,是这一章要纠正的第一个误解。
客户端声明「我关心哪些主题」,服务端据此过滤要推给它的事件。
订阅请求里必须能带上游标。只带主题不带位置的订阅,等于每次重连都从「现在」开始,掉线期间的一切自动消失。
标记「我收到哪里了」的单调递增位置。
时间戳和数据库自增序列都不可靠,前者不严格单调,后者会因为提交顺序留下空洞。链上坐标是这里天然的最优解:严格单调、全局唯一、重启后依然有效,而且客户端可以拿它去链上核对。
客户端带着游标回来,服务端把它之后的事件按顺序补发。
它成立的前提是服务端保留了一个可重放窗口。窗口是有限的,所以协议里必须有「你的游标太旧了」这条消息——否则客户端会停在一个看起来正常、实际上永远缺数据的状态里。
消费方比生产方慢时,压力沿着链路往回传导。
处理它只有一个正确方向:降级,不是等待。 每个连接一个有界队列,满了就清空并要求对方重新拉快照。阻塞会让一个慢客户端拖垮所有人,无限缓冲会让进程在高峰期被杀掉。
双方定期互发的探活消息,用来发现对方已经不在了。
必须双向。只有一方发心跳,另一方就可能长期维持一条半开连接:一边以为还连着、往里写,另一边早已消失。这类连接会安静地堆积,直到耗尽资源。
动手
全程本地与测试网,只读。 这一章不签名、不广播、不花钱。
验收标准是一条可以自动跑的断言:
断线 30 秒后重连收到的事件序列,
必须与全程不断线的对照客户端收到的序列完全一致(数量、顺序、内容)。接在 T17 的待发表上,用链上坐标做游标。
不要另起一套事件源。推送要发的事件就是 outbox 里的行,投递条件沿用那一章的「区块高度足够深」。
游标直接用 (chain_id, block_number, log_index) 编码成一个可比较的字符串。这一步选对了,后面的补齐查询就是一条普通的范围查询。
实现五种消息,reset 不要省。
按上面那张协议表实现 subscribe、snapshot、event、reset 和双向心跳。revoke 可以先不做——因为你只投递足够深的事件。
补齐查询就是一条 SQL:
select * from outbox
where chain_id = $1
and (block_number, log_index) > ($2, $3) -- 客户端给的游标
and block_number <= $safeHeight
order by block_number, log_index
limit 1000;超过 limit 时不要一次性发完再继续,分批发,每批之间检查连接还活着没有。
把接缝顺序写对。
按「先订阅、缓冲、再取快照、按游标放行」的顺序实现。
写完做一个针对性的实验:在取快照的那一行代码前后各加一秒延迟,然后在那一秒里手动触发一个事件。这个事件必须出现在客户端。 如果没有,说明你的顺序写反了——而这正是线上「偶尔少几条」的来源。
验收:断线 30 秒。
开两个客户端,订阅同样的主题。一个全程保持连接当对照组,另一个在中途断开 30 秒再带着游标重连。
跑完之后比较两边收到的事件序列:
- 数量一致吗?
- 顺序一致吗?
- 逐条内容一致吗?
- 重连的那个有没有收到重复的?收到了的话,去重逻辑生效了吗?把这四项写成一个脚本,它就是你这个系统的回归测试。
做一个慢客户端,把背压逼出来。
写一个客户端,每读一条就 sleep 一秒。然后往系统里灌一万条事件。
先用无界缓冲跑一遍,观察服务端进程的内存。看着它一路往上涨,是这个 Lab 里最有教育意义的三分钟。
然后换成有界队列(比如一千条),队列满时清空并发 reset。再跑一遍,确认两件事:内存平稳;其他正常客户端完全不受影响。
在弱网下实测。
用你操作系统提供的网络模拟能力(各平台都有对应的命令行工具,先查一下帮助),给本地连接加上延迟、丢包和带宽限制。
在这三档下各跑一遍验收脚本:
- 延迟 500 毫秒
- 丢包 10%
- 带宽限制到很低,逼出背压丢包那一档最容易暴露心跳超时设得太短的问题:连接被反复判死重连,而实际上它还活着。
触发一次重组,确认没有假事件泄漏。
用 T17 的 Lab 里那个伪造重组的办法,让索引器回滚一次。
然后检查客户端:被回滚掉的那些高度上的事件,一条都不应该收到过。 因为它们在 outbox 里还没到安全高度,回滚时被一起删了。
如果你想体验反面:把安全高度调成 0 再试一次,你会亲眼看到一条「不存在的转账」推到客户端上——那就是 T1 第一章提出的那个问题的现场。
同时开 200 个连接,重启服务端。
观察它们的重连时间分布。如果全部挤在同一秒,说明你的退避没有加抖动。
加上抖动再试一次,重连应该被打散在一个区间里。这个实验在 200 个连接上就能看出差别,而线上是几万个。
AI Lab
分三步,第二步和第三步才是真正的考题。
第一步:
设计一个链上事件推送的订阅协议,要求断线重连后不丢事件。
列出所有消息类型、字段、以及客户端和服务端各自的状态机。
第二步:
现在逐个回答这些场景,说明你的协议里哪一部分在起作用:
1. 客户端断线 30 秒后重连,中间有 12 条事件。
2. 客户端断线 3 小时后重连,服务端的重放窗口只有 1 小时。
3. 一个客户端读取速度只有推送速度的十分之一,持续 10 分钟。
4. 服务端重启,5 万个客户端同时断开。
5. 客户端在「取快照」和「建立订阅」之间的那一瞬间,正好发生了一条事件。
6. 已经推给客户端的 3 条事件,因为链重组而不再有效。
第三步:
上面六个场景里,哪几个是你的协议目前处理不了的?
对每一个,说明最小的改动是什么,以及这个改动的代价。第一步模型通常给得相当完整,因为「带游标的订阅协议」是一个被反复写过的题目。要盯的是游标的选择:它很可能给你时间戳,而时间戳在多实例、时钟调整、同毫秒多事件这三种情况下都会出错。
第二步是真正的分水岭。六个场景里,第 2、3、5、6 条是高频失手项:窗口过期没有 reset、背压用了无限缓冲、接缝顺序写反、重组完全没提。第 4 条会暴露它有没有想到抖动。
第三步是这个 Lab 最值钱的部分,而且它考的不是知识,是诚实。一个好的回答会明确承认「我的协议处理不了第 3 条和第 6 条」并给出改法;一个差的回答会坚持说全都处理得了。让模型指出自己方案的缺口,比让它给方案有用得多——这条规律 T1 的 AI Lab 里就说过一次。
最后必须你自己做:在弱网下实测。 这一章所有的失败模式都只在真实网络条件下出现,而模型对「它的协议在 10% 丢包下表现如何」这类问题,只能给你一个听起来很合理的猜测。
AI 说完之后,你必须自己验证
- 它的游标是什么:时间戳和数据库自增序列都会出问题,追问它为什么不用链上坐标
- 协议里有没有「游标太旧、窗口已过期」这条消息;没有的话,客户端会静默地永远缺数据
- 心跳是不是双向的,两边都有超时判断吗
- 重连退避有没有随机抖动;没有的话,一次重启会让所有客户端在同一秒回来
- 快照和增量的接缝顺序:是先订阅后快照,还是先快照后订阅;自己构造一个恰好落在接缝上的事件来验证
- 背压策略是阻塞、无限缓冲还是有界丢弃加 reset;前两种都要退回
- 有没有处理链重组造成的事件撤销,或者明确说明只投递足够深的事件
- 它用的每一个库的 API、帧类型、关闭码,都要查文档核对是否真的存在
- 最后在弱网下实测:延迟、丢包、带宽限制三档各跑一遍,和不断线的对照组比对事件序列
真实案例
开头那个场景。客户端重连时先取快照再建立订阅,两者之间的那一瞬间发生的事件没有任何人在接收。
它的特征是概率极低但永不消失:一个活跃用户一天遇到一两次,客服收到零散投诉,而开发本地从来复现不出来——因为你需要恰好在那几十毫秒里有一个事件。
修复只是把顺序换过来:先订阅并缓冲,再取快照,再按游标放行缓冲里的事件。
推送用的是阻塞式写入:写不进去就等。某个客户端在极弱的网络下挂着,发送缓冲区满了,那次写入一直不返回。
如果所有连接共用一个投递循环,整个服务的推送就此停住,所有用户一起卡死。
换成无限缓冲也只是把故障形态换了一个:内存一路涨,在流量高峰时被系统杀掉——也就是在最不能出事的时候出事。
正确做法是有界队列加降级。
心跳只有服务端往客户端发,服务端没有检查客户端是否回应。
用户的网络在链路中间被切断(移动网络切换、NAT 超时),客户端那边连接已经没了,服务端这边还认为它活着,继续往它的队列里写。这类连接会越积越多,占着内存和主题索引,直到资源耗尽。
教训:心跳要双向,两边都要有超时判死。 这是一个写起来五分钟、不写就迟早出事的东西。
推送直接接在链头上,区块一出就推。重组之后,那几条事件对应的交易不再存在,但客户端上已经显示了。
这和 T17 开头那条投诉是同一个问题的两个出口:一个走通知,一个走界面。
两种解法都可以,但必须选一个:要么只投递足够深的事件,要么实现撤销消息。 什么都不做,等于向用户展示一段不存在的历史。
改一个变量
挂起时连接被静默断开,恢复时客户端以为自己还连着,往一条死连接上写,写了很久才发现。
处理办法是把「应用回到前台」当成一个明确的重连触发点:主动断开重连,带上游标补齐,而不是等 TCP 自己超时。
同时要接受一个现实:挂起可能持续几小时,游标大概率已经超出重放窗口。所以 reset 加重新拉快照这条路径,在移动端不是异常路径,是主路径。
五条连接、五份订阅、五倍扇出。而且它们的游标可能各不相同,补齐时会各查各的。
两条路:在客户端侧用一个共享的后台上下文把五个页面收敛成一条连接;或者在服务端接受这个成本,但把「按用户限连接数」写进限流规则——否则一个用户开着几十个标签页,就能占掉不成比例的资源。
这也是 T15 说的「限流要有两套」在推送场景的具体形态。
每个在线用户每秒一次查询。一万人在线就是每秒一万次,而其中绝大多数返回「没有变化」。
这时候你会自然地想做两件事:加一层缓存挡住重复查询,以及让请求挂住等待变化再返回。做完这两件事,你实际上已经把长连接重新实现了一遍,只是形式不同。
结论不是「轮询不好」,而是:轮询的成本随延迟要求急剧上升,而长连接的成本随连接数线性上升。 选哪个取决于你的延迟要求和在线规模,而不是取决于哪个听起来更先进。
三件事会同时压过来:连接本身的内存、主题到连接的索引查找、以及每条事件的扇出写入。
通常的演化顺序是:先把内容瘦下来(改成只推变化信号,让客户端自己拉,把扇出的成本转成可缓存的读取),再按主题分片(一个实例只负责一部分主题),最后才考虑更复杂的分发层。
不变的是那条规则:分发层只负责快,数据库负责对。 补齐永远回数据库按游标查——这是你在一百万连接下唯一不会崩掉的那条路径。
带走的问题
它解决什么问题?这一章解决的是「怎样让用户看到的东西和链上保持一致,即使他的网络不可靠」。注意这里有两个独立的目标:快(延迟)和不丢(正确性)。只做到第一个的系统,上线时看起来很好,问题全部由用户替你发现。
谁在支付?实时推送的成本是按在线连接数和扇出量计费的,而它往往是一个产品里最先失控的一项:连接内存、扇出写入、重放窗口的存储。把「只推变化信号」和「只推足够深的事件」放进设计,不只是为了正确,也是为了让这笔账算得过来。
谁承担风险?推送错了或者少了,用户看到的是一个和链上不一致的界面,而他没有任何办法知道自己看到的不完整。这类错误的代价不是那几条丢失的事件,是用户从此不再相信你的界面——于是他每次都去区块浏览器核对,你的产品就退化成了一个不必要的中间层。
本章自测
因为长连接只解决了「通道」,没有解决「他不在的时候怎么办」。
推完就忘的模式里,服务端不保存任何关于这条消息的状态:那一刻谁连着谁收到,没连着的就永远错过,而且服务端日志上看起来完全正常——发送函数返回成功了。
要不丢,需要三样东西:单调游标、可重放窗口、重连时带游标补齐的协议。缺任何一个都会丢,而且是静默地丢。
时间戳不严格单调:同一毫秒可能有多条,时钟调整会倒退,多实例之间还会漂移。
数据库自增序列的问题更隐蔽:序列号在插入时分配,事务却不按分配顺序提交。拿到 103 的事务先提交,读取者把游标推到 103,然后拿到 102 的事务才提交——102 永远不会被读到,而且不会有任何报错。
链上坐标(链 ID、区块高度、日志序号)天然严格单调、全局唯一、重启后依然有效,还能被客户端拿去链上核对。在这个场景里它几乎是免费的最优解。
先订阅并缓冲,再取快照,最后把缓冲里游标大于快照游标的事件放出去。
反过来(先取快照再订阅)会丢掉两者之间那一瞬间发生的事件,而且这个 bug 在本地几乎不可能复现——需要恰好在那几十毫秒里有一个事件。它在线上表现为「偶尔少几条,刷新又有了」。
验证办法是主动制造那一瞬间:在取快照的代码前后加延迟,在延迟窗口里触发一个事件,确认它出现在客户端。
给它降级,不要等它。
阻塞等待会让一个慢客户端拖垮整个服务;无限缓冲会让进程在流量高峰时被系统杀掉。两者的共同问题是:它们试图「不丢任何一条」,结果代价由所有人分摊。
正确做法是每个连接一个有界队列,满了就清空它的全部排队事件并发一条 reset,让它自己重新拉快照。注意不要只丢最老的几条——那会在客户端造成一个它察觉不到的空洞。
记成一句话:背压的出口永远是降级,不是等待。
没有标准答案,检查这几件事:
- 游标是什么?严格单调吗?重启后还有效吗?
- 重放窗口多长?超出之后有没有一条明确的
reset? - 快照和增量的接缝顺序对吗?有没有构造一个落在接缝上的事件验证过?
- 背压策略是有界队列加降级吗?慢客户端会影响其他人吗?
- 心跳是双向的吗?两边都有超时判死吗?
- 重连退避有抖动吗?
- 推的是确认深度之内还是之外的事件?之外的话,撤销怎么做?
- 掉线 30 秒的对照实验跑过吗?弱网三档都跑过吗?
最后一条最重要:这些问题里有一半,只有在真实网络条件下才会给出诚实的答案。
一句话带走
实时系统必须可重放:用户掉线再回来,不能丢掉中间发生的事件。