求助:multiprocessing.Process.start调用后主进程无并行问题排查
问题分析
你看到的“主进程同步等待子进程完成”是假象,核心原因有两个:
- 大进程间数据传递的阻塞:你的每个chunk包含约166万元素,在Windows的
spawn启动模式下,proc.start()会同步完成chunk的序列化(pickle)与跨进程传递,这个拷贝过程会阻塞主进程,导致子进程看起来是逐个启动执行。 - 输出缓冲导致的日志误导:子进程的
print输出存在缓冲,主进程与子进程的日志输出没有同步,进一步让你误以为进程是串行执行的。
另外,你的代码存在一个功能缺陷:multiprocessing.Process不会返回函数执行结果,当前代码完全没有收集乘法后的计算结果。
解决方案
1. 改用multiprocessing.Pool(推荐)
Pool会自动管理进程池和数据分发,比手动创建Process更高效,还能直接收集结果:
import math import time import numpy as np from multiprocessing import Pool N = 10000000 THREADS_AMOUNT = 6 THREAD_LIST_LEN = math.ceil(N / THREADS_AMOUNT) def multiply_vector(vector, num_multiply=2): print('start multiplication') result = [num * num_multiply for num in vector] print('end multiplication') return result def main(): random_list = np.random.rand(N).tolist() chunks = [random_list[i * THREAD_LIST_LEN: (i + 1) * THREAD_LIST_LEN] for i in range(THREADS_AMOUNT)] start = time.perf_counter() with Pool(THREADS_AMOUNT) as pool: results = pool.map(multiply_vector, chunks) # 合并所有子进程的结果 final_result = [] for res in results: final_result.extend(res) end = time.perf_counter() print(f"{end - start:.5f}") if __name__ == '__main__': main()
2. 用共享内存减少数据拷贝(超大数据场景)
如果数据量极大,避免跨进程拷贝数据,改用共享内存存储原数据:
import math import time import numpy as np from multiprocessing import Process, Array from ctypes import c_double N = 10000000 THREADS_AMOUNT = 6 THREAD_LIST_LEN = math.ceil(N / THREADS_AMOUNT) def multiply_shared(start_idx, end_idx, shared_array, num_multiply=2): print('start multiplication') for i in range(start_idx, end_idx): shared_array[i] *= num_multiply print('end multiplication') def main(): random_array = np.random.rand(N) # 创建无锁的共享内存数组 shared_array = Array(c_double, random_array, lock=False) start = time.perf_counter() procs = [] for i in range(THREADS_AMOUNT): start_idx = i * THREAD_LIST_LEN end_idx = min((i + 1) * THREAD_LIST_LEN, N) proc = Process(target=multiply_shared, args=(start_idx, end_idx, shared_array)) procs.append(proc) proc.start() for proc in procs: proc.join() # 从共享内存中读取最终结果 final_result = np.frombuffer(shared_array.get_obj(), dtype=np.float64) end = time.perf_counter() print(f"{end - start:.5f}") if __name__ == '__main__': main()
3. 快速验证并行效果
把N改成1000后运行原代码,你会看到多个子进程的start multiplication日志同时出现,证明原代码本身是支持并行的,只是大数据传递的开销造成了串行的假象。
内容的提问来源于stack exchange,提问作者Questioner
相关产品推荐
相关产品推荐

