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

Python中Asyncio多API并发访问及ETH WebSocket监听卡顿问题

问题根源分析
  1. 同步HTTP调用阻塞异步事件循环
    你使用的web3.eth.get_transaction是基于同步HTTPProvider的方法,即使包裹在async函数里,它仍然会阻塞整个异步事件循环直到HTTP请求完成。当每秒有数十笔交易时,每一笔交易的同步请求都会让WebSocket消息处理暂停,导致消息堆积,最终程序卡顿。

  2. 串行处理交易导致瓶颈
    你在获取每个交易哈希后立即await交易详情请求,意味着必须等前一笔交易的请求完成才能处理下一笔WebSocket消息。这种串行处理方式无法应对高频的NFT交易流量,很快就会形成请求队列,拖慢整个程序。

  3. 错误处理过于宽泛
    代码中的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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.16 03:20:15