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

基于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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.18 06:52:42