如何利用Numpy结合多核计算实现向量运算并行化?
嘿,我来帮你搞定这个多核加速的问题!你的代码里那个逐行循环在处理这么大的矩阵时,肯定会有点慢对吧?毕竟单线程跑这么多计算太浪费多核CPU的性能了。下面给你几个实用的方案,既能结合多核算力,又能和Numpy无缝配合:
其实你的运算完全可以不用写循环——Numpy的广播机制能帮你把逐行的操作转换成矩阵级的矢量化运算,而且如果你的Numpy是链接了支持多核的BLAS库(比如Anaconda默认的MKL,或者pip安装的OpenBLAS),底层会自动用多个CPU核心来并行计算,这是最省心的优化方式。
修改后的代码:
import numpy as np num_row, num_col = 6000, 13572 ss = np.ones((num_row, num_col), dtype=np.complex128) # 把ph和fre转成列向量,方便和tau(行向量)做广播运算 ph = np.random.standard_normal(num_row)[:, np.newaxis] fre = np.random.standard_normal(num_row)[:, np.newaxis] tau = np.random.standard_normal(num_col) # 一步完成所有行的运算,无循环 exponents = 1j * (ph + fre * tau) ss *= np.exp(exponents)
原理很简单:列向量ph、fre会和行向量tau自动广播成和ss一样形状的矩阵,整个指数运算和乘法都是批量完成的,Numpy底层的优化会帮你搞定多核利用。
如果某些场景下没法完全矢量化,或者你想保留循环的逻辑,那Numba绝对是你的好帮手——它能把Python循环编译成高效的机器码,还支持一键开启多核并行。
修改后的代码:
import numpy as np from numba import jit, prange num_row, num_col = 6000, 13572 ss = np.ones((num_row, num_col), dtype=np.complex128) ph = np.random.standard_normal(num_row) fre = np.random.standard_normal(num_row) tau = np.random.standard_normal(num_col) # 用numba的JIT编译,开启并行模式 @jit(nopython=True, parallel=True) def apply_exponents(ss, ph, fre, tau): # 用prange代替range,告诉numba要并行这个循环 for idx in prange(num_row): ss[idx, :] *= np.exp(1j*(ph[idx] + fre[idx]*tau)) apply_exponents(ss, ph, fre, tau)
这里的prange是关键,它会让Numba把循环的任务分配到多个CPU核心上,nopython=True则是让代码完全脱离Python解释器,速度提升非常明显。
如果上面两种方法都不适用,比如你需要更细粒度的任务控制,那可以用Python的multiprocessing库手动把矩阵拆分成多个块,让每个进程处理一块,最后再把结果拼接回去。
修改后的代码:
import numpy as np from multiprocessing import Pool, cpu_count # 定义每个进程要执行的任务函数 def process_chunk(args): chunk_ss, chunk_ph, chunk_fre, tau = args for idx in range(chunk_ss.shape[0]): chunk_ss[idx, :] *= np.exp(1j*(chunk_ph[idx] + chunk_fre[idx]*tau)) return chunk_ss num_row, num_col = 6000, 13572 ss = np.ones((num_row, num_col), dtype=np.complex128) ph = np.random.standard_normal(num_row) fre = np.random.standard_normal(num_row) tau = np.random.standard_normal(num_col) # 根据CPU核心数拆分数据块 num_chunks = cpu_count() chunk_size = num_row // num_chunks chunks = [] for i in range(num_chunks): start = i * chunk_size # 最后一块要处理剩余的行 end = start + chunk_size if i != num_chunks - 1 else num_row chunks.append((ss[start:end], ph[start:end], fre[start:end], tau)) # 启动多进程池处理所有块 with Pool(num_chunks) as pool: results = pool.map(process_chunk, chunks) # 把处理后的块拼接回原矩阵 ss = np.vstack(results)
不过要注意:多进程会有数据拷贝的开销,如果你的矩阵特别大,这个方法的效率可能不如前两种,所以优先考虑前两个方案哦!
内容的提问来源于stack exchange,提问作者LowQualityDelivery

