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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.21 03:48:41