如何以类形式使用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
相关产品推荐
相关产品推荐

