使用multiprocessing.Pool多进程时无法占满可用CPU资源问题咨询
多进程处理RREF矩阵无提速问题解决方案
核心问题原因
- 进程间通信开销占比过高:
get_rref_matrices在主进程生成矩阵后,需要通过pickle序列化传递给子进程,处理完的结果还要再序列化传回主进程,如果check_valid本身计算量很小,通信开销会完全抵消多进程的收益,子进程大部分时间处于等待任务状态,CPU利用率自然上不去。 - 子进程资源争抢:如果numpy使用了MKL、OpenBLAS等支持多线程的后端,每个子进程调用numpy时会自动启动多个线程,多进程叠加多线程会导致CPU核心争抢,整体效率下降。
- 任务调度开销过大:
imap本身需要维持输入输出的顺序对应,同时不合适的chunksize会导致频繁的任务调度,进一步消耗主进程资源。
优化方案
1. 关闭numpy内置多线程
你已经使用多进程做并行计算,numpy自带的多线程会造成资源竞争,在代码最开头加入环境变量配置,强制每个进程仅使用单线程:
import os os.environ['MKL_NUM_THREADS'] = '1' os.environ['OPENBLAS_NUM_THREADS'] = '1' os.environ['NUMEXPR_NUM_THREADS'] = '1' os.environ['OMP_NUM_THREADS'] = '1'
2. 增大任务粒度,减少通信次数
不要单条传递矩阵,改为批量生成、批量处理矩阵,大幅降低进程间通信的频次:
import multiprocessing as mp import numpy as np def check_valid(matrix): # 原有校验逻辑 if all_checks_passed: return matrix.copy() return None # 批量处理函数,一次性处理一批矩阵 def process_batch(matrix_batch): valid_list = [] for mat in matrix_batch: res = check_valid(mat) if res is not None: valid_list.append(res) return valid_list # 批量生成矩阵的生成器 def batch_generator(batch_size=10000): batch = [] for mat in get_rref_matrices(5): batch.append(mat) if len(batch) == batch_size: yield batch batch = [] if batch: yield batch if __name__ == '__main__': subgroups = [] with mp.Pool() as pool: # 用imap_unordered替代imap,不需要维持输入顺序的场景下可以大幅减少调度开销 for batch_res in pool.imap_unordered(process_batch, batch_generator(), chunksize=5): subgroups.extend(batch_res)
3. 进一步优化(可选)
如果get_rref_matrices本身生成速度较慢,成为瓶颈,可以把RREF生成逻辑也放到子进程中实现,每个子进程负责生成某一分类的RREF矩阵同时完成校验,完全省略主进程分发矩阵的开销。
内容的提问来源于stack exchange,提问作者MrLatinNerd
相关产品推荐
相关产品推荐

