Python异步WebSocket订阅函数:如何传递接收数据至程序其他模块
如何让Async WebSocket订阅函数的数据被程序其他部分获取?
我实现了一个订阅WebSocket数据流的async函数,该函数通过Bitfinex交易所的WebSocket订阅BTCUSD和ETHUSD交易对的订单簿频道,当前仅将接收的message数据打印到标准输出。请问如何让这些数据能够被程序的其他部分获取?示例代码如下:
import json import asyncio import websockets async def subscribe(ws_host, subscribe_request): async with websockets.connect(ws_host) as ws: request = json.dumps(subscribe_request) await ws.send(request) while True: try: message = await ws.recv() print(message) except websockets.exceptions.ConnectionClosed: print("Connection was closed") if __name__ == '__main__': ws_host = 'wss://api.bitfinex.com/ws/2' subscribe_request_btc = dict( event='subscribe', channel='book', symbol='tBTCUSD', prec='P0', freq='F1', len='25' ) subscribe_request_eth = dict( event='subscribe', channel='book', symbol='tETHUSD', prec='P0', freq='F1', len='25' ) loop = asyncio.get_event_loop() tasks = [subscribe(ws_host, subscribe_request_btc), subscribe(ws_host, subscribe_request_eth)] loop.run_until_complete(asyncio.wait(tasks)) loop.close()
这里有几个实用的方法可以让你WebSocket订阅到的数据被程序其他部分获取,都是异步环境下安全的方案:
1. 使用asyncio.Queue(推荐)
asyncio.Queue是异步环境中安全传递数据的首选方式,它支持生产者-消费者模式,你的WebSocket订阅函数作为生产者把消息放入队列,程序的其他部分作为消费者从队列中取出数据。
修改后的代码示例:
import json import asyncio import websockets async def subscribe(ws_host, subscribe_request, data_queue): async with websockets.connect(ws_host) as ws: request = json.dumps(subscribe_request) await ws.send(request) while True: try: message = await ws.recv() # 将消息放入队列,附带交易对标识方便区分 await data_queue.put((subscribe_request['symbol'], message)) except websockets.exceptions.ConnectionClosed: print("Connection was closed") break async def process_data(data_queue): """模拟程序其他部分处理数据的任务""" while True: symbol, message = await data_queue.get() print(f"Processing {symbol} data: {message[:100]}...") # 这里可以添加自定义逻辑:解析订单簿、存储到数据库、计算深度等 data_queue.task_done() if __name__ == '__main__': ws_host = 'wss://api.bitfinex.com/ws/2' subscribe_request_btc = dict( event='subscribe', channel='book', symbol='tBTCUSD', prec='P0', freq='F1', len='25' ) subscribe_request_eth = dict( event='subscribe', channel='book', symbol='tETHUSD', prec='P0', freq='F1', len='25' ) # 创建异步队列 data_queue = asyncio.Queue() loop = asyncio.get_event_loop() tasks = [ subscribe(ws_host, subscribe_request_btc, data_queue), subscribe(ws_host, subscribe_request_eth, data_queue), process_data(data_queue) ] loop.run_until_complete(asyncio.gather(*tasks)) loop.close()
2. 使用回调函数
如果你希望数据一到达就触发特定逻辑,可以传入一个回调函数,订阅函数收到消息后直接调用这个函数,把数据传递过去。
示例代码:
import json import asyncio import websockets async def subscribe(ws_host, subscribe_request, callback): async with websockets.connect(ws_host) as ws: request = json.dumps(subscribe_request) await ws.send(request) while True: try: message = await ws.recv() # 调用回调函数,传递交易对和消息数据 await callback(subscribe_request['symbol'], message) except websockets.exceptions.ConnectionClosed: print("Connection was closed") break async def handle_message(symbol, message): """自定义的消息处理回调函数""" print(f"Received {symbol} update: {message[:80]}...") # 这里可以添加业务逻辑:解析增量更新、维护本地订单簿快照等 if __name__ == '__main__': ws_host = 'wss://api.bitfinex.com/ws/2' subscribe_request_btc = dict( event='subscribe', channel='book', symbol='tBTCUSD', prec='P0', freq='F1', len='25' ) subscribe_request_eth = dict( event='subscribe', channel='book', symbol='tETHUSD', prec='P0', freq='F1', len='25' ) loop = asyncio.get_event_loop() tasks = [ subscribe(ws_host, subscribe_request_btc, handle_message), subscribe(ws_host, subscribe_request_eth, handle_message) ] loop.run_until_complete(asyncio.wait(tasks)) loop.close()
3. 自定义数据存储类
如果需要保存最新的订单簿状态,可以创建一个异步安全的类来存储数据,订阅函数更新这个类的属性,其他部分直接访问该类获取最新数据。
示例代码:
import json import asyncio import websockets from dataclasses import dataclass, field from typing import Dict @dataclass class OrderBookStorage: """异步安全的订单簿存储类""" books: Dict[str, any] = field(default_factory=dict) async def update_book(self, symbol, message): # 实际场景建议解析Bitfinex的消息格式,维护结构化的订单簿 # 这里暂时直接存储原始消息作为示例 self.books[symbol] = message def get_latest_book(self, symbol): return self.books.get(symbol, None) async def subscribe(ws_host, subscribe_request, storage): async with websockets.connect(ws_host) as ws: request = json.dumps(subscribe_request) await ws.send(request) while True: try: message = await ws.recv() # 更新存储类中的数据 await storage.update_book(subscribe_request['symbol'], message) except websockets.exceptions.ConnectionClosed: print("Connection was closed") break async def monitor_order_books(storage): """模拟程序其他部分定期获取最新订单簿""" while True: btc_book = storage.get_latest_book('tBTCUSD') eth_book = storage.get_latest_book('tETHUSD') print(f"Latest BTC book: {btc_book[:60]}..." if btc_book else "No BTC data yet") print(f"Latest ETH book: {eth_book[:60]}..." if eth_book else "No ETH data yet") await asyncio.sleep(5) # 每5秒检查一次 if __name__ == '__main__': ws_host = 'wss://api.bitfinex.com/ws/2' subscribe_request_btc = dict( event='subscribe', channel='book', symbol='tBTCUSD', prec='P0', freq='F1', len='25' ) subscribe_request_eth = dict( event='subscribe', channel='book', symbol='tETHUSD', prec='P0', freq='F1', len='25' ) # 创建存储实例 order_book_storage = OrderBookStorage() loop = asyncio.get_event_loop() tasks = [ subscribe(ws_host, subscribe_request_btc, order_book_storage), subscribe(ws_host, subscribe_request_eth, order_book_storage), monitor_order_books(order_book_storage) ] loop.run_until_complete(asyncio.gather(*tasks)) loop.close()
注意事项
- 所有方案都保证了异步环境下的安全性,避免了协程间的数据竞争问题
- Bitfinex的订单簿消息包含初始快照和增量更新两种格式,实际使用时建议解析消息并维护正确的本地订单簿状态
- 如果需要持久化数据,可以在处理逻辑中添加数据库写入操作
内容的提问来源于stack exchange,提问作者Max Mikhaylov
相关产品推荐
相关产品推荐

