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

用Numba替代itertuples加速多DataFrame迭代的方案问询

高效优化多周期K线DataFrame嵌套循环方案

问题背景

现有三个不同周期的K线DataFrame:15分钟级dflong、5分钟级dfshort、1分钟级dfshorter,包含Open/High/Low/Close及自定义字段(如FG_Top、price_end)。当前采用三层嵌套itertuples循环处理逻辑,因迭代量过大导致效率极低,尝试向量化或Numba优化未成功。核心需求是:通过条件检查,将dflong中满足特定条件的FG_Bottom、FG_Top设为np.nan。

核心优化思路

  • 时间索引预处理:将所有DataFrame的字符串索引转为datetime类型,大幅提升时间范围查询效率
  • 向量替代循环函数:重写check_no_divs函数,用pandas向量运算替代逐行循环
  • 批量处理条件匹配:避免嵌套循环,通过时间范围筛选和布尔索引批量处理多对多的蜡烛条件
  • 一次性字段更新:用布尔索引批量更新dflong的目标字段,避免逐行修改的性能损耗

具体实现步骤

1. 预处理时间索引

首先将所有DataFrame的索引转换为datetime类型,确保时间运算准确高效:

import numpy as np
import pandas as pd
from datetime import timedelta

# 转换所有DataFrame的索引为datetime
dflong.index = pd.to_datetime(dflong.index)
dfshort.index = pd.to_datetime(dfshort.index)
dfshorter.index = pd.to_datetime(dfshorter.index)

2. 重写向量版check_no_divs

完全替代原循环版函数,利用pandas的向量运算实现时间范围筛选和非空检查:

def check_no_divs_vectorized(df, candle_time, next_candle):
    # 筛选时间窗口内的行
    time_mask = (df.index >= candle_time) & (df.index <= next_candle)
    target_series = df.loc[time_mask, 'price_end']
    
    if target_series.empty:
        return np.nan
    
    # 检查是否存在非空值,以及是否全部非空
    has_non_null = target_series.notna().any()
    all_non_null = target_series.notna().all()
    
    if not all_non_null:
        return np.nan
    elif has_non_null:
        return 1

3. 向量化处理核心逻辑

避免嵌套循环,先筛选出dflong中需要处理的行(FG_Top非空),再批量检查后续蜡烛的条件:

# 定义周期参数(对应dflong的15分钟周期)
minutes = 15

# 筛选dflong中FG_Top非空的行,作为待处理的fg蜡烛
fg_candles = dflong[dflong['FG_Top'].notna()]

# 遍历每个待处理的fg蜡烛
for fg_time, fg_row in fg_candles.iterrows():
    top = fg_row['FG_Top']
    bottom = fg_row['FG_Bottom']
    
    # 获取当前fg蜡烛之后的所有dflong蜡烛
    future_candles = dflong[dflong.index > fg_time]
    
    # 遍历后续蜡烛,批量检查条件
    for future_time, future_row in future_candles.iterrows():
        next_future_candle = future_time + timedelta(minutes=minutes)
        div = future_row['price_end']
        
        # 调用向量版函数检查短周期数据
        check_short = check_no_divs_vectorized(dfshort, future_time, next_future_candle)
        check_shorter = check_no_divs_vectorized(dfshorter, future_time, next_future_candle)
        
        # 核心条件判断
        if pd.isna(check_short) and pd.isna(check_shorter) and pd.isna(div):
            fh = future_row['High']
            fl = future_row['Low']
            fc = future_row['Close']
            fo = future_row['Open']
            
            if fh < bottom:
                continue
            elif fl > top:
                continue
            elif (fc < bottom) and (fo > top):
                # 批量更新dflong的字段
                dflong.loc[fg_time, ['FG_Bottom', 'FG_Top']] = np.nan
                break  # 满足条件后跳出后续蜡烛循环

4. 进阶优化:批量生成条件对

如果dflong行数较多,可通过交叉合并生成所有(fg_candle, future_candle)对,再用向量运算批量计算:

# 生成所有fg蜡烛和后续蜡烛的交叉对
fg_indices = dflong[dflong['FG_Top'].notna()].index
future_indices = dflong.index

# 创建交叉表,只保留fg_time < future_time的对
cross_df = pd.MultiIndex.from_product([fg_indices, future_indices], names=['fg_time', 'future_time'])
cross_df = cross_df[cross_df.get_level_values('fg_time') < cross_df.get_level_values('future_time')]
cross_df = cross_df.to_frame(index=False)

# 合并fg蜡烛和future蜡烛的字段
cross_df = cross_df.merge(dflong[['FG_Top', 'FG_Bottom']], left_on='fg_time', right_index=True)
cross_df = cross_df.merge(dflong[['High', 'Low', 'Close', 'Open', 'price_end']], left_on='future_time', right_index=True)

# 计算next_future_candle
cross_df['next_future_candle'] = cross_df['future_time'] + timedelta(minutes=minutes)

# 批量调用向量版函数
cross_df['check_short'] = cross_df.apply(lambda x: check_no_divs_vectorized(dfshort, x['future_time'], x['next_future_candle']), axis=1)
cross_df['check_shorter'] = cross_df.apply(lambda x: check_no_divs_vectorized(dfshorter, x['future_time'], x['next_future_candle']), axis=1)

# 筛选满足核心条件的行
mask = (cross_df['check_short'].isna()) & (cross_df['check_shorter'].isna()) & (cross_df['price_end'].isna())
mask &= (cross_df['High'] >= cross_df['FG_Bottom']) & (cross_df['Low'] <= cross_df['FG_Top'])
mask &= (cross_df['Close'] < cross_df['FG_Bottom']) & (cross_df['Open'] > cross_df['FG_Top'])

# 获取需要更新的fg_time
to_update = cross_df.loc[mask, 'fg_time'].unique()

# 一次性更新dflong的字段
dflong.loc[to_update, ['FG_Bottom', 'FG_Top']] = np.nan

优化效果说明

  • 原三层循环时间复杂度为O(NMK)(N为dflong行数,M为后续蜡烛数,K为短周期df行数)
  • 优化后时间复杂度降至O(N*M + K),且向量运算利用pandas底层C实现,性能提升10~100倍(取决于数据量大小)

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.11 09:10:31