基于asyncio+websockets的加密货币流数据处理阻塞问题求助
解决方案与库推荐
核心解决思路
把数据流采集和套利计算拆分为独立的异步任务,用asyncio.Queue做低延迟数据中转——别担心队列延迟,asyncio队列是单进程异步实现,开销远低于进程间通信,完全满足加密货币套利的时效性要求。
具体实现方案
1. 异步任务拆分
- 每个交易所的WebSocket流作为独立异步任务,仅负责接收、快速解析数据,然后将数据推入队列。
- 单独开一个套利计算异步任务,从队列中异步获取数据,实时更新行情缓存并执行套利判断。
2. 代码示例
import asyncio import websockets import ujson # 用ujson提升解析速度 BINANCE_WS_URL = "wss://stream.binance.com:9443/ws/btcusdt@trade" OKX_WS_URL = "wss://ws.okx.com:8443/ws/v5/public" ARBITRAGE_THRESHOLD = 0.1 # 自定义套利阈值 def parse_binance_msg(msg): data = ujson.loads(msg) return {"price": float(data["p"]), "timestamp": data["T"]} def parse_okx_msg(msg): data = ujson.loads(msg) if data["event"] == "trade": return {"price": float(data["data"][0]["px"]), "timestamp": data["data"][0]["ts"]} return None async def binance_ws_worker(queue: asyncio.Queue): async with websockets.connect(BINANCE_WS_URL) as ws: async for msg in ws: parsed_data = parse_binance_msg(msg) if parsed_data: await queue.put(("binance", parsed_data)) async def okx_ws_worker(queue: asyncio.Queue): subscribe_msg = ujson.dumps({ "op": "subscribe", "args": [{"channel": "trade", "instId": "BTC-USDT"}] }) async with websockets.connect(OKX_WS_URL) as ws: await ws.send(subscribe_msg) async for msg in ws: parsed_data = parse_okx_msg(msg) if parsed_data: await queue.put(("okx", parsed_data)) async def arbitrage_calculator(queue: asyncio.Queue): price_cache = {"binance": None, "okx": None} while True: exchange, data = await queue.get() price_cache[exchange] = data["price"] # 双端行情齐全时计算套利 if price_cache["binance"] and price_cache["okx"]: spread = abs(price_cache["binance"] - price_cache["okx"]) if spread > ARBITRAGE_THRESHOLD: # 异步执行套利逻辑,避免阻塞事件循环 await execute_arbitrage(price_cache) queue.task_done() async def execute_arbitrage(prices): # 这里实现你的套利下单逻辑,必须用异步HTTP请求(如aiohttp) print(f"发现套利机会:Binance {prices['binance']}, OKX {prices['okx']}") async def main(): # 设置队列上限,做背压防止内存溢出 queue = asyncio.Queue(maxsize=100) tasks = [ asyncio.create_task(binance_ws_worker(queue)), asyncio.create_task(okx_ws_worker(queue)), asyncio.create_task(arbitrage_calculator(queue)), ] await asyncio.gather(*tasks) if __name__ == "__main__": asyncio.run(main())
3. 关键优化点
- 用
ujson替代标准库json,提升行情数据解析速度,减少延迟。 - 队列设置
maxsize,当套利计算跟不上数据流时,自动阻塞WebSocket任务的put操作,避免内存泄漏。 - 套利逻辑必须异步实现,禁止使用同步HTTP请求,否则会阻塞整个事件循环。
推荐WebSocket库
- websockets:你当前使用的库,官方维护、稳定性高,异步支持完善,文档清晰,适合专注WebSocket数据流的场景。
- aiohttp:同时支持异步HTTP和WebSocket,如果你需要在套利逻辑中调用交易所API下单,用aiohttp可以统一技术栈,减少依赖。
内容的提问来源于stack exchange,提问作者Jimmy Router
相关产品推荐
相关产品推荐

