You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

ETH内存池待处理交易WSS连接异常(含web3验证报错)

以太坊Mempool交易订阅报错问题及修复方案

问题描述

通过WSS连接节点获取ETH Mempool待处理交易,仅接收消息时运行正常,但添加web3交易验证逻辑后,运行数分钟出现以下错误:

WebSocket库报错

ERROR:root:An error occurred: object of type 'NoneType' has no len()

websockets库报错

Future exception was never retrieved

future: <Future finished exception=ConnectionClosedError(None, None, None)>

websockets.exceptions.ConnectionClosedError: no close frame received or sent

原代码问题分析

WebSocket版本核心问题

  1. 订阅请求发送时机错误:ws.run_forever()是阻塞调用,后续的ws.send(subscribe_request)永远无法执行,导致从未真正订阅事件,消息处理逻辑混乱。
  2. 同步阻塞交易查询:直接在消息回调中调用web3.eth.get_transaction()会阻塞WebSocket消息循环,节点因长时间未收到心跳主动断开连接。
  3. 空值未校验:get_transaction可能返回None(交易已打包或被移出Mempool),后续隐式调用len()操作时触发类型错误。

websockets版本核心问题

  1. 同步API阻塞异步循环:web3.eth.get_transaction()是同步方法,在异步循环中调用会阻塞事件循环,导致WebSocket ping/pong超时,节点断开连接。
  2. 异常捕获不完整:仅捕获TransactionNotFound,未处理get_transaction返回None的场景,也未处理JSON解析失败等异常。
  3. 心跳配置缺失:未设置ping_interval和ping_timeout,无法维持长连接活性。

修复后的代码

WebSocket版本(线程异步处理查询)

import time
import json
import ssl
import logging
import threading
from colored import colored
import websocket
from web3 import Web3
from web3.exceptions import TransactionNotFound

web3 = Web3(Web3.HTTPProvider("http://localhost:7546"))

def get_tx(tx_hash):
    try:
        tx = web3.eth.get_transaction(tx_hash)
        if tx is not None:
            print(tx)
    except TransactionNotFound:
        pass
    except Exception as err:
        logging.error(f"交易查询错误: {err}")

def on_message(ws, message):
    timing = time.perf_counter()
    try:
        if 'eth_subscription' in message:
            tx_data = json.loads(message)
            tx_hash = tx_data['params']['result']
            # 用线程异步处理查询,避免阻塞WebSocket循环
            threading.Thread(target=get_tx, args=(tx_hash,), daemon=True).start()
    except Exception as err:
        logging.error(f"消息处理错误: {err}")
    finally:
        print(colored(str(time.perf_counter() - timing), 'magenta'))

def on_error(ws, error):
    logging.error(f"WebSocket错误: {error}")

def on_close(ws, close_status_code, close_msg):
    print("### 连接已关闭 ###")

def on_open(ws):
    print("连接已建立")
    # 连接成功后立即发送订阅请求
    subscribe_request = {
        "id": 1,
        "method": "eth_subscribe",
        "params": ["newPendingTransactions"],
        "jsonrpc": "2.0"
    }
    ws.send(json.dumps(subscribe_request))

def main():
    try:
        websocket.enableTrace(True, level='INFO')
        ws = websocket.WebSocketApp(
            "ws://localhost:7546",
            on_open=on_open,
            on_message=on_message,
            on_error=on_error,
            on_close=on_close
        )
        ws.run_forever(reconnect=5, sslopt={"cert_reqs": ssl.CERT_NONE})
    except Exception as err:
        logging.error(f"主进程错误: {err}")

if __name__ == '__main__':
    main()

websockets版本(异步Web3客户端)

import asyncio
import json
import logging
from web3 import AsyncWeb3
from web3.exceptions import TransactionNotFound
import websockets

node_ws_url = "ws://localhost:7546"
async_web3 = AsyncWeb3(AsyncWeb3.AsyncHTTPProvider("http://localhost:7546"))

async def get_tx(tx_hash):
    try:
        tx = await async_web3.eth.get_transaction(tx_hash)
        if tx is not None:
            print(tx)
    except TransactionNotFound:
        pass
    except Exception as e:
        logging.error(f"交易查询错误: {e}")

async def get_event():
    async with websockets.connect(node_ws_url, ping_interval=30, ping_timeout=60) as ws:
        subscribe_msg = json.dumps({
            "jsonrpc": "2.0",
            "id": 1,
            "method": "eth_subscribe",
            "params": ["newPendingTransactions"]
        })
        await ws.send(subscribe_msg)
        subscription_response = await ws.recv()
        print(subscription_response)

        while True:
            try:
                message = await ws.recv()
                response = json.loads(message)
                if 'params' in response and 'result' in response['params']:
                    tx_hash = response['params']['result']
                    # 异步执行查询,不阻塞WebSocket循环
                    asyncio.create_task(get_tx(tx_hash))
            except websockets.exceptions.ConnectionClosed as e:
                logging.error(f"连接已关闭: {e}")
                break
            except json.JSONDecodeError as e:
                logging.error(f"无效JSON消息: {e}")
            except Exception as e:
                logging.error(f"未知错误: {e}")

if __name__ == "__main__":
    asyncio.run(get_event())

关键修复点

  • 修正订阅时机:WebSocket版本在on_open回调中发送订阅请求,确保连接建立后立即完成订阅。
  • 异步处理查询:WebSocket版本用线程、websockets版本用异步Web3+asyncio.create_task,避免阻塞消息循环导致心跳超时。
  • 空值校验:对get_transaction返回结果做None判断,避免后续操作触发类型错误。
  • 配置心跳机制:websockets版本添加ping_interval和ping_timeout,维持长连接活性。
  • 完善异常捕获:增加JSON解析错误等场景的捕获,避免单个异常导致整个循环崩溃。

内容的提问来源于stack exchange,提问作者404

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.07.03 15:29:51