You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

求助: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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.08.07 21:40:21