Python使用WebSocket实现计数器累加及数据存入DataFrame问题求助
问题解答
现有代码的核心问题
- 计数变量作用域错误:
count定义在on_message内部,每次触发回调都会被重置为0,无法实现累计统计消息次数的需求。 - 消息处理逻辑错误:
message是单条推送的字符串数据,直接遍历会逐个字符计数,完全不符合统计消息条数统计的需求。 - JSON解析逻辑缺失:收到的原始消息是JSON字符串,未做反序列化无法提取结构化数据存入DataFrame。
实现思路合理性说明
你的思路完全可行,不过需要注意:推荐先将收到的消息和对应计数先暂存到普通Python列表中,最后统一转DataFrame,避免频繁对DataFrame做追加操作,DataFrame是不可变结构,频繁追加性能极低。
修正后可运行代码
import websocket import json import pandas as pd # 全局变量存储数据和计数器 msg_count = 0 data_list = [] # 订阅参数 sub_param = json.dumps({'op': 'subscribe', 'channel': 'trades', 'market': 'BTC-PERP'}) def on_open(wsapp): wsapp.send(sub_param) def on_message(wsapp, message): global msg_count, data_list # 累计计数,每收到一次消息计数+1 msg_count += 1 # 解析JSON消息 msg_data = json.loads(message) # 只处理实际交易数据,过滤心跳等非业务消息 if msg_data.get('type') == 'update' and 'data' in msg_data: for trade in msg_data['data']: # 追加计数和交易数据到列表 trade['msg_seq'] = msg_count data_list.append(trade) # 测试打印当前计数 print(f"当前累计收到消息数:{msg_count},累计交易数据条数:{len(data_list)}") # 可自定义停止条件,比如收到100条消息就停止并导出数据 if msg_count >= 100: wsapp.close() # 列表转DataFrame df = pd.DataFrame(data_list) print(df.head()) # 可导出到本地CSV # df.to_csv('ftx_trades.csv', index=False) def on_error(wsapp, error): print(error) wsapp = websocket.WebSocketApp("wss://ftx.com/ws/", on_message=on_message, on_open=on_open, on_error=on_error) wsapp.run_forever()
代码说明
- 用全局变量
msg_count做累计计数,每触发一次on_message计数加1,对应单次响应计数准确。 - 收到消息先做JSON反序列化,只提取业务交易数据,避免无效数据存入。
- 数据先存在列表中,收够指定条数后统一转DataFrame,性能更高也更灵活。
内容的提问来源于stack exchange,提问作者Jordan
相关产品推荐
相关产品推荐

