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

Python循环创建子进程失效,并行化任务无法执行求助

并行化余弦距离计算任务失败排查

我尝试对一项耗时的余弦距离计算任务做并行化处理,但始终无法正常运行,没法定位问题原因。

初始代码(Process+Queue实现)

from scipy.spatial import distance
import numpy as np
from scipy import linalg
import os
from multiprocessing import Queue
from multiprocessing import Process
from typing import Tuple
    
num_cpu = os.cpu_count()
arr_1 = np.random.random((20000, 1000))
arr_2 = np.random.random((20000, 1000))
   
def cosine_similarity( a: np.ndarray, b: np.ndarray):
    assert len(a) == len(b)
    len_a = linalg.norm(a)
    len_b = linalg.norm(b)
    # Check if one vector is all zeros. Possible, not probable
    if len_a == 0 or len_b == 0:
        return 0
    return a.dot(b) / (len_a * len_b)

def cosine_distance(a: np.ndarray, b: np.ndarray):
    assert len(a) == len(b)
    return 1 - cosine_similarity(a, b)

def run_task(start: int, end: int, queue_out: Queue):
    print("Called with " + str(start) + " to " + str(end))
    if end > len(arr_1):
        end = len(arr_1)
    for i in range(start, end):
        min_dist = 2.0
        labeled_data = arr_1[i]
        for unlabeled_data in arr_2:
            min_dist = min(min_dist, cosine_distance(labeled_data, unlabeled_data))

        queue_out.put((i, min_dist))

step = len(arr_1) // num_cpu

t_s = []
queue = Queue()
for i in range(0, len(arr_1), step):
    print("Calling with " + str(i) + " to " + str(i + step))
    p = Process(target=run_task, args=(i, i + step, queue))
    p.start()
    t_s.append(p)

for p in t_s:
    p.join()

问题现象

  • 调用queue.qsize()结果为0,代码瞬间执行完毕,仅输出主进程的“Calling with...”,无“Called with”输出;
  • 手动调用run_task(0, 1000, queue)时,任务正常运行约3分钟,queue.qsize()为1000;
  • 查看t_s中的进程,均显示已停止且退出码为1。

改写后的代码(Pool.map实现)

def run_task(index):
    print("Called index " + str(index))
    min_dist = 2.0
    labeled_data = arr_1[index]
    for unlabeled_data in arr_2:
        min_dist = min(min_dist, cosine_distance(labeled_data, unlabeled_data))

    (index, min_dist)

a = 0
with Pool(processes=num_cpu) as pool:
    a = pool.imap(run_task, range(len(arr_1)))

for i in a:
    print(f"showing the result as it is ready {i}")

问题现象

出现同样问题,子进程任务未被调用。

请问我哪里操作错误?


内容的提问来源于stack exchange,提问作者SinisterMJ

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.20 09:24:27