Python中计算观测行非空点斜率的高效实现方案咨询
高效计算带NaN的滞后序列斜率(针对大数据量)
问题场景
面对30万行数据、80个类似price的变量,需要计算每行滞后值的斜率。原始逐行循环的方法在小数据量下可行,但大数据量下速度极慢,需要更高效的向量化实现方案。
原始数据构造与滞后列生成
import pandas as pd import numpy as np df = pd.DataFrame({'date':[1,2,3,4,5,6,7,8], 'price':[4.95, 5.04, 4.88, 4.22, 5.67, 5.89, 5.50, 5.12]}) pd.set_option('display.max_columns', None) for lag in range(1,7): df[f'price_lag{lag}M'] = df['price'].shift(lag) print(df)
输出结果:
date price price_lag1M price_lag2M price_lag3M price_lag4M price_lag5M price_lag6M 0 1 4.95 NaN NaN NaN NaN NaN NaN 1 2 5.04 4.95 NaN NaN NaN NaN NaN 2 3 4.88 5.04 4.95 NaN NaN NaN NaN 3 4 4.22 4.88 5.04 4.95 NaN NaN NaN 4 5 5.67 4.22 4.88 5.04 4.95 NaN NaN 5 6 5.89 5.67 4.22 4.88 5.04 4.95 NaN 6 7 5.50 5.89 5.67 4.22 4.88 5.04 4.95 7 8 5.12 5.50 5.89 5.67 4.22 4.88 5.04
原始低效实现(仅适用于小数据)
逐行过滤非空值后计算斜率,但Python循环在大数据量下性能极差:
vars_to_consider = [f'price_lag{i}M' for i in range(1,7)] for i in range(len(df)): Y = df.loc[i, vars_to_consider].values idx = np.where(~np.isnan(Y))[0] if len(idx) < 2: df.loc[i, 'price_trend_6M'] = np.nan else: df.loc[i, 'price_trend_6M'] = np.polyfit(np.arange(len(idx)), Y[idx], 1)[0].round(4) df = df.drop(vars_to_consider, axis=1) print(df)
输出结果:
date price price_trend_6M 0 1 4.95 NaN 1 2 5.04 NaN 2 3 4.88 -0.0900 3 4 4.22 0.0350 4 5 5.67 0.2350 5 6 5.89 -0.0620 6 7 5.50 -0.1694 7 8 5.12 -0.1937
高效向量化实现
核心思路是利用NumPy向量化运算替代Python循环,结合线性斜率的数学公式(斜率 = 协方差(X,Y)/方差(X)),避免逐行调用np.polyfit。
单变量处理代码
vars_to_consider = [f'price_lag{i}M' for i in range(1,7)] # 定义滞后对应的X坐标(lag1对应1,lag2对应2...) X = np.arange(1,7) X_mat = np.tile(X, (len(df), 1)) # 生成非空值掩码 mask = ~df[vars_to_consider].isna().values # 计算每行有效数据的数量 count = mask.sum(axis=1) # 初始化斜率数组 slopes = np.full(len(df), np.nan) # 仅处理有效数据≥2的行 valid_rows = count >= 2 if valid_rows.any(): # 提取有效X和Y值并重塑形状 X_valid = X_mat[valid_rows][mask[valid_rows]].reshape(-1, count[valid_rows][0]) Y_valid = df[vars_to_consider].values[valid_rows][mask[valid_rows]].reshape(-1, count[valid_rows][0]) # 计算均值 X_mean = X_valid.mean(axis=1, keepdims=True) Y_mean = Y_valid.mean(axis=1, keepdims=True) # 计算协方差与方差 cov = ((X_valid - X_mean) * (Y_valid - Y_mean)).sum(axis=1) var = ((X_valid - X_mean)**2).sum(axis=1) # 计算斜率,避免除以0 slopes[valid_rows] = np.where(var != 0, cov / var, np.nan) # 赋值并保留4位小数 df['price_trend_6M'] = slopes.round(4) df = df.drop(vars_to_consider, axis=1)
批量处理80个变量的封装
将逻辑封装为函数,批量处理所有目标变量:
def calculate_trend(df, var_name, lag_num=6): # 生成滞后列 lag_cols = [f'{var_name}_lag{i}M' for i in range(1, lag_num+1)] for lag in range(1, lag_num+1): df[f'{var_name}_lag{lag}M'] = df[var_name].shift(lag) # 向量化计算斜率 X = np.arange(1, lag_num+1) X_mat = np.tile(X, (len(df), 1)) mask = ~df[lag_cols].isna().values count = mask.sum(axis=1) slopes = np.full(len(df), np.nan) valid_rows = count >= 2 if valid_rows.any(): X_valid = X_mat[valid_rows][mask[valid_rows]].reshape(-1, count[valid_rows][0]) Y_valid = df[lag_cols].values[valid_rows][mask[valid_rows]].reshape(-1, count[valid_rows][0]) X_mean = X_valid.mean(axis=1, keepdims=True) Y_mean = Y_valid.mean(axis=1, keepdims=True) cov = ((X_valid - X_mean) * (Y_valid - Y_mean)).sum(axis=1) var = ((X_valid - X_mean)**2).sum(axis=1) slopes[valid_rows] = np.where(var != 0, cov / var, np.nan) df[f'{var_name}_trend_{lag_num}M'] = slopes.round(4) df = df.drop(lag_cols, axis=1) return df # 替换为你的80个变量列表 vars_list = ['price', 'var1', 'var2', ..., 'var79'] for var in vars_list: df = calculate_trend(df, var)
性能优势说明
- 完全基于NumPy底层的向量化运算,避免Python循环的性能损耗
- 批量处理有效行,减少重复计算逻辑
- 直接用数学公式替代
np.polyfit函数调用,降低额外开销
内容的提问来源于stack exchange,提问作者Tejas
相关产品推荐
相关产品推荐

