WebSocket 订阅是一条没有游标的消防水管:服务器只推送新事件,从不重放你错过的事件,而关闭帧也不会告诉你丢失了多少条消息。正确的做法是把每条通知的区块号(以及用于重组安全的区块哈希)视为单调检查点,持久化最后看到的游标,并在重连时对缺口执行有界回填查询,然后才重新信任实时流。由于回填和实时投递可能重叠,你必须按 (交易哈希, 日志索引) 或 (区块号, 日志索引) 去重,并在回填高水位线到达之前缓冲实时事件。同一模式可泛化到 Ethereum、Solana、Sui 和 Substrate,只是原语不同。本文给出可运行的 TypeScript 实现、验证循环、故障模式表,以及与纯轮询的权衡。
为什么重连后的订阅会静默丢事件
WebSocket 订阅是一条没有游标的消防水管。根据 Ethereum JSON-RPC PubSub 规范,带 logs 参数的 eth_subscribe 会返回一个订阅 ID,然后在新区块被挖出时推送通知。服务器不会重放你套接字断开期间发生的事件,关闭帧也不会告诉你丢失了多少条消息。因此,天真的重连会产生一条带空洞的流。
这与请求/响应式 RPC 调用有本质区别:后者的失败会以错误形式显现。订阅失败则是不可见的:你的处理程序只是停止接收,然后在某个更晚的区块恢复。如果你在索引转账、清算或治理事件,带空洞的日志流比没有流更糟,因为下游消费者默认数据是完整的。
断开连接的原因本身已在 为什么 WebSocket RPC 会断开以及如何修复 中介绍;本文假设你已经能重连,并聚焦于更难的问题:证明你没有漏掉任何东西,并确定性地恢复缺口。
- 订阅只推送新事件;不存在服务器端重放缓冲区。
- 关闭帧不携带序列号或丢失消息计数。
- 重连而不做对账 = 静默缺口。
- 对正确性敏感的消费者必须把订阅视为提示,而非账本。
游标心智模型:把区块号当作单调检查点
把每条通知的区块号视为单调检查点。把你已完整处理的最高区块号持久化为 lastCursor。对于重组安全的设计,还要持久化区块哈希,以便检测会使游标失效的链重组,遵循 EIP-1898 描述的区块身份约定以及 EIP-234 的区块号语义。
重连时,规则是:在回填 [lastCursor + 1 .. latest] 之前,不要信任实时流。回填是对该范围的有界 eth_getLogs 查询。只有在回填完成且你把 lastCursor 推进到回填高水位线之后,才应恢复消费实时通知。
这把不可靠的推送通道变成了可靠的拉取加推送管道。订阅给你低延迟;回填给你完整性。两者单独都不够。
- 持久化
lastCursor(区块号),可选持久化lastHash(区块哈希)。 - 回填范围为
[lastCursor + 1 .. latest],以latest为界。 - 只有在回填被持久处理后,才推进游标。
- 重组安全:如果
lastHash不再匹配链,则回退游标。
顺序风险:在回填完成前缓冲实时事件
一个微妙的故障模式:重新订阅后,实时流可能开始投递超前于回填高水位线的区块。如果你立即处理实时事件,可能会在处理区块 95(回填)之前先处理区块 100(实时),导致状态乱序和重复计数。
规则是在回填完成前缓冲实时事件。把传入的通知收集到队列中,运行回填,然后把队列与回填结果合并、去重,并按区块顺序处理。只有在队列排空且游标推进后,才切换到直接实时处理。
这个缓冲窗口很短(几秒),但它决定了索引器是正确的,还是会偶尔重复计算一笔转账。无论你在 Ethereum 还是任何其他链上,同样的纪律都适用。
- 实时流可能开始超前于回填高水位线。
- 回填期间把实时事件缓冲到队列中。
- 合并、去重、按区块号排序,然后处理。
- 只有在队列排空后才切换到直接实时处理。
按 (交易哈希, 日志索引) 或 (区块号, 日志索引) 去重
回填和实时投递可能重叠。一条在区块 100 实时到达的日志,如果回填范围延伸到区块 100,也可能出现在你的 eth_getLogs 回填结果中。不去重的话,你会处理它两次。
Ethereum 日志的规范去重键是 (transactionHash, logIndex)。对于没有日志索引的链,使用 (blockNumber, eventIndex) 或确定性事件 ID。维护一个短期去重集合(或数据库中的唯一约束),至少覆盖回填窗口加安全余量。
去重不是可选项。它是让推送与拉取之间的重叠变得安全的机制。如果你跳过它,你的回填反而会引入它本要防止的重复计数。
- Ethereum:去重键
(transactionHash, logIndex)。 - 通用:
(blockNumber, eventIndex)或确定性事件 ID。 - 去重集合的保留时间至少与回填窗口一样长。
- 数据库唯一约束是最稳健的去重方式。
可运行的 TypeScript:订阅、检测重连、回填、去重、恢复
以下示例使用 ethers v6 和持久化游标。它订阅日志、持久化 lastBlock、检测重连、运行有界 getLogs 回填、去重,然后才恢复实时处理。把端点替换为你的提供商的 WebSocket URL;空闲超时和最大回填范围等提供商特定行为因提供商而异。
注意 buffer 数组:回填期间到达的实时事件会被排队,而不是被处理。seen 集合按 (transactionHash, logIndex) 去重。backfill 函数以 latest 为界,并且应当感知速率限制。
import { ethers } from "ethers";
const WS_URL = process.env.WS_URL!;
const ADDRESS = process.env.ADDRESS!; // contract to watch
const MAX_RANGE = 2000; // bounded backfill window
let lastBlock = Number(process.env.LAST_BLOCK ?? 0);
let backfilling = false;
const buffer: ethers.Log[] = [];
const seen = new Set<string>();
function key(l: ethers.Log) {
return `${l.transactionHash}:${l.index}`;
}
async function processLog(l: ethers.Log) {
const k = key(l);
if (seen.has(k)) return;
seen.add(k);
// TODO: persist to your store
console.log("processed", k, "block", l.blockNumber);
if (l.blockNumber > lastBlock) lastBlock = l.blockNumber;
}
async function backfill(provider: ethers.Provider) {
backfilling = true;
const latest = await provider.getBlockNumber();
let from = lastBlock + 1;
while (from <= latest) {
const to = Math.min(from + MAX_RANGE - 1, latest);
const logs = await provider.getLogs({ address: ADDRESS, fromBlock: from, toBlock: to });
logs.sort((a, b) => a.blockNumber - b.blockNumber || a.index - b.index);
for (const l of logs) await processLog(l);
from = to + 1;
}
// drain buffered live events
buffer.sort((a, b) => a.blockNumber - b.blockNumber || a.index - b.index);
for (const l of buffer) await processLog(l);
buffer.length = 0;
backfilling = false;
}
async function main() {
const provider = new ethers.WebSocketProvider(WS_URL);
provider.on("error", () => {});
provider.websocket.on("close", async () => {
console.warn("socket closed; reconnecting");
await backfill(provider);
});
provider.on({ address: ADDRESS }, async (l: ethers.Log) => {
if (backfilling) buffer.push(l);
else await processLog(l);
});
await backfill(provider); // initial catch-up
}
main().catch(console.error);验证循环:杀掉套接字、注入事件、断言恰好一次恢复
未经测试的缺口恢复设计不可信。构建一个验证循环,故意杀掉套接字,在中断期间注入一个已知事件,并断言该事件被恰好一次恢复。这是证明你的回填和去重逻辑真正有效的唯一方法。
循环步骤:(1) 启动订阅者并记录 lastBlock;(2) 强制关闭 WebSocket;(3) 在断开期间发送一笔会发出已知事件的交易;(4) 允许重连和回填运行;(5) 断言该事件在你的存储中恰好出现一次,且 lastBlock 已推进越过它。
针对你自己的端点运行这个循环,并把结果记录到表格中。不要依赖供应商公布的延迟或可靠性数字;测量你自己的。
- 以编程方式强制关闭套接字(例如
provider.websocket.close())。 - 在中断窗口内注入一个已知事件。
- 断言恰好一次恢复以及游标推进。
- 在空闲超时、提供商部署和负载均衡器重置等场景下重复测试。
结果表:针对你自己的端点测量缺口恢复
使用下表记录你自己的测量结果。用你的端点和工作负载的结果填写它。提供商特定数字因提供商而异,因此你自己的测量是唯一可靠的指南。
每个场景至少运行十次,并记录最坏情况,而不是平均值。缺口窗口才是关键:如果你的回填范围超过提供商的最大 eth_getLogs 范围,你必须分块。
- 场景 | 断开原因 | 缺口(区块) | 回填时间(毫秒) | 恢复事件数 | 重复数 | 通过/失败
- 空闲超时 | N 分钟无流量 | | | | |
- 提供商部署 | 服务器端重启 | | | | |
- 负载均衡器重置 | 连接被断开 | | | | |
- 强制关闭 | 客户端主动杀掉 | | | | |
- 重组 | 链重组 | | | | |
模式泛化:Solana、Sui 和 Substrate
同样的纪律适用于不同原语的各条链。在 Solana 上,使用基于 slot 的游标:订阅程序或账户,持久化最后处理的 slot,并在重连时用 getSignaturesForAddress 和 getTransaction 回填。在 Sui 上,使用 checkpoint 游标和 suix_queryEvents 回填缺口。在 Substrate 上,使用区块号游标和每个区块的 system_events。
原语不同,但心智模型完全相同:单调游标、有界回填、去重、在回填完成前缓冲实时事件。如果你理解了 Ethereum 的情况,你就理解了所有这些。
关于这些链的端点选择和故障转移,请参阅 RPC 端点指南 和 RPC 节点监控、指标与故障转移。
- Solana:slot 游标 +
getSignaturesForAddress回填。 - Sui:checkpoint 游标 +
suix_queryEvents回填。 - Substrate:区块号游标 +
system_events回填。 - 同样的纪律,不同的原语。
限制回填范围:速率限制、退避与缺口窗口
缺口窗口很重要。空闲超时、提供商部署和负载均衡器重置可能产生从几秒到几分钟的缺口。你的回填范围必须有界且感知速率限制,否则会触达提供商限制并恢复失败。
对 eth_getLogs 调用分块(例如每次 2000 个区块),并在速率限制错误上应用指数退避。退避与你的重连逻辑相互作用:如果你重连过于激进,可能会重连到被限流的状态并再次失败。退避模式请参阅 RPC 超时错误:原因与修复。
如果你的缺口超过提供商的最大回填范围,你必须分块。如果超过你的保留窗口,你必须回退到完整重新同步。在需要之前就了解你的限制。
- 对回填调用分块,以遵守提供商的最大范围。
- 在速率限制错误上应用指数退避。
- 不要激进地重连到被限流的状态。
- 了解你的保留窗口;超过时回退到完整重新同步。
WebSocket 加回填 vs 纯轮询:为正确性做选择
纯轮询(定时执行 eth_getLogs)更简单,并且如果你持久化游标,它天生无缺口,但它增加延迟,在高频率下可能更昂贵。WebSocket 加回填给你低延迟加完整性,代价是代码更复杂。
当延迟重要且你能正确实现去重和缓冲时,选择 WebSocket 加回填。当正确性至关重要且延迟可容忍,或者你的提供商的 WebSocket 可靠性不确定时,选择纯轮询。对比详见 eth_subscribe 日志 vs 轮询过滤器。
混合方案通常最好:用 WebSocket 做低延迟提示,加上周期性轮询对账作为安全网。这能捕获你的重连逻辑漏掉的任何缺口。
- 轮询:更简单,配合游标无缺口,延迟更高。
- WebSocket 加回填:低延迟,更复杂。
- 混合:WebSocket 提示 + 周期性轮询对账。
- 根据延迟容忍度和正确性要求来选择。
故障模式与排查
下表列出常见故障模式及其修复方法。当你的缺口恢复没有按预期工作时使用它。
如果看到重复,说明你的去重键错误或去重集合太短。如果看到事件缺失,说明你的回填范围错误或游标推进过早。如果看到乱序处理,说明你没有在回填期间缓冲实时事件。
- 重复 | 去重键错误或集合太短 | 使用
(txHash, logIndex),延长集合。 - 事件缺失 | 回填范围错误或游标过早推进 | 验证
[lastCursor+1 .. latest],回填后再推进。 - 乱序 | 实时事件未缓冲 | 回填期间缓冲,按区块排序。
- 被限流 | 回填过于激进 | 分块调用,指数退避。
- 重组损坏 | 没有区块哈希检查 | 持久化
lastHash,不匹配时回退。 - 静默缺口 | 完全没有回填 | 实现游标 + 回填。
局限、权衡与后续步骤
这个模式有局限。它假设你的提供商支持在缺口范围上执行 eth_getLogs;有些提供商会限制范围或结果数量。它假设你的游标是持久的;如果你的进程在处理和持久化之间崩溃,你可能会重复处理或跳过。它假设没有超出保留窗口的深度重组;更深的重组需要完整重新同步。
权衡是复杂性:你是在构建一个小型对账引擎,而不仅仅是一个订阅者。对于正确性至关重要的消费者,这种复杂性是合理的。对于低风险的仪表盘,纯轮询可能就足够了。
后续步骤:针对你自己的端点实现验证循环,填写结果表,并查看 RPC 定价 和 API 服务 以了解成本影响。更广泛的概览请参阅 OnFinality Learn 中心。
更新条件:当你的提供商更改最大回填范围、当你添加新链、或当你观察到比保留窗口更深的重组时,重新审视这个设计。
- 提供商最大回填范围可能强制分块。
- 游标持久性需要事务性持久化。
- 超出保留窗口的深度重组需要完整重新同步。
- 在提供商变更、新增链或深度重组时重新审视。