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

Python并行化优化咨询:多周期计算未提速及CPU利用率低问题

Python并行化未提速问题排查与解决方案

问题描述

我正在尝试对以下Python代码进行并行化处理,但对Python并行化技术并不熟悉。目标是为每个i迭代启动并行进程——单个i迭代耗时约10分钟,需处理247个周期。

原始串行代码如下:

periodMisIdx_CTrad = []
periodADF_CTrad = []

for i in range(numPeriods):    
    firstMisIdx = []
    lowestADF = []

    for j in range(periodNumAssets[i]):
        init = normRets_np[i][0][:,Q_trad_new[i][j]]
        form = normRets_np[i][1][:,Q_trad_new[i][j]]
        trade = normRets_np[i][2][:,Q_trad_new[i][j]]
            
        # init       
        initCopC, _ = fitCvines(init, structures)
            
        # form
        misIdxForm = compute_mispricing_index(form,initCopC)
        adftests = []
        for k in range(len(structures)):
            temp = adfuller(misIdxForm[k])
            adftests.append(temp[0])
        minadf = min(adftests)
        structNum = adftests.index(minadf)
        formCopC = [pv.Vinecop(form, structure = pv.CVineStructure(structures[structNum]))]
        misIdxForm = misIdxForm[structNum]
        
        # trade
        misIdxTrade = np.squeeze(compute_mispricing_index(trade, formCopC))
            
        firstMisIdx.append( [misIdxForm, misIdxTrade] )
        lowestADF.append(minadf)
    
    periodMisIdx_CTrad.append(firstMisIdx)
    periodADF_CTrad.append(lowestADF)

我尝试了以下并行化代码,但运行2个周期仍耗时约20分钟,完全未提速:

numPeriods = 2

from multiprocessing.pool import ThreadPool as Pool

def process_period(i):
    firstMisIdx = []
    lowestADF = []

    for j in range(periodNumAssets[i]):

        init = normRets_np[i][0][:,Q_trad_new[i][j]]
        form = normRets_np[i][1][:,Q_trad_new[i][j]]
        trade = normRets_np[i][2][:,Q_trad_new[i][j]]
            
        # init       
        initCopC, _ = fitCvines(init, structures)
            
        # form
        misIdxForm = compute_mispricing_index(form,initCopC)
        adftests = []
        for k in range(len(structures)):
            temp = adfuller(misIdxForm[k])
            adftests.append(temp[0])
        minadf = min(adftests)
        structNum = adftests.index(minadf)
        formCopC = [pv.Vinecop(form, structure = pv.CVineStructure(structures[structNum]))]
        misIdxForm = misIdxForm[structNum]
        
        # trade
        misIdxTrade = np.squeeze(compute_mispricing_index(trade, formCopC))
            
        firstMisIdx.append( [misIdxForm, misIdxTrade] )
        lowestADF.append(minadf)

    return firstMisIdx, lowestADF

if __name__ == '__main__':
    pool = Pool()
    

    results = pool.map(process_period, range(numPeriods))


    periodMisIdx_CTrad = [result[0] for result in results]
    periodADF_CTrad = [result[1] for result in results]

当前CPU利用率仅8-12%,硬件为12核CPU、24GB GPU内存,我是否忽略了可提速的错误?


解决方案

1. 核心问题:误用ThreadPool而非进程池

你使用的ThreadPool是线程池,受Python全局解释器锁(GIL)限制,无法利用多核CPU执行CPU密集型任务。你的代码中fitCvines、compute_mispricing_index等操作均为CPU密集型,线程池本质还是单线程串行执行,因此耗时与原代码无差异,CPU利用率极低。

解决方法:改用multiprocessing.Pool(进程池),每个进程拥有独立的Python解释器,能真正实现多核并行:

# 替换导入语句
from multiprocessing import Pool

if __name__ == '__main__':
    # 指定进程数,建议设为CPU核心数(比如12)
    pool = Pool(processes=12)
    results = pool.map(process_period, range(numPeriods))
    pool.close()
    pool.join()  # 等待所有进程完成

    periodMisIdx_CTrad = [result[0] for result in results]
    periodADF_CTrad = [result[1] for result in results]

2. 额外优化建议

  • 数据传递优化:进程间数据拷贝存在开销,如果normRets_np、Q_trad_new等数据量极大,建议用multiprocessing.Manager实现共享内存,或提前拆分数据后传入进程,避免重复拷贝。
  • GPU利用:若fitCvines或compute_mispricing_index支持GPU加速(如基于PyTorch/TensorFlow实现),可将部分任务转移到GPU执行,进一步提升效率。
  • 分层次并行:除了对i(周期)并行,还可在每个周期内对j(资产)再做并行处理,比如在process_period内部启动子进程池(注意避免过度并行导致资源竞争)。
  • 函数内部串行优化:adfuller循环可改为并行执行,减少单进程内的串行瓶颈。

3. 效果验证

替换进程池后,运行2个周期的耗时应接近10分钟(而非20分钟),CPU利用率会显著提升至多核负载水平(如12核CPU利用率接近100%)。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.01 07:54:58