Python中如何跨函数使用WebSocket数据 解决run_forever阻塞问题
问题根源
你的判断完全正确,程序阻塞就是ws.run_forever()方法导致的。这个方法是websocket-client库实现的同步阻塞事件循环,启动后会持续占用当前线程的执行权,监听消息、处理回调,直到WebSocket连接断开才会返回,因此后续的业务逻辑代码永远得不到执行机会。
最优实现方案
对于你当前的同步代码写法,给WebSocket连接单独创建独立守护线程是成本最低、最稳定的方案,不需要重构现有消息接收逻辑,主线程可以正常执行数据校验、业务判断等逻辑,两个线程通过共享内存变量交换数据即可,只需要加个简单的线程锁避免读写冲突。
修正后可运行代码示例
import json import websocket import os import threading import time from datetime import datetime # ========== 初始化原代码缺失的全局变量 ========== SOCK = "wss://你的WebSocket接口地址" # 替换成实际的ws地址 compare_coins = "./coin_price.json" # 替换成实际的文件存储路径 min_price = {} # 共享的价格存储字典 # 线程锁,避免ws线程写数据的时候,主线程读数据读到半更新的脏值 data_lock = threading.Lock() def on_open(ws): print("WebSocket连接成功") def on_close(ws, close_status_code, close_msg): print(f"WebSocket连接断开,状态码:{close_status_code},信息:{close_msg}") def on_message(ws,message): global min_price json_message = json.loads(message) cs = json_message # 加锁更新共享数据 with data_lock: for coin in cs: min_price[coin["s"]] = { "lastPrice": coin["c"], "lowPrice": coin["l"]} # 提示:原代码每收到一条消息就重写一次文件,IO开销极高,建议降低写盘频率,或者直接读内存数据即可 with open(compare_coins, "w") as file: json.dump(min_price, file, indent=4) # 注意:原代码这里的candle变量未定义,是残留的K线逻辑,如果当前ws返回的不是K线数据请删除这段,否则会报错 # is_candle_closed = candle['x'] # if is_candle_closed: # print(json.dumps(candle, indent=2)) def check_below(): volatile_coins = {} # 读数据的时候也加锁 with data_lock: # 做一次数据拷贝,避免锁持有时间太长影响ws写入 price_snapshot = min_price.copy() for coin in price_snapshot: if float(price_snapshot[coin]["lastPrice"]) >= float(price_snapshot[coin]["lowPrice"]): volatile_coins[coin] = price_snapshot[coin]["lastPrice"] return volatile_coins, len(volatile_coins), price_snapshot if __name__ == "__main__": ws = websocket.WebSocketApp(SOCK, on_open=on_open,on_close=on_close, on_message=on_message) # 启动守护线程跑ws事件循环,主程序退出时该线程会自动终止,不会残留 ws_thread = threading.Thread(target=ws.run_forever, daemon=True) ws_thread.start() # 主线程正常跑你的业务逻辑,不会被阻塞 while True: # 等待ws连接完成第一次数据拉取 if not min_price: time.sleep(1) continue # 比如每隔3秒执行一次价格校验 volatile_coins, count, price_data = check_below() print(f"当前符合条件的币种数量:{count}") # 这里写后续业务逻辑即可 time.sleep(3)
其他可选方案
如果你不想用多线程,可以替换成异步WebSocket客户端库(比如websockets、aiohttp的ws模块),基于asyncio协程在同一个线程里同时跑WebSocket消息监听和业务逻辑,但这种方案需要把所有相关代码都改成异步写法,改造成本较高,对于当前场景没有必要。
原代码注意事项
- 每收到一条ws消息就全量重写本地json文件的IO开销极高,高消息频率下会占用大量磁盘性能,建议降低写盘频率(比如每10秒写一次、或者数据更新超过10条再写),如果业务逻辑不需要持久化,直接读内存里的
min_price即可。 - 原代码中
is_candle_closed = candle['x']段的candle变量没有在回调函数内定义,属于残留的其他业务逻辑,如果当前对接的WebSocket接口返回的不是K线数据,这段代码运行时会直接抛错,需要根据实际接口返回内容调整。
内容的提问来源于stack exchange,提问作者Dinizov
相关产品推荐
相关产品推荐

