如何创建固定长度DataFrame:追加新行同时删除首行(Binance K线场景)
固定容量DataFrame实现方案(对接Binance实时K线)
针对你需要维护一个永不增长、固定行数的DataFrame,每次追加最新Binance K线数据时自动删除最旧行的需求,以下是两种高效实现方案:
方案1:简洁通用版(基于pandas.concat + tail)
适合大多数场景,代码简洁易维护,无需复杂索引管理:
步骤1:初始化配置与空DataFrame
import pandas as pd from binance import ThreadedWebsocketManager # 定义要保留的最大K线数量(根据你的存储需求调整) MAX_KLINE_COUNT = 100 # 初始化DataFrame,字段与Binance K线返回数据一一对应 kline_columns = [ 'timestamp', 'open', 'high', 'low', 'close', 'volume', 'close_time', 'quote_volume', 'trade_count', 'taker_buy_base', 'taker_buy_quote', 'ignore' ] df = pd.DataFrame(columns=kline_columns)
步骤2:实时K线处理函数
def process_kline(msg): global df kline_data = msg['k'] # 仅处理已闭合的完整K线(可选,根据需求调整) if not kline_data['x']: return # 解析新K线为DataFrame行 new_kline = pd.DataFrame({ 'timestamp': [kline_data['t']], 'open': [float(kline_data['o'])], 'high': [float(kline_data['h'])], 'low': [float(kline_data['l'])], 'close': [float(kline_data['c'])], 'volume': [float(kline_data['v'])], 'close_time': [kline_data['T']], 'quote_volume': [float(kline_data['q'])], 'trade_count': [int(kline_data['n'])], 'taker_buy_base': [float(kline_data['V'])], 'taker_buy_quote': [float(kline_data['Q'])], 'ignore': [kline_data['B']] }) # 追加新行并保留最后MAX_KLINE_COUNT行(自动删除最旧行) df = pd.concat([df, new_kline], ignore_index=True).tail(MAX_KLINE_COUNT) # 验证输出(可选) print(f"当前K线数量: {len(df)}") print("最新K线:\n", df.iloc[-1])
步骤3:启动Binance实时数据流
# 启动Binance WebSocket(公共K线无需API密钥) ws_manager = ThreadedWebsocketManager() ws_manager.start_kline_socket( callback=process_kline, symbol='BTCUSDT', # 目标交易对 interval='1m' # K线周期 ) ws_manager.join()
方案2:性能优化版(预分配空间+循环覆盖)
适合高频K线场景(如1秒级),避免频繁DataFrame拼接,性能更优:
步骤1:预分配固定大小DataFrame
import pandas as pd from binance import ThreadedWebsocketManager MAX_KLINE_COUNT = 100 kline_columns = [ 'timestamp', 'open', 'high', 'low', 'close', 'volume', 'close_time', 'quote_volume', 'trade_count', 'taker_buy_base', 'taker_buy_quote', 'ignore' ] # 预分配MAX_KLINE_COUNT行的DataFrame,用NaN填充 df = pd.DataFrame(index=range(MAX_KLINE_COUNT), columns=kline_columns) current_index = 0 # 当前要覆盖的行索引
步骤2:高性能K线处理函数
def process_kline_optimized(msg): global df, current_index kline_data = msg['k'] if not kline_data['x']: return # 直接覆盖当前索引的行,无需拼接 df.loc[current_index] = [ kline_data['t'], float(kline_data['o']), float(kline_data['h']), float(kline_data['l']), float(kline_data['c']), float(kline_data['v']), kline_data['T'], float(kline_data['q']), int(kline_data['n']), float(kline_data['V']), float(kline_data['Q']), kline_data['B'] ] # 更新索引,循环覆盖最旧行 current_index = (current_index + 1) % MAX_KLINE_COUNT # 获取有效K线(过滤未填充的NaN行) valid_df = df.dropna() print(f"当前有效K线数量: {len(valid_df)}") print("最新K线:\n", valid_df.iloc[-1])
步骤3:启动数据流(同方案1)
ws_manager = ThreadedWebsocketManager() ws_manager.start_kline_socket( callback=process_kline_optimized, symbol='BTCUSDT', interval='1s' ) ws_manager.join()
关键注意事项
- 完整K线判断:Binance WebSocket会实时推送K线的更新数据,通过
kline_data['x']可以判断是否为闭合的完整K线,按需选择是否处理中间更新 - 数据类型转换:Binance返回的数值均为字符串,需转换为float/int类型才能进行后续计算
- 存储控制:两种方案都能严格控制DataFrame的最大行数,确保存储消耗稳定在低水平
内容的提问来源于stack exchange,提问作者harish sarwade
相关产品推荐
相关产品推荐

