如何用CCXT并行获取各交易所同一时刻BTC买卖报价?
并行获取多交易所BTC实时买卖价(解决串行时间差问题)
问题背景
串行调用CCXT的fetchTicker接口时,每个请求都会产生网络延迟,导致不同交易所的报价并非同一时刻的数据。比如原串行代码:
import ccxt binance = ccxt.binance() coinbase = ccxt.coinbase() ftx = ccxt.ftx() kraken = ccxt.kraken() kucoin = ccxt.kucoin() stat_list = binance.fetchTicker('BTC/USDT') print('Binance', stat_list['bid'], stat_list['ask']) stat_list = coinbase.fetchTicker('BTC-USD') print('Coinbase', stat_list['bid'], stat_list['ask']) stat_list = ftx.fetchTicker('BTC/USDT') print('FTX', stat_list['bid'], stat_list['ask']) stat_list = kraken.fetchTicker('BTC/USDT') print('Kraken', stat_list['bid'], stat_list['ask']) stat_list = kucoin.fetchTicker('BTC/USDT') print('Kucoin', stat_list['bid'], stat_list['ask'])
输出的报价存在时间差,无法反映同一时刻的市场状态:
Binance 20451.04 20451.33 Coinbase 20343.33 20552.61 FTX 20451.0 20452.0 Kraken 20454.5 20454.6 Kucoin 20450.9 20451.0
解决方案1:使用线程池(ThreadPoolExecutor)
针对IO密集型的网络请求,线程池可以同时发起多个请求,大幅缩短总耗时,让报价更接近同一时刻。
代码实现
import ccxt from concurrent.futures import ThreadPoolExecutor def fetch_exchange_ticker(exchange_name, exchange_class, symbol): try: exchange = exchange_class() ticker = exchange.fetchTicker(symbol) return (exchange_name, ticker['bid'], ticker['ask']) except Exception as e: return (exchange_name, None, f"Error: {str(e)}") # 定义要查询的交易所和对应交易对 exchanges_to_query = [ ("Binance", ccxt.binance, "BTC/USDT"), ("Coinbase", ccxt.coinbase, "BTC-USD"), ("FTX", ccxt.ftx, "BTC/USDT"), ("Kraken", ccxt.kraken, "BTC/USDT"), ("Kucoin", ccxt.kucoin, "BTC/USDT"), ] # 用线程池并行执行 with ThreadPoolExecutor(max_workers=len(exchanges_to_query)) as executor: results = list(executor.map(lambda x: fetch_exchange_ticker(*x), exchanges_to_query)) # 打印结果 for result in results: name, bid, ask = result if isinstance(ask, str): print(f"{name} {ask}") else: print(f"{name} {bid} {ask}")
说明
- 线程池的
max_workers设置为交易所数量,确保所有请求同时发起 - 加入异常处理,避免单个交易所请求失败导致整个程序崩溃
- 所有请求几乎同时发出,报价时间差可控制在毫秒级
解决方案2:使用CCXT异步接口(推荐)
CCXT原生支持异步IO,使用asyncio配合异步版本的交易所接口,效率比线程池更高,尤其适合大量交易所查询场景。
代码实现
import asyncio import ccxt.async_support as ccxt async def fetch_exchange_ticker_async(exchange_name, exchange_class, symbol): try: async with exchange_class() as exchange: ticker = await exchange.fetchTicker(symbol) return (exchange_name, ticker['bid'], ticker['ask']) except Exception as e: return (exchange_name, None, f"Error: {str(e)}") async def main(): exchanges_to_query = [ ("Binance", ccxt.binance, "BTC/USDT"), ("Coinbase", ccxt.coinbase, "BTC-USD"), ("FTX", ccxt.ftx, "BTC/USDT"), ("Kraken", ccxt.kraken, "BTC/USDT"), ("Kucoin", ccxt.kucoin, "BTC/USDT"), ] # 创建所有异步任务 tasks = [fetch_exchange_ticker_async(*args) for args in exchanges_to_query] # 并行执行所有任务 results = await asyncio.gather(*tasks) # 打印结果 for result in results: name, bid, ask = result if isinstance(ask, str): print(f"{name} {ask}") else: print(f"{name} {bid} {ask}") if __name__ == "__main__": asyncio.run(main())
说明
- 使用
ccxt.async_support导入异步版本的CCXT模块 - 用
async with自动管理交易所连接的生命周期 asyncio.gather同时执行所有异步任务,耗时仅取决于最慢的那个请求- 异步IO无需线程切换开销,性能更优
内容的提问来源于stack exchange,提问作者Will-NotGiveUp
相关产品推荐
相关产品推荐

