以太坊的 eth_subscribe 通过 WebSocket 提供基于推送的实时事件流,非常适合低延迟通知,但在断开连接时会丢失事件且无法重放。当需要可审计性或必须在连接中断时不遗漏事件时,轮询过滤器或一次性 eth_getLogs 查询更为合适。
推送模型:eth_subscribe 如何工作
以太坊 JSON-RPC 通过 WebSocket 传输暴露了发布-订阅(pubsub)机制,这在官方的 Ethereum execution-apis 和 geth 的实时事件 页面中有文档说明。与请求-响应调用不同,eth_subscribe 打开一个持久的推送通道:客户端发送订阅请求,节点回复一个订阅 ID(例如 0x9cef478923ff2bf24f5f1e1a1b3f8f7b),此后每当订阅的事件发生时,节点会发送通知对象。
每个通知都是一个 JSON-RPC 消息,包含 jsonrpc: "2.0"、method: "eth_subscription" 以及一个包含 subscription ID 和 result 负载的 params 对象。客户端不会为每个事件发送请求;它只是从套接字读取消息。这与轮询有根本的不同,轮询中客户端反复向节点请求新数据。
订阅与 WebSocket 连接绑定。如果连接断开,服务器端的订阅将被销毁(或变得不可达)。这是大多数实际订阅失败的根源,也决定了我们稍后将介绍的重连策略。
- 传输:WebSocket(WS 或 WSS)。
eth_subscribe不能通过普通 HTTP 使用。 - 发起:
eth_subscribe(channel, params)返回一个订阅 ID。 - 通知:服务器推送
eth_subscription消息,包含params.subscription和params.result。 - 终止:
eth_unsubscribe(subscriptionId)停止流;关闭套接字也会结束它。
四个标准订阅通道
以太坊 JSON-RPC 规范定义了四个标准通道。每个通道都有不同的结果负载和用例。
newHeads:每当链上添加新的规范区块头时触发。结果是区块头对象(类似于eth_getBlockByNumber且交易参数为false)。请注意,在链重组期间,您可能会收到后来成为孤块的区块头;您的客户端必须通过比较区块哈希并必要时回滚状态来处理重组。logs:每当有日志匹配您提供的过滤器时触发。过滤器语法与eth_getLogs相同:您可以指定address(单个地址或数组)和topics(一个数组,其中每个位置可以是null表示任意,或一个替代主题值的数组)。结果是包含address、topics、data、blockNumber、transactionHash、logIndex等的日志对象。newPendingTransactions:当节点交易池中添加新的待处理交易时触发。结果取决于节点配置(例如 geth 的--rpc.evmtimeout或类似--ws.fulltx的标志),可以是交易哈希(默认)或完整交易对象。此通道对于内存池监控很有用,但可能会很嘈杂。syncing:当节点的同步状态改变时触发。结果是同步对象(类似于eth_sync),当节点完全同步时为false。对于节点健康监控很有用。
订阅生命周期与重连问题
订阅的生命周期很简单:创建它,接收事件,最终销毁它。但典型的失败是 WebSocket 断开。当连接断开时,服务器端的订阅就消失了。如果您的客户端天真地重连并继续从旧的订阅 ID 读取,您将什么也收不到。更糟的是,如果您不重新订阅,您将在间隙期间静默地错过事件。
健壮的模式是:重连时,始终创建新的订阅,然后回填从您处理的最后一个区块到当前头部之间错过的事件。这就是为什么生产系统通常将 logs 订阅与从最后看到的区块开始的定期 eth_getLogs 查询配对。订阅提供低延迟通知,而轮询回填填补了断开后的间隙。
此外,您必须处理套接字的异步特性。使用专用的读取循环(例如 Go 中的 goroutine,或 Python 中的异步任务),以便解析和处理通知永远不会阻塞新订阅请求或心跳的发送。许多 WebSocket 库会缓冲消息,但如果您阻塞读取循环,您可能会错过消息或导致背压。
- 始终将 WebSocket 断开视为“重新订阅”——切勿重用订阅 ID。
- 跟踪最后处理的区块(来自日志的
blockNumber或区块头的number)以启用回填。 - 重连后,再次调用
eth_subscribe,然后从lastBlock+1到currentHead查询eth_getLogs以填补间隙。 - 使用异步读取循环以避免阻塞套接字。
- 当您真正完成订阅时,调用
eth_unsubscribe以避免泄漏服务器端资源。服务器通常限制每个连接的同时订阅数量(有文档记录 / 因提供商而异)。
可复现示例:订阅 newHeads 和日志
以下 Python 示例使用 websockets 库连接到公共以太坊节点(例如 Cloudflare-ETH,或您自己的端点)。它订阅 newHeads 和特定合约地址和主题的 logs,打印订阅 ID 和前几个通知,然后取消订阅。
- 预期输出:两个订阅 ID(十六进制字符串),然后是一系列通知对象。
newHeads结果包含区块头;logs结果包含日志条目。 - 如果您没有看到任何日志通知,您的过滤器可能不匹配最近的事件,或者节点可能落后于链头。
import asyncio
import json
import websockets
WS_URL = "wss://cloudflare-eth.com" # replace with your provider's WSS endpoint
async def main():
async with websockets.connect(WS_URL) as ws:
# Subscribe to newHeads
await ws.send(json.dumps({"jsonrpc": "2.0", "id": 1, "method": "eth_subscribe", "params": ["newHeads"]}))
resp = json.loads(await ws.recv())
head_sub = resp["result"]
print("newHeads subscription:", head_sub)
# Subscribe to logs for a contract (e.g., USDC Transfer event)
# address: USDC contract, topics: Transfer event signature
filter_params = {
"address": "0xA0b86991c6218b36c1d19D4a2e9Eb0cE3606eB48",
"topics": ["0xddf252ad1be2c89b69c2b068fc378daa952ba7f163c4a11628f55a4df523b3ef"]
}
await ws.send(json.dumps({"jsonrpc": "2.0", "id": 2, "method": "eth_subscribe", "params": ["logs", filter_params]}))
resp = json.loads(await ws.recv())
log_sub = resp["result"]
print("logs subscription:", log_sub)
# Read a few notifications
for _ in range(3):
msg = json.loads(await ws.recv())
if msg.get("method") == "eth_subscription":
print("Notification:", msg["params"]["subscription"], msg["params"]["result"])
# Unsubscribe
await ws.send(json.dumps({"jsonrpc": "2.0", "id": 3, "method": "eth_unsubscribe", "params": [head_sub]}))
resp = json.loads(await ws.recv())
print("Unsubscribed head:", resp)
await ws.send(json.dumps({"jsonrpc": "2.0", "id": 4, "method": "eth_unsubscribe", "params": [log_sub]}))
resp = json.loads(await ws.recv())
print("Unsubscribed logs:", resp)
asyncio.run(main())轮询过滤器与推送订阅:决策表
在推送订阅和轮询过滤器之间进行选择取决于您对持久性、间隙处理、传输和成本的要求。下表总结了权衡。
- 持久性:推送订阅是短暂的;如果连接断开,您将丢失流。轮询过滤器(通过
eth_newFilter+eth_getFilterChanges)在节点上保持状态,但该状态也会在超时后过期(在 geth 中通常为 5 分钟,有文档记录 / 因提供商而异)。 - 间隙处理:使用推送,您必须在重连后回填。使用轮询,您可以从特定区块范围查询,但必须管理过滤器的生命周期。
- 传输:推送需要 WebSocket;轮询也可以通过 HTTP 工作。
- 成本:推送对于高频事件是高效的,因为您不需要发送重复请求。如果轮询过于频繁,可能会浪费资源,但它更简单且更可审计。
- 可审计性:使用
eth_getLogs从已知区块进行轮询可提供可验证的历史记录。推送订阅是最后消息胜出,并且在断开连接时会丢失,因此不适合作为可审计的事件源。
故障排除:为什么我没有收到事件?
当您的订阅似乎匹配但您什么也没收到时,请检查以下清单。
- 主题语法错误:确保您的
topics数组使用null作为通配符,使用数组表示 OR 条件。例如,["0x...", null]匹配任何第二个主题,而[["0x...", "0x..."]]匹配第一个位置的两个主题中的任何一个。 - 承诺/头部滞后:节点可能落后于链头。检查
eth_syncing;如果它返回同步对象,请等待直到它为false。 - 通过 HTTPS 而不是 WSS 订阅:
eth_subscribe仅适用于 WebSocket。如果您使用 HTTPS 端点,您将收到错误或没有事件。 - 节点未暴露 pubsub:某些提供商或私有节点禁用 pubsub。请查阅节点文档或尝试公共 WSS 端点。
- 静默丢弃的订阅:如果 WebSocket 连接静默断开(例如由于网络超时),您的订阅将消失。实现心跳或重连逻辑。
- 区块 gas/日志修剪:轻端点或历史记录有限的提供商可能会修剪日志。这有文档记录 / 因提供商而异。如果您需要历史日志,请使用带区块范围的
eth_getLogs,但要注意提供商限制。
局限性与权衡:何时不使用订阅
推送订阅非常适合实时仪表板、内存池监控和低延迟重要的事件驱动应用程序。然而,它们本质上有丢失性:如果您的客户端离线,您将错过事件,并且没有重放机制。服务器不会为断开的客户端排队消息。
对于需要完整、可审计事件历史记录的应用程序(例如索引、会计或法律合规),请依赖带有明确区块范围的 eth_getLogs。轮询过滤器可以是一个中间选择,但它们也有节点端状态和过期限制。
始终设计您的系统以优雅地处理重连。一种常见模式是运行订阅以获取实时更新,并运行定期回填作业,从最后处理的区块到当前头部查询 eth_getLogs。这确保即使订阅中断几秒钟也不会错过任何事件。
还要考虑成本:每个订阅都会消耗服务器资源。如果您有许多客户端,您可能会达到提供商对并发订阅的限制(有文档记录 / 因提供商而异)。完成后请及时使用 eth_unsubscribe。
后续步骤与进一步阅读
既然您了解了推送模型,您可以将其应用于您自己的应用程序。要深入了解相关主题,请探索以下资源:
- 处理以太坊 WebSocket 断开连接 – 实用的重连策略。
- 使用 eth_getLogs 和主题过滤器进行一次性日志查询 – 用于回填和历史查询。
- 监控 RPC 端点和节点健康 – 跟踪同步状态和延迟。
- Base WebSocket RPC 指南 – Base 上的类似模式。
- 选择以太坊 RPC 节点(RPC 助手) – 选择支持 pubsub 的提供商。
- 以太坊网络概述 和 OnFinality 学习中心 获取更多指南。