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

为带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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.25 08:15:02