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版本核心问题
- 订阅请求发送时机错误:
ws.run_forever()是阻塞调用,后续的ws.send(subscribe_request)永远无法执行,导致从未真正订阅事件,消息处理逻辑混乱。 - 同步阻塞交易查询:直接在消息回调中调用
web3.eth.get_transaction()会阻塞WebSocket消息循环,节点因长时间未收到心跳主动断开连接。 - 空值未校验:
get_transaction可能返回None(交易已打包或被移出Mempool),后续隐式调用len()操作时触发类型错误。
websockets版本核心问题
- 同步API阻塞异步循环:
web3.eth.get_transaction()是同步方法,在异步循环中调用会阻塞事件循环,导致WebSocket ping/pong超时,节点断开连接。 - 异常捕获不完整:仅捕获
TransactionNotFound,未处理get_transaction返回None的场景,也未处理JSON解析失败等异常。 - 心跳配置缺失:未设置
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
相关产品推荐
相关产品推荐

