如何为各交易所对象创建独立线程实现异步并发请求
解决交易所API请求并发执行的问题
你的代码存在核心问题:每个线程都在重复执行所有交易所的请求,这不仅没实现“每个交易所请求同时执行”的目标,反而造成资源浪费和重复任务。而且asyncio本身就是为单线程I/O并发设计的,ccxta的异步客户端正好适配这种场景,完全不需要额外套多线程。
修正后的代码
import asyncio import time import ccxta async def async_client(exchange): client = getattr(ccxta, exchange)() try: tickers = await client.fetch_tickers() print(f"完成 {exchange} 的数据拉取") return (exchange, tickers) finally: await client.close() async def main(): exchanges = ['binance', 'bitget', 'bitmart', 'bitvavo', 'bybit', 'gate', 'huobi', 'kucoin', 'mexc', 'okx'] tic = time.time() # 并发调度所有交易所的异步请求 results = await asyncio.gather(*[async_client(exchange) for exchange in exchanges], return_exceptions=True) # 整理结果,区分成功与失败请求 successful_results = {exch: data for exch, data in results if not isinstance(data, Exception)} failed_exchanges = [exch for exch, data in results if isinstance(data, Exception)] print(f"总耗时: {time.time() - tic:.2f} 秒") print(f"成功拉取 {len(successful_results)} 个交易所数据") if failed_exchanges: print(f"拉取失败的交易所: {failed_exchanges}") return successful_results if __name__ == '__main__': asyncio.run(main())
关键说明
- 移除冗余多线程:
asyncio.gather会在同一个事件循环里自动调度所有异步任务,实现真正的I/O并发,所有交易所的API请求会同时发起。 - 优化异常处理:保留
return_exceptions=True避免单个请求失败导致整体任务中断,同时在结果处理中明确区分成功和失败的请求,方便后续排查问题。 - 确保资源清理:用
finally块保证客户端连接无论请求成功或失败都会被关闭,避免资源泄漏。
特殊场景下的多线程用法(非必要)
如果你的场景中存在同步阻塞的交易所客户端(比如使用ccxt而非ccxta),可以用asyncio.to_thread将同步调用包装为异步任务,借助线程池处理阻塞操作:
import ccxt async def sync_client_wrapper(exchange): client = getattr(ccxt, exchange)() try: # 将同步请求放入线程池执行,不阻塞事件循环 tickers = await asyncio.to_thread(client.fetch_tickers) return (exchange, tickers) finally: client.close()
内容的提问来源于stack exchange,提问作者Xena
相关产品推荐
相关产品推荐

