Logo
新用户订阅 RPC,首月享 6.5 折优惠查看优惠
OnFinality Learn
RPC 故障排查阅读约 14 分钟

RPC WebSocket 重连不丢数据:缺口检测与回填

重连 WebSocket 很简单,难的是证明你没有漏掉任何数据。了解针对 Ethereum、Solana、Sui 和 Substrate 的基于游标的缺口检测与有界回填。

TL;DR

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,并在重连时用 getSignaturesForAddressgetTransaction 回填。在 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 中心

更新条件:当你的提供商更改最大回填范围、当你添加新链、或当你观察到比保留窗口更深的重组时,重新审视这个设计。

  • 提供商最大回填范围可能强制分块。
  • 游标持久性需要事务性持久化。
  • 超出保留窗口的深度重组需要完整重新同步。
  • 在提供商变更、新增链或深度重组时重新审视。

永远不用担心基础设施

OnFinality 消除了 DevOps 的繁重工作,让您能够更聪明、更快地构建。

开始