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

如何以类形式使用websocket-client与Thread并跨文件调用获取实时结果

问题分析与解决方案

原代码无法获取实时流计算结果,核心问题出在类方法调用错误、返回值逻辑错误以及实时数据同步时机这几点,以下是具体修复步骤:

1. 修复类方法调用错误

在on_message方法中,调用类内的df_import和best_conin时必须通过self引用,否则会被当成全局函数调用,触发NameError。

原错误代码:

def on_message(self, ws, message):
    out = json.loads(message)
    df_import(out)  # 错误:未通过self调用实例方法
    print('============')
    best_conin()    # 错误:未通过self调用实例方法

修复后:

def on_message(self, ws, message):
    out = json.loads(message)
    self.df_import(out)
    print('============')
    self.best_conin()

2. 修复best_conin方法的返回值逻辑

原方法中tt = print(tables[k],v),而print()函数返回值为None,导致整个方法最终返回无效值。需改为返回计算后的币种收益列表。

原错误代码片段:

ls = []
for k, v in bv.items():
    ls.append([tables[k],v])
    tt = print(tables[k],v)

return tt  # 返回的是print的None值

修复后:

ls = []
for k, v in bv.items():
    coin_info = [tables[k], round(v*100, 2)]  # 转成百分比更直观
    ls.append(coin_info)
    print(f"{tables[k]}: {round(v*100, 2)}%")  # 仅用于控制台输出

return ls  # 返回实际计算结果列表

3. 修复DataFrame列赋值的潜在问题

直接用df_.c = df_.c.astype(float)可能因列名c与DataFrame内置属性冲突报错,改为显式列索引访问:

原错误代码:

df_.c = df_.c.astype(float)

修复后:

df_['c'] = df_['c'].astype(float)

4. 优化实时结果的获取方式

原调用代码在实例化后立即调用best_conin,此时websocket尚未收到数据,数据库为空,无法得到有效结果。推荐两种优化方式:

方式一:延迟调用(测试用)

修改调用代码,等待数据流入后再获取结果:

from videos_algovibes.cotacao import CotarMoedas
import time

if __name__ == "__main__":
    rst = CotarMoedas()
    time.sleep(5)  # 等待5秒让websocket获取数据
    vai = rst.best_conin()
    print("实时计算结果:", vai)

方式二:回调机制(生产推荐)

在类中添加回调函数支持,每次计算完成后主动通知外部:

修改CotarMoedas类:

class CotarMoedas:
    def __init__(self, result_callback=None):
        self.endpoint = 'wss://stream.binance.com:9443/ws/!miniTicker@arr'
        self.result_callback = result_callback  # 保存外部回调函数
        tr = Thread(target=self.call_ws, daemon=True)  # 设置为守护线程,随主程序退出
        tr.start()

    # ... 其他方法省略 ...

    def best_conin(self):
        engine = create_engine('sqlite:///COINS.db')
        dfs = pd.read_sql("""SELECT name FROM sqlite_master WHERE type='table'""", engine)
        tables = dfs.name.to_list()
        returns = []

        for table in tables:
            df__ = pd.read_sql(table, engine)
            if len(df__) < 2:  # 数据不足无法计算收益率,跳过
                returns.append(0)
                continue
            ret_ = (df__['c'].pct_change() + 1).prod() - 1
            returns.append(ret_)

        rst = pd.Series(returns).nlargest(10)
        rst = rst.sort_values(ascending=False)  # 明确降序排序
        bv = rst.to_dict()
        ls = []
        for k, v in bv.items():
            coin_info = [tables[k], round(v*100, 2)]
            ls.append(coin_info)
            print(f"{tables[k]}: {round(v*100, 2)}%")
        
        # 如果有回调函数,触发回调返回结果
        if self.result_callback:
            self.result_callback(ls)
        
        return ls

修改调用代码,传入回调函数:

from videos_algovibes.cotacao import CotarMoedas

def handle_result(result):
    print("\n实时更新的Top10币种收益:")
    for coin, pct in result:
        print(f"{coin}: {pct}%")

if __name__ == "__main__":
    rst = CotarMoedas(result_callback=handle_result)
    # 保持主程序运行,等待websocket消息
    try:
        while True:
            pass
    except KeyboardInterrupt:
        print("程序已终止")

5. 其他优化点

  • 在call_ws中设置线程为守护线程,避免主程序退出后线程残留
  • 在best_conin中添加数据长度判断,避免数据不足时计算出错
  • 修复拼写错误:reutrns改为returns(原代码笔误)

内容的提问来源于stack exchange,提问作者marreco

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.07 23:50:41