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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.25 10:19:54