子线程启动Binance WebSocket无法持续接收消息问题求助
问题:子线程中启动Binance WebSocket仅接收第一条消息,后续无消息
在子线程中启动Binance WebSocket连接时,即便不读取latest_price_info属性,WebSocket也只能收到第一条消息,后续完全收不到任何消息。
子线程启动的测试代码
import threading import time from binance_websocket import BinanceWebSocket lock = threading.Lock() def manage_websocket(binance_ws): binance_ws.ws.run_forever() if __name__ == "__main__": symbol = "BTCUSDT" interval = "1m" binance_ws = BinanceWebSocket(symbol, interval) websocket_thread = threading.Thread(target=manage_websocket, args=(binance_ws, )) websocket_thread.start() print("thread started") while True: print("while in icine girdi") time.sleep(10)
BinanceWebSocket类实现代码
import websocket import json import time import threading class BinanceWebSocket: def __init__(self, symbol, interval): self.symbol = symbol self.interval = interval self.ws_url = f"wss://stream.binance.com:9443/ws/{self.symbol}@kline_{self.interval}" self.ws = websocket.WebSocketApp(self.ws_url, on_message=self.on_message, on_error=self.on_error, on_close=self.on_close) self.ws.on_open = self.on_open self.latest_price_info = None self.lock = threading.Lock() # Initialize the lock def set_latest_price_info(self, price): with self.lock: self.latest_price_info = price def get_latest_price_info(self): with self.lock: return self.latest_price_info def on_message(self, ws, message): print("The message is: ", message) data = json.loads(message) ''' kline = data['k'] close = kline['c'] self.set_latest_price_info(close) # Use the instance method to set the latest price ''' print("Latest current price updated: ", self.get_latest_price_info()) def on_error(self, ws, error): print(f"Error: {error}") def on_close(self, ws, close_status_code, close_msg): print("WebSocket connection closed") self.reconnect() def on_open(self, ws): print("WebSocket connection opened") subscription_payload = { "method": "SUBSCRIBE", "params": [ f"{self.symbol}@kline_{self.interval}" ], "id": 1 } self.ws.send(json.dumps(subscription_payload)) # Binance disconnects websocket connections every 24h, therefore reconnecting when disconnected def reconnect(self): while True: try: print("Reconnecting...") time.sleep(5) # Delay before reconnecting self.ws = websocket.WebSocketApp(self.ws_url, on_message=self.on_message, on_error=self.on_error, on_close=self.on_close) self.ws.on_open = self.on_open self.ws.run_forever() except Exception as e: print(f"Reconnection failed: {e}")
主线程启动正常工作的代码
from binance_websocket import BinanceWebSocket import time symbol = "btcusdt" interval = "1m" # You can adjust the interval here ws = BinanceWebSocket(symbol, interval) ws.ws.run_forever() try: while True: lastPriceFrom_ws = ws.latest_price_info if lastPriceFrom_ws is not None: print("Latest Current Price:", lastPriceFrom_ws) time.sleep(1) # Adjust the sleep time if needed except KeyboardInterrupt: print("WebSocket connection stopped.")
排查过程与核心需求
- 最初怀疑主线程读取和子线程更新
latest_price_info时发生死锁,在类中添加了锁,但问题未解决。 - 尝试将子线程与主线程
join,依然无效,程序运行时WebSocket仅处于等待状态,无消息接收也无报错。 - 核心需求:通过WebSocket持续获取Binance交易所价格数据,WebSocket在子线程中持续更新
latest_price_info属性,主线程可随时读取该值;因主线程启动WebSocket会阻塞执行,必须采用子线程方案。
内容的提问来源于stack exchange,提问作者202286ege
相关产品推荐
相关产品推荐

