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

专业程序如何处理WebSocket海量数据?Python队列积压求解

解决WebSocket接收海量交易所数据的消息积压与并行处理方案

当前代码的核心问题

你的代码采用单协程串行处理模式,同时使用multiprocessing.Manager.dict做跨进程数据存储,这两个点是导致消息积压的主要原因:

  • 单协程既要负责WebSocket消息接收,又要同步处理每条消息的JSON解析、字典更新,一旦消息量超过单协程处理能力,队列必然积压;
  • Manager.dict本质是跨进程通信的代理对象,每次读写都要经过IPC(进程间通信),开销远大于普通字典,进一步拖慢处理速度。

专业数据筛选网站的并行处理核心方案

1. 解耦消息接收与处理:生产者-消费者模式

将WebSocket消息的接收和业务处理拆分为独立模块,用队列做缓冲:

  • 生产者:单个/多个协程专门负责从WebSocket接收消息,直接丢入异步队列(asyncio.Queue)或进程队列(multiprocessing.Queue);
  • 消费者:启动多个协程/进程从队列中取消息并行处理,通过调整消费者数量匹配消息产出速度;
  • 进阶优化:按币种哈希值分片,让不同消费者固定处理某一类币种,避免多线程/协程对同一份数据的锁竞争。

2. 替换高开销数据存储

  • 单进程多协程场景:用普通Python字典加asyncio.Lock做线程安全的本地存储,完全避免IPC开销;
  • 多进程场景:使用内存数据库做共享存储,或让每个进程维护自身的币种数据分片,最后按需合并。绝对不要用multiprocessing.Manager的容器类。

3. 提升数据解析与处理效率

  • 用ujson替代标准库json:ujson的解析速度是标准库的3-5倍,能大幅降低消息解析的CPU占用;
  • 批量处理:攒够N条消息后再批量解析、批量更新数据,减少锁操作或IO操作的次数;
  • 跳过无效消息:提前过滤不符合格式的消息,避免无效处理。

4. 多WebSocket连接分流订阅

交易所的单WebSocket连接通常有订阅数量上限,全币种订阅可以拆分到多个连接中:

  • 比如将1000个币种分成10组,每组100个,用10个独立的WebSocket连接分别订阅,每个连接对应一个生产者协程,分散单连接的消息压力。

5. 异步化所有阻塞操作

如果需要做持久化(如写数据库、文件),绝对不要在消费者协程中做同步阻塞操作,而是将这些任务丢到专门的异步任务池(如asyncio.Semaphore控制并发)或进程池处理,避免阻塞消息处理流程。

优化后的代码示例(单进程多协程)

import websockets
import ujson
import asyncio

async def process_message_batch(queue, book_tickers, lock):
    while True:
        # 批量获取消息(最多10条,可调整)
        messages = []
        for _ in range(10):
            try:
                messages.append(await asyncio.wait_for(queue.get(), timeout=0.1))
            except asyncio.TimeoutError:
                break
        if not messages:
            continue
        
        # 批量解析与更新
        batch_updates = {}
        for msg in messages:
            try:
                data = ujson.loads(msg)
                symbol = data['s']
                batch_updates[symbol] = {
                    'best_bid': data['b'],
                    'best_ask': data['a']
                }
            except Exception:
                continue
        
        # 批量更新字典,减少锁竞争
        async with lock:
            book_tickers.update(batch_updates)

async def ws_producer(add_symbols, queue):
    async for websocket in websockets.connect('wss://fstream.binance.com/ws'):
        # 发送订阅请求
        await websocket.send(ujson.dumps({
            "method": "SUBSCRIBE",
            "params": [el.lower() + '@bookTicker' for el in add_symbols],
            "id": 1
        }))
        # 忽略订阅确认消息
        await websocket.recv()
        
        # 持续接收消息并放入队列
        async for message in websocket:
            # 队列满时丢弃旧消息(或根据需求调整策略)
            if queue.full():
                try:
                    queue.get_nowait()
                except asyncio.QueueEmpty:
                    pass
            await queue.put(message)

async def main(add_symbols):
    # 初始化队列(设置上限防止内存溢出)
    queue = asyncio.Queue(maxsize=2000)
    book_tickers = {}
    lock = asyncio.Lock()
    
    # 启动3个消费者协程(根据CPU核心数调整)
    consumer_tasks = [
        asyncio.create_task(process_message_batch(queue, book_tickers, lock))
        for _ in range(3)
    ]
    
    # 启动生产者协程
    producer_task = asyncio.create_task(ws_producer(add_symbols, queue))
    
    # 等待所有任务运行
    await asyncio.gather(producer_task, *consumer_tasks)

if __name__ == "__main__":
    # 示例币种列表,替换为你的全币种列表
    target_symbols = ["BTCUSDT", "ETHUSDT", "BNBUSDT", "ADAUSDT"]
    asyncio.run(main(target_symbols))

内容的提问来源于stack exchange,提问作者Fresto

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.08 00:18:14