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

子线程启动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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.12 23:54:51