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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.29 07:18:25