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

如何将WebSocket代码改写为异步/多线程?解决无实时输出问题

问题描述

原同步WebSocket代码(用于获取AAX实时行情):

import websocket
import json


STREAM_HOST = 'wss://realtime.aax.com/marketdata/v2/'


def on_open(ws):
    ws.send(json.dumps({"e": "subscribe", "stream": ["BTCUSDT@book_50", "BTCUSDT@trade", "tickers"]}))


def on_message(ws, message):
    print(message)


def on_error(ws, error):
    print(error)


def on_close(ws):
    print("Connection closed")


if __name__ == "__main__":
    websocket.enableTrace(True)
    ws = websocket.WebSocketApp(STREAM_HOST,
                                on_open=on_open,
                                on_message=on_message,
                                on_error=on_error,
                                on_close=on_close)
    ws.run_forever(ping_interval=1)

尝试的异步代码(无输出且立即结束):

import json
import websockets
import asyncio


STREAM_HOST = 'wss://realtime.aax.com/marketdata/v2/'
SUBSCRIBE = {"e": "subscribe", "stream": ["BTCUSDT@book_50", "BTCUSDT@trade", "tickers"]}


async def hello():
    async with websockets.connect(STREAM_HOST) as websocket:
        await websocket.send(json.dumps(SUBSCRIBE))
        await websocket.recv()


asyncio.run(hello())

请问这段异步代码遗漏了什么?如何编写正确的异步/多线程版本来获取并显示实时数据?


问题原因

你的异步代码只调用了一次await websocket.recv(),它只会接收一条消息(通常是订阅成功的确认消息),之后函数执行完毕,async with块结束,WebSocket连接关闭,程序自然退出。而原同步代码是run_forever()持续监听消息,所以需要在异步版本里持续循环接收消息。

另外,AAX的WebSocket服务需要定期发送心跳包维持连接,原同步代码设置了ping_interval=1,异步版本也需要处理心跳。


正确的异步实现
import json
import websockets
import asyncio

STREAM_HOST = 'wss://realtime.aax.com/marketdata/v2/'
SUBSCRIBE = {"e": "subscribe", "stream": ["BTCUSDT@book_50", "BTCUSDT@trade", "tickers"]}

async def listen_to_stream():
    async with websockets.connect(STREAM_HOST, ping_interval=1) as websocket:
        # 发送订阅请求
        await websocket.send(json.dumps(SUBSCRIBE))
        print("已发送订阅请求")
        
        # 持续循环接收消息
        while True:
            try:
                message = await websocket.recv()
                print(message)
            except websockets.exceptions.ConnectionClosed:
                print("连接已关闭,尝试重连...")
                break
            except Exception as e:
                print(f"发生错误: {e}")
                break

if __name__ == "__main__":
    asyncio.run(listen_to_stream())

说明

  • 使用while True循环持续调用await websocket.recv(),保证一直接收实时消息
  • 在websockets.connect()中设置ping_interval=1,和原同步代码一致,定期发送心跳维持连接
  • 增加异常捕获,处理连接关闭或其他错误情况

多线程实现

如果你想用多线程版本,可以基于原websocket库,把WebSocket的运行放到子线程中,主线程可以做其他操作:

import websocket
import json
import threading

STREAM_HOST = 'wss://realtime.aax.com/marketdata/v2/'

def on_open(ws):
    ws.send(json.dumps({"e": "subscribe", "stream": ["BTCUSDT@book_50", "BTCUSDT@trade", "tickers"]}))
    print("已发送订阅请求")

def on_message(ws, message):
    print(message)

def on_error(ws, error):
    print(f"发生错误: {error}")

def on_close(ws):
    print("连接已关闭")

def run_websocket():
    websocket.enableTrace(True)
    ws = websocket.WebSocketApp(STREAM_HOST,
                                on_open=on_open,
                                on_message=on_message,
                                on_error=on_error,
                                on_close=on_close)
    ws.run_forever(ping_interval=1)

if __name__ == "__main__":
    # 启动WebSocket子线程
    ws_thread = threading.Thread(target=run_websocket)
    ws_thread.daemon = True  # 设置为守护线程,主线程退出时子线程也退出
    ws_thread.start()
    
    # 主线程可以在这里执行其他任务,比如等待用户输入退出
    input("按回车键退出程序...\n")

说明

  • 把原run_forever()的逻辑放到子线程中,避免阻塞主线程
  • 设置子线程为守护线程,保证主线程退出时子线程也会终止
  • 主线程可以添加其他业务逻辑,比如处理用户输入、记录数据等

内容的提问来源于stack exchange,提问作者Nia Grace

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.21 06:18:29