如何将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
相关产品推荐
相关产品推荐

