asyncio搭配websockets时传递接收变量的更优实现方案
问题场景
通过websockets搭配asyncio从Bitstamp流式获取BTC/USD订单簿价格数据时,最初采用全局变量反复覆写的方式存储最新卖一价,供其他业务函数调用。该方案可维护性差,需要更合理的跨协程数据传递实现。
原有实现代码:
import asyncio import websockets import json bitstamp_asks = [] async def bitstamp_connect(): global bitstamp_asks uri = "wss://ws.bitstamp.net/" subscription = { "event": "bts:subscribe", "data": { "channel": "order_book_btcusd" } } async with websockets.connect(uri) as ws: await ws.send(json.dumps(subscription)) subscription_response = await ws.recv() print(json.loads(subscription_response)['event']) while True: response = await ws.recv() bitstamp_data = json.loads(response)['data'] bitstamp_asks = float(bitstamp_data['asks'][0][0]) print(bitstamp_asks) await asyncio.sleep(0.0001) async def store_data(): while True: print(f'The stored value is: {bitstamp_asks}') await asyncio.sleep(0.1) asyncio.get_event_loop().run_until_complete(asyncio.wait([bitstamp_connect(),store_data()]))
优化方案
全局变量的问题在于状态完全暴露,后续扩展多交易对、多数据源时很容易出现变量污染,也不方便做数据校验、异常处理。针对这种「单生产者、多消费者,只需要保留最新值」的流式数据场景,有两种非常稳妥的实现方式:
方案1:封装类管理连接与共享状态(最轻量,适合仅需要读最新值的场景)
asyncio本身是单线程事件循环模型,只要不在更新共享值的代码片段中插入await,就不会出现读写竞态。把连接逻辑、数据存储都封装到类里,共享值作为实例属性,后续所有要读数据的函数都直接调用类实例的属性即可,完全不需要全局变量。
改进代码:
import asyncio import websockets import json from typing import Optional class BitstampStream: def __init__(self, trading_pair: str = "btcusd"): self.trading_pair = trading_pair self.latest_ask: Optional[float] = None self.latest_bid: Optional[float] = None self._ws: Optional[websockets.WebSocketClientProtocol] = None self._subscription = { "event": "bts:subscribe", "data": {"channel": f"order_book_{trading_pair}"} } async def connect(self): uri = "wss://ws.bitstamp.net/" async with websockets.connect(uri) as self._ws: await self._ws.send(json.dumps(self._subscription)) sub_resp = await self._ws.recv() print(f"订阅{self.trading_pair}结果:{json.loads(sub_resp)['event']}") while True: resp = await self._ws.recv() data = json.loads(resp)["data"] # 更新值的过程没有await,不会被事件循环打断,无竞态 self.latest_ask = float(data["asks"][0][0]) self.latest_bid = float(data["bids"][0][0]) async def store_data(stream: BitstampStream): while True: if stream.latest_ask is not None: print(f'最新卖一价: {stream.latest_ask}, 最新买一价: {stream.latest_bid}') await asyncio.sleep(0.1) async def main(): btc_stream = BitstampStream("btcusd") # 后续要加eth数据流直接再实例化即可:eth_stream = BitstampStream("ethusd") await asyncio.gather( btc_stream.connect(), store_data(btc_stream) ) if __name__ == "__main__": asyncio.run(main())
方案2:用asyncio.Queue做生产者-消费者解耦(适合需要消费每一条更新、或者多消费者独立处理数据的场景)
如果后续需要对每一条推送的订单簿数据做处理(比如存数据库、算指标),直接用asyncio自带的Queue作为数据管道,生产者收到数据就往队列里塞,所有消费者独立从队列取数据即可。如果只需要最新值、不需要历史数据,每次塞新数据前把队列里的旧值清空就行,避免队列堆积无效的过期价格。
核心实现片段:
# 初始化时创建队列,设置maxsize=1实现单槽,只保留最新值 self.data_queue = asyncio.Queue(maxsize=1) # 生产者收到数据后更新队列 while True: resp = await self._ws.recv() data = json.loads(resp)["data"] latest_data = { "ask": float(data["asks"][0][0]), "bid": float(data["bids"][0][0]), "timestamp": data["timestamp"] } # 队列满了就先把旧值取出来丢掉,保证永远存的是最新的 if self.data_queue.full(): try: self.data_queue.get_nowait() except asyncio.QueueEmpty: pass await self.data_queue.put(latest_data) # 消费者从队列读数据即可 async def store_data(queue: asyncio.Queue): while True: latest = await queue.get() print(f"收到最新价格:{latest}") queue.task_done()
额外注意点
- 原有代码循环中添加的
await asyncio.sleep(0.0001)完全多余,await ws.recv()本身会在无新消息时自动让出事件循环控制权,额外加sleep只会平白增加数据延迟。 - Python 3.7+版本不要再用旧的
get_event_loop().run_until_complete写法,直接用asyncio.run()即可,框架会自动处理事件循环的创建与回收。 - 后续如果需要加断线重连逻辑,直接在
connect方法外层套异常捕获循环即可,类封装的结构下新增这类逻辑的成本远低于全局变量写法。
内容的提问来源于stack exchange,提问作者spinosaurus7
相关产品推荐
相关产品推荐

