Python多进程任务:需超10个进程的矩阵运算对称性验证
解决方案:多进程任务实现与资源监控
一、实现超过10个工作进程的多进程方案
问题根源
如果仅直接分配10个矩阵对验证任务,即便进程池设置了超过10的max_workers,实际同时运行的进程最多只有10个(受任务数量限制)。要让超过10个进程全部启用并工作,需将每个矩阵对的验证任务拆分为多个子任务,确保进程池有足够任务分配给所有worker进程。
1. 任务拆分思路
将单个矩阵对的验证拆分为分片运算(比如按矩阵的行/列分片),把原本10个主任务扩展为N×10个子任务(N>1),让所有worker进程都能被调度执行。
2. 代码实现示例
方法1:使用multiprocessing.Pool
import multiprocessing import numpy as np # 替换为你的实际运算逻辑 def fn(mat1, mat2): return np.dot(mat1, mat2) # 子任务:验证单个分片的运算结果是否相等 def verify_subtask(mat_a, mat_b, slice_idx, slice_size): a_slice = mat_a[slice_idx*slice_size : (slice_idx+1)*slice_size] b_slice = mat_b[slice_idx*slice_size : (slice_idx+1)*slice_size] res1 = fn(a_slice, b_slice) res2 = fn(b_slice, a_slice) return np.array_equal(res1, res2) # 拆分单个矩阵对为多个子任务 def split_matrix_task(mat_a, mat_b, num_subtasks=5): slice_size = mat_a.shape[0] // num_subtasks return [(mat_a, mat_b, i, slice_size) for i in range(num_subtasks)] if __name__ == "__main__": # 生成10对测试矩阵(100x100方阵) A = [np.random.rand(100, 100) for _ in range(10)] B = [np.random.rand(100, 100) for _ in range(10)] # 设置进程池大小为15(严格超过10) pool_size = 15 pool = multiprocessing.Pool(processes=pool_size) # 生成所有子任务 all_subtasks = [] for mat_a, mat_b in zip(A, B): all_subtasks.extend(split_matrix_task(mat_a, mat_b)) # 执行子任务并汇总结果 results = pool.starmap(verify_subtask, all_subtasks) pair_results = [] idx = 0 for _ in range(10): pair_res = all(results[idx:idx+5]) pair_results.append(pair_res) idx += 5 print("各矩阵对验证结果:", pair_results) pool.close() pool.join()
方法2:使用concurrent.futures.ProcessPoolExecutor
import concurrent.futures import numpy as np def fn(mat1, mat2): return np.dot(mat1, mat2) def verify_subtask(mat_a, mat_b, slice_idx, slice_size): a_slice = mat_a[slice_idx*slice_size : (slice_idx+1)*slice_size] b_slice = mat_b[slice_idx*slice_size : (slice_idx+1)*slice_size] res1 = fn(a_slice, b_slice) res2 = fn(b_slice, a_slice) return np.array_equal(res1, res2) def split_matrix_task(mat_a, mat_b, num_subtasks=5): slice_size = mat_a.shape[0] // num_subtasks return [(mat_a, mat_b, i, slice_size) for i in range(num_subtasks)] if __name__ == "__main__": A = [np.random.rand(100, 100) for _ in range(10)] B = [np.random.rand(100, 100) for _ in range(10)] pool_size = 15 with concurrent.futures.ProcessPoolExecutor(max_workers=pool_size) as executor: all_subtasks = [] for mat_a, mat_b in zip(A, B): all_subtasks.extend(split_matrix_task(mat_a, mat_b)) # 提交并执行子任务 futures = [executor.submit(verify_subtask, *task) for task in all_subtasks] results = [future.result() for future in concurrent.futures.as_completed(futures)] # 汇总矩阵对验证结果 pair_results = [] idx = 0 for _ in range(10): pair_res = all(results[idx:idx+5]) pair_results.append(pair_res) idx += 5 print("各矩阵对验证结果:", pair_results)
二、进程资源监控方法
使用psutil库可以监控进程池内每个worker进程的CPU、内存占用情况,步骤如下:
1. 监控实现示例
import multiprocessing import numpy as np import psutil import time def fn(mat1, mat2): time.sleep(0.5) # 模拟运算耗时 return np.dot(mat1, mat2) def verify_subtask(mat_a, mat_b, slice_idx, slice_size): a_slice = mat_a[slice_idx*slice_size : (slice_idx+1)*slice_size] b_slice = mat_b[slice_idx*slice_size : (slice_idx+1)*slice_size] res1 = fn(a_slice, b_slice) res2 = fn(b_slice, a_slice) return np.array_equal(res1, res2) # 监控进程资源的函数 def monitor_workers(pool, interval=0.2): worker_pids = [proc.pid for proc in pool._pool] print("进程资源监控(PID, CPU%, 内存MB):") while any(psutil.pid_exists(pid) for pid in worker_pids): for pid in worker_pids: if psutil.pid_exists(pid): proc = psutil.Process(pid) cpu = proc.cpu_percent() mem = proc.memory_info().rss / (1024 * 1024) print(f"PID {pid}: CPU {cpu:.1f}%, 内存 {mem:.2f} MB") time.sleep(interval) print("---") if __name__ == "__main__": A = [np.random.rand(100, 100) for _ in range(10)] B = [np.random.rand(100, 100) for _ in range(10)] pool_size = 15 pool = multiprocessing.Pool(processes=pool_size) # 启动监控线程(用线程避免额外进程开销) monitor_thread = multiprocessing.Process(target=monitor_workers, args=(pool,)) monitor_thread.start() # 生成并执行子任务 all_subtasks = [] for mat_a, mat_b in zip(A, B): slice_size = mat_a.shape[0] //5 all_subtasks.extend([(mat_a, mat_b, i, slice_size) for i in range(5)]) results = pool.starmap(verify_subtask, all_subtasks) # 汇总结果 pair_results = [] idx =0 for _ in range(10): pair_res = all(results[idx:idx+5]) pair_results.append(pair_res) idx +=5 print("各矩阵对验证结果:", pair_results) pool.close() pool.join() monitor_thread.join()
2. 说明
- 先安装
psutil:执行pip install psutil - 监控线程会定期输出每个worker进程的资源使用情况,直到所有worker进程结束
- Windows系统下必须将多进程代码放在
if __name__ == "__main__":块内,避免子进程重复导入模块的问题
关键注意事项
- 任务拆分粒度要平衡:过细会增加进程间通信开销,过粗可能无法填满所有worker进程
- 进程池大小建议根据CPU核心数调整(比如核心数的1.5-2倍),同时确保严格超过10
内容的提问来源于stack exchange,提问作者Daniele Pittari
相关产品推荐
相关产品推荐

