Python滑动窗口协整测试循环优化:提升函数运行效率
滑动窗口协整交易对筛选函数性能优化方案
问题描述
我计划在滑动窗口场景下开展实验,每次试验处理一批数据并通过协整测试筛选交易对,但当前find_cointegrated_pairs(data)函数运行耗时过长,求优化方案。
原代码依赖与数据获取
依赖导入
import numpy as np import pandas as pd import yfinance as yf from statsmodels.tsa.vector_ar.vecm import coint_johansen from itertools import combinations # from statsmodels.regression.linear_model import OLS
标的数据获取
df1 = pd.read_html('https://en.wikipedia.org/wiki/List_of_S%26P_500_companies')[0] df = yf.download(df1.Symbol.to_list()) df = df["Close"]
原性能瓶颈函数
def find_cointegrated_pairs(data): data = data.replace(0, np.nan) data = pd.DataFrame(data).dropna(axis=1) combo = [] columns = data.columns for combination in combinations(columns, 2): combo.append(combination) n = len(combo) keys = combo pairs = [] for i in range(n): df = pd.concat([data[keys[i][0]].diff(), data[keys[i][1]].diff()], axis=1).reset_index(drop=True).astype(float) df = df.dropna(axis=0) result = coint_johansen(df, 1, k_ar_diff=1) tracevalues = result.lr1 critical_values = result.cvt if (tracevalues > critical_values[:, 1]).all(): model = data[[keys[i][1], keys[i][0]]].fillna(0).astype(float).cov() prams = model.iloc[0,1]/model.iloc[0,0] spread = data[keys[i][1]] - prams * data[keys[i][0]] signal = (spread - spread.mean())/spread.std() out = pd.concat([pd.Series(keys[i][1]), pd.Series(keys[i][0]), pd.Series(signal.iloc[-1])], axis=1) pairs.append(out) return pd.concat(pairs) find_cointegrated_pairs(df)
针对性优化方案
1. 提前全局预处理,避免重复计算
原函数每次循环都对单对数据做差分、缺失值处理,可一次性完成全量预处理:
def find_cointegrated_pairs_optimized(data): # 一次性完成全局清洗与差分 data_clean = data.replace(0, np.nan).dropna(axis=1) data_diff = data_clean.diff().dropna(axis=0) # 提前完成差分+去缺失 combo = list(combinations(data_clean.columns, 2)) # 简化组合生成 pairs = [] for pair in combo: # 直接从预处理数据取列,避免重复拼接与转换 pair_data = data_diff[[pair[0], pair[1]]] result = coint_johansen(pair_data, 1, k_ar_diff=1) if (result.lr1 > result.cvt[:, 1]).all(): # 省略不必要的astype(float)(原始数据已为数值型) cov_matrix = data_clean[[pair[1], pair[0]]].cov() prams = cov_matrix.iloc[0,1]/cov_matrix.iloc[0,0] spread = data_clean[pair[1]] - prams * data_clean[pair[0]] signal = (spread - spread.mean())/spread.std() # 用列表拼接替代Series,减少对象创建开销 pairs.append([pair[1], pair[0], signal.iloc[-1]]) # 一次性转换为DataFrame,避免多次concat开销 return pd.DataFrame(pairs, columns=['Asset1', 'Asset2', 'Signal'])
2. 并行化处理独立计算任务
每个交易对的协整测试是完全独立的,利用多核CPU并行执行可大幅缩短耗时:
from concurrent.futures import ProcessPoolExecutor def process_single_pair(pair, data_clean, data_diff): pair_data = data_diff[[pair[0], pair[1]]] result = coint_johansen(pair_data, 1, k_ar_diff=1) if (result.lr1 > result.cvt[:,1]).all(): cov_matrix = data_clean[[pair[1], pair[0]]].cov() prams = cov_matrix.iloc[0,1]/cov_matrix.iloc[0,0] spread = data_clean[pair[1]] - prams * data_clean[pair[0]] signal = (spread - spread.mean())/spread.std() return [pair[1], pair[0], signal.iloc[-1]] return None def find_cointegrated_pairs_parallel(data): data_clean = data.replace(0, np.nan).dropna(axis=1) data_diff = data_clean.diff().dropna(axis=0) combo = list(combinations(data_clean.columns, 2)) # 多进程并行处理 with ProcessPoolExecutor() as executor: results = executor.map(process_single_pair, combo, [data_clean]*len(combo), [data_diff]*len(combo)) # 过滤无效结果,一次性生成DataFrame valid_results = [res for res in results if res is not None] return pd.DataFrame(valid_results, columns=['Asset1', 'Asset2', 'Signal'])
3. 前置过滤减少协整测试量
通过相关性筛选提前排除大概率不协整的组合(比如仅保留相关系数绝对值>0.8的对):
def find_cointegrated_pairs_filtered(data): data_clean = data.replace(0, np.nan).dropna(axis=1) data_diff = data_clean.diff().dropna(axis=0) # 前置相关性过滤 corr_matrix = data_clean.corr() high_corr_pairs = [] cols = corr_matrix.columns for i in range(len(cols)): for j in range(i+1, len(cols)): if abs(corr_matrix.iloc[i,j]) > 0.8: high_corr_pairs.append((cols[i], cols[j])) # 仅对高相关组合做协整测试 pairs = [] for pair in high_corr_pairs: pair_data = data_diff[[pair[0], pair[1]]] result = coint_johansen(pair_data, 1, k_ar_diff=1) if (result.lr1 > result.cvt[:,1]).all(): cov_matrix = data_clean[[pair[1], pair[0]]].cov() prams = cov_matrix.iloc[0,1]/cov_matrix.iloc[0,0] spread = data_clean[pair[1]] - prams * data_clean[pair[0]] signal = (spread - spread.mean())/spread.std() pairs.append([pair[1], pair[0], signal.iloc[-1]]) return pd.DataFrame(pairs, columns=['Asset1', 'Asset2', 'Signal'])
内容的提问来源于stack exchange,提问作者Waqas Jamil
相关产品推荐
相关产品推荐

