Python中Asyncio多API并发访问及ETH WebSocket监听卡顿问题
问题根源分析
同步HTTP调用阻塞异步事件循环
你使用的web3.eth.get_transaction是基于同步HTTPProvider的方法,即使包裹在async函数里,它仍然会阻塞整个异步事件循环直到HTTP请求完成。当每秒有数十笔交易时,每一笔交易的同步请求都会让WebSocket消息处理暂停,导致消息堆积,最终程序卡顿。串行处理交易导致瓶颈
你在获取每个交易哈希后立即await交易详情请求,意味着必须等前一笔交易的请求完成才能处理下一笔WebSocket消息。这种串行处理方式无法应对高频的NFT交易流量,很快就会形成请求队列,拖慢整个程序。错误处理过于宽泛
代码中的except:会捕获所有异常(包括键盘中断、超时等),且没有任何日志输出,这会导致你无法发现请求失败、超时等潜在问题,进一步加剧程序的不稳定。
修复方案
1. 切换到异步Web3 Provider
使用AsyncWeb3和AsyncHTTPProvider替代同步版本,让交易详情请求变为非阻塞的异步调用:
from web3 import AsyncWeb3, AsyncHTTPProvider import asyncio import json from websockets import connect alchemy_http_url = 'https://eth-mainnet.g.alchemy.com/v2/gPXApmtTzATb-g4sFtpEXuL5-8i9uaqF' alchemy_ws_url = 'wss://eth-mainnet.g.alchemy.com/v2/gPXApmtTzATb-g4sFtpEXuL5-8i9uaqF' # 初始化异步Web3实例 async_w3 = AsyncWeb3(AsyncHTTPProvider(alchemy_http_url)) options721 = {'topics': [async_w3.sha3(text='Transfer(address,address,uint256)').hex()]} request_721 = {"jsonrpc":"2.0", "id": 1, "method": "eth_subscribe", "params": ["logs", options721]} request_string_721 = json.dumps(request_721)
2. 并发处理交易请求
使用asyncio.Semaphore限制并发请求数量(避免触发Alchemy的速率限制),并通过asyncio.create_task将交易详情请求作为独立任务并行处理,不阻塞WebSocket消息监听:
async def get_transaction(tx_hash, semaphore): async with semaphore: try: tx = await async_w3.eth.get_transaction(tx_hash) return tx['hash'] except Exception as e: print(f"获取交易失败 {tx_hash}: {str(e)}") return None async def process_transaction(tx_hash, semaphore): result = await asyncio.wait_for(get_transaction(tx_hash, semaphore), timeout=60) if result: print(result) async def get_event_721(): # 限制最大并发请求数(根据Alchemy配额调整,比如10) semaphore = asyncio.Semaphore(10) async with connect(alchemy_ws_url) as ws: await ws.send(request_string_721) while True: try: message = await asyncio.wait_for(ws.recv(), timeout=60) event = json.loads(message) if len(event['params']['result']['topics']) == 4: tx_hash = event['params']['result']['transactionHash'] print(tx_hash) # 创建异步任务处理交易,不阻塞WebSocket监听 asyncio.create_task(process_transaction(tx_hash, semaphore)) except asyncio.TimeoutError: print("WebSocket接收超时,尝试保持连接") except Exception as e: print(f"WebSocket处理错误: {str(e)}") asyncio.run(get_event_721())
3. 优化错误处理
移除宽泛的except:,改为捕获特定异常并输出日志,方便排查问题。
关键改进点
- 异步HTTP请求不会阻塞事件循环,WebSocket可以持续接收新消息
- 并发处理交易请求,大幅提升处理效率
- 并发限制避免触发API速率限制
- 明确的错误日志帮助定位问题
内容的提问来源于stack exchange,提问作者Saihhold Chiu
相关产品推荐
相关产品推荐

