为带For循环的MACD计算函数实现并发优化(进程/线程)
为多品种MACD计算代码添加并发支持
需求:为交易品种MACD处理代码添加进程式/线程式并发,以提升大量品种或多时间框架信号处理时的效率。优先支持多进程+多线程组合,单独实现任一方式也可接受。
一、多进程实现版本
针对多品种独立计算的场景,使用multiprocessing.Pool实现进程级并行,充分利用多核CPU资源。
import numpy as np import pandas as pd import MetaTrader5 as mt5 from multiprocessing import Pool # 全局配置参数 marketFeedArrayLength = 364320 timeFrame = mt5.TIMEFRAME_M1 symbols_MultiTFMACD = ['AUDCAD', 'AUDJPY', 'AUDNZD', 'AUDUSD', 'CADJPY', 'EURAUD', 'EURCAD'] m1_macdLine_fastEMA_period_1 = 1 m1_macdLine_slowEMA_period_480 = 480 m1_macdLine_slowEMA_period_840 = 840 def process_symbol(symbol): """单品种MACD计算与文件更新逻辑,作为进程池任务单元""" # 读取收盘价数据 closeDataForEMA = np.load(f"{symbol}.npy")[:,4] m1_closeDataForEMA = np.flip(closeDataForEMA[-1::-timeFrame]) # 计算MACD权重与EMA weightsFast_1 = np.exp(np.linspace(-1., 0., m1_macdLine_fastEMA_period_1)) weightsSlow_480 = np.exp(np.linspace(-1., 0., m1_macdLine_slowEMA_period_480)) weightsSlow_840 = np.exp(np.linspace(-1., 0., m1_macdLine_slowEMA_period_840)) # 归一化权重 weightsFast_1 /= weightsFast_1.sum() weightsSlow_480 /= weightsSlow_480.sum() weightsSlow_840 /= weightsSlow_840.sum() # 计算EMA m1_fastEMA_1 = np.pad(np.convolve(weightsFast_1, m1_closeDataForEMA, mode='valid'), (m1_macdLine_fastEMA_period_1 - 1, 0)) m1_slowEMA_480 = np.pad(np.convolve(weightsSlow_480, m1_closeDataForEMA, mode='valid'), (m1_macdLine_slowEMA_period_480 - 1, 0)) m1_slowEMA_840 = np.pad(np.convolve(weightsSlow_840, m1_closeDataForEMA, mode='valid'), (m1_macdLine_slowEMA_period_840 - 1, 0)) # 计算MACD线 m1_macdLine_1_480 = m1_fastEMA_1 - m1_slowEMA_480 m1_macdLine_1_840 = m1_fastEMA_1 - m1_slowEMA_840 # 更新.npy文件 npyToPandas = pd.DataFrame(np.load(f"{symbol}.npy")) npyToPandas['8'] = m1_macdLine_1_480 npyToPandas['9'] = m1_macdLine_1_840 np.save(symbol, npyToPandas.to_numpy(), allow_pickle=True, fix_imports=False) print(f"完成 {symbol} 的MACD计算与文件更新") def multiTF_MACD_parallel(): # 创建进程池,默认使用CPU核心数 with Pool() as pool: pool.map(process_symbol, symbols_MultiTFMACD) if __name__ == "__main__": multiTF_MACD_parallel()
核心说明
- 将单品种处理逻辑抽离为独立函数,作为进程池的执行单元
- 利用
Pool.map自动分发任务到多个进程,并行处理所有货币对 - 必须通过
if __name__ == "__main__":包裹主调用,避免Windows系统下的进程创建异常 - 每个进程独立操作对应品种的文件,无资源冲突
二、多进程+多线程组合版本
针对单品种需计算大量指标的场景,在进程并行的基础上,为单品种内的多EMA计算添加线程并行,进一步提升资源利用率。
import numpy as np import pandas as pd import MetaTrader5 as mt5 from multiprocessing import Pool import threading # 全局配置参数 marketFeedArrayLength = 364320 timeFrame = mt5.TIMEFRAME_M1 symbols_MultiTFMACD = ['AUDCAD', 'AUDJPY', 'AUDNZD', 'AUDUSD', 'CADJPY', 'EURAUD', 'EURCAD'] m1_macdLine_fastEMA_period_1 = 1 m1_macdLine_slowEMA_period_480 = 480 m1_macdLine_slowEMA_period_840 = 840 def calculate_ema(close_data, period, result_dict, key): """单个周期EMA计算逻辑,供线程调用""" weights = np.exp(np.linspace(-1., 0., period)) weights /= weights.sum() ema = np.pad(np.convolve(weights, close_data, mode='valid'), (period - 1, 0)) result_dict[key] = ema def process_symbol(symbol): """单品种处理逻辑,内部用线程并行计算多周期EMA""" # 读取收盘价数据 closeDataForEMA = np.load(f"{symbol}.npy")[:,4] m1_closeDataForEMA = np.flip(closeDataForEMA[-1::-timeFrame]) # 用字典存储线程计算结果(线程无法直接返回值) ema_results = {} threads = [] # 创建线程并行计算不同周期的EMA threads.append(threading.Thread(target=calculate_ema, args=(m1_closeDataForEMA, m1_macdLine_fastEMA_period_1, ema_results, 'fast'))) threads.append(threading.Thread(target=calculate_ema, args=(m1_closeDataForEMA, m1_macdLine_slowEMA_period_480, ema_results, 'slow_480'))) threads.append(threading.Thread(target=calculate_ema, args=(m1_closeDataForEMA, m1_macdLine_slowEMA_period_840, ema_results, 'slow_840'))) # 启动并等待所有线程完成 for t in threads: t.start() for t in threads: t.join() # 计算MACD线 m1_macdLine_1_480 = ema_results['fast'] - ema_results['slow_480'] m1_macdLine_1_840 = ema_results['fast'] - ema_results['slow_840'] # 更新.npy文件 npyToPandas = pd.DataFrame(np.load(f"{symbol}.npy")) npyToPandas['8'] = m1_macdLine_1_480 npyToPandas['9'] = m1_macdLine_1_840 np.save(symbol, npyToPandas.to_numpy(), allow_pickle=True, fix_imports=False) print(f"完成 {symbol} 的MACD计算与文件更新") def multiTF_MACD_parallel(): with Pool() as pool: pool.map(process_symbol, symbols_MultiTFMACD) if __name__ == "__main__": multiTF_MACD_parallel()
核心说明
- 将单个EMA计算抽离为线程任务,并行处理同一品种的多周期EMA计算
- 用字典存储线程计算结果,解决线程无法直接返回值的问题
- 多进程处理不同品种,多线程处理单品种内的多指标计算,最大化CPU利用率
三、注意事项
- 进程数控制:若磁盘IO成为瓶颈,可通过
Pool(processes=4)手动限制进程数,避免磁盘读写过载 - 文件冲突:确保每个进程仅操作对应品种的文件,禁止多进程同时读写同一文件
- numpy线程优化:numpy默认启用多线程加速,若与自定义线程冲突,可通过
np.set_num_threads(1)禁用内部多线程 - MT5连接:若需调用MT5接口,需在每个进程内单独初始化MT5连接,禁止跨进程共享MT5实例
内容的提问来源于stack exchange,提问作者rn kim
相关产品推荐
相关产品推荐

