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

如何创建固定长度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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.18 19:25:22