专业程序如何处理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
相关产品推荐
相关产品推荐

