求助:从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
相关产品推荐
相关产品推荐

