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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.09.27 10:36:04