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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.30 05:01:05