Python并行化卷积运算问题:多进程耗时异常及优化方案咨询
多进程优化fftconvolve并行效率问题
我需要实现一个数据拟合函数,形式为F = conv(R1,T1) + conv(R2,T2) + conv(R3,T3),这个函数会被多次调用。为了优化性能,我尝试用multiprocessing.Pool并行执行3次scipy.signal.fftconvolve卷积操作,但测试发现并行耗时(Time 2=0.028s)反而比串行循环(Time 1=0.016s)更长;但用time.sleep模拟任务时,并行耗时符合预期(约等于单次耗时)。我需要能让3次函数调用并行执行、总耗时接近单次调用时长的正确实现方案。
测试代码如下:
import time import numpy as np from scipy.signal import fftconvolve from multiprocessing import Pool q = np.arange(0,5,0.0001) theory = np.exp(-1*q) q2 = np.arange(-5,5,0.0001) resolution = np.exp(-1/2*(q2**2)) def f( i ): return fftconvolve( theory , resolution , mode='same' ) # time.sleep(1) start = time.time() f( 0 ) end = time.time() print(f'Time 0 = { end-start }') start = time.time() for i in range(0,3): f( i ) end = time.time() print(f'Time 1 = { end-start }') pool = Pool( processes= 3 ) start = time.time() pool.map( f , [0, 1, 2, ]) end = time.time() print(f'Time 2 = { end-start }')
测试结果:
Time 0 = 0.005838155746459961 Time 1 = 0.01575756072998047 Time 2 = 0.028168439865112305
问题原因分析
- 进程启动与数据序列化开销:
multiprocessing.Pool创建子进程需要额外启动成本,且每次调用函数时,theory和resolution这类大数组会被序列化(pickle)传递给子进程,这个序列化/反序列化的时间可能远超过卷积本身的耗时,导致并行总耗时更高。 - BLAS多线程冲突:
scipy.fftconvolve底层依赖的BLAS库(如OpenBLAS、MKL)本身默认是多线程并行的。当用multiprocessing开启多个进程时,每个进程内的BLAS又会启动多线程,导致CPU核心被过度占用,线程频繁切换反而降低效率。
解决方案
方案1:禁用BLAS多线程+优化进程初始化
先关闭BLAS的多线程,避免线程冲突;同时通过进程初始化函数提前传递大数组,避免重复序列化开销:
import os # 禁用BLAS多线程,根据你使用的库选择对应环境变量 os.environ["OMP_NUM_THREADS"] = "1" os.environ["OPENBLAS_NUM_THREADS"] = "1" os.environ["MKL_NUM_THREADS"] = "1" import time import numpy as np from scipy.signal import fftconvolve from multiprocessing import Pool q = np.arange(0,5,0.0001) theory = np.exp(-1*q) q2 = np.arange(-5,5,0.0001) resolution = np.exp(-1/2*(q2**2)) def f( i ): return fftconvolve( theory , resolution , mode='same' ) # 预热一次,消除首次调用的初始化开销对计时的影响 f(0) start = time.time() f( 0 ) end = time.time() print(f'Time 0 = { end-start }') start = time.time() for i in range(0,3): f( i ) end = time.time() print(f'Time 1 = { end-start }') # 初始化子进程时传递大数组,避免重复序列化 def init_worker(theory_arr, resolution_arr): global theory, resolution theory = theory_arr resolution = resolution_arr if __name__ == '__main__': pool = Pool(processes=3, initializer=init_worker, initargs=(theory, resolution)) start = time.time() pool.map(f, [0,1,2]) end = time.time() print(f'Time 2 = { end-start }')
方案2:使用线程池替代多进程
fftconvolve底层是C实现的BLAS操作,会释放GIL,因此用线程池可以避免进程启动和序列化的开销,效率更高:
import os os.environ["OMP_NUM_THREADS"] = "1" os.environ["OPENBLAS_NUM_THREADS"] = "1" import time import numpy as np from scipy.signal import fftconvolve from concurrent.futures import ThreadPoolExecutor q = np.arange(0,5,0.0001) theory = np.exp(-1*q) q2 = np.arange(-5,5,0.0001) resolution = np.exp(-1/2*(q2**2)) def f( i ): return fftconvolve( theory , resolution , mode='same' ) # 预热 f(0) start = time.time() f( 0 ) end = time.time() print(f'Time 0 = { end-start }') start = time.time() for i in range(0,3): f( i ) end = time.time() print(f'Time 1 = { end-start }') start = time.time() with ThreadPoolExecutor(max_workers=3) as executor: executor.map(f, [0,1,2]) end = time.time() print(f'Time 2 = { end-start }')
方案3:批量卷积计算
如果多个卷积的输入结构一致,可以用向量化方式一次性计算所有卷积,彻底避免并行开销:
import time import numpy as np from scipy.signal import fftconvolve q = np.arange(0,5,0.0001) theory = np.exp(-1*q) # 模拟3个输入R1/R2/R3,堆叠成三维数组 theories = np.stack([theory]*3, axis=0) q2 = np.arange(-5,5,0.0001) resolution = np.exp(-1/2*(q2**2)) start = time.time() # 利用广播一次性计算3个卷积 results = fftconvolve(theories, resolution[np.newaxis, :], mode='same') end = time.time() print(f'Time batch = { end-start }') # results的每个维度对应一个卷积结果
内容的提问来源于stack exchange,提问作者Richard Williams
相关产品推荐
相关产品推荐

