Logo
新用户订阅 RPC,首月享 6.5 折优惠查看优惠
OnFinality Learn
网络与协议指南阅读约 12 分钟

以太坊 eth_subscribe:日志与区块头订阅对比轮询过滤器

了解以太坊的 eth_subscribe 推送模型:newHeads、日志、生命周期、重连处理,以及何时应改用轮询。

TL;DR

以太坊的 eth_subscribe 通过 WebSocket 提供基于推送的实时事件流,非常适合低延迟通知,但在断开连接时会丢失事件且无法重放。当需要可审计性或必须在连接中断时不遗漏事件时,轮询过滤器或一次性 eth_getLogs 查询更为合适。

推送模型:eth_subscribe 如何工作

以太坊 JSON-RPC 通过 WebSocket 传输暴露了发布-订阅(pubsub)机制,这在官方的 Ethereum execution-apisgeth 的实时事件 页面中有文档说明。与请求-响应调用不同,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.subscriptionparams.result
  • 终止:eth_unsubscribe(subscriptionId) 停止流;关闭套接字也会结束它。

四个标准订阅通道

以太坊 JSON-RPC 规范定义了四个标准通道。每个通道都有不同的结果负载和用例。

  • newHeads:每当链上添加新的规范区块头时触发。结果是区块头对象(类似于 eth_getBlockByNumber 且交易参数为 false)。请注意,在链重组期间,您可能会收到后来成为孤块的区块头;您的客户端必须通过比较区块哈希并必要时回滚状态来处理重组。
  • logs:每当有日志匹配您提供的过滤器时触发。过滤器语法与 eth_getLogs 相同:您可以指定 address(单个地址或数组)和 topics(一个数组,其中每个位置可以是 null 表示任意,或一个替代主题值的数组)。结果是包含 addresstopicsdatablockNumbertransactionHashlogIndex 等的日志对象。
  • 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+1currentHead 查询 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

后续步骤与进一步阅读

既然您了解了推送模型,您可以将其应用于您自己的应用程序。要深入了解相关主题,请探索以下资源:

永远不用担心基础设施

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

开始