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

Python多进程重复执行任务且耗时更长的问题排查

多进程并行计算的异常问题排查

问题描述

我有一个元组列表,需要在循环中执行计算任务。因计算耗时较长,尝试用多进程分配到多个CPU核心提速,但遇到两个异常:

  • 所有核心均被占用,但执行相同的计算任务——比如期望核心分别处理1+1、2+2、3+3,实际所有核心都跑这三项计算,结果重复三次。
  • 多进程模式的耗时反而比单进程更长,完全不符合预期。

测试代码

原始代码

from multiprocessing import Process
import numpy as np
import itertools
import time
import math

def func(list_products):
    for i in range(len(list_products)):
        # 示例仅保留循环结构,实际为耗时计算
        i

""" 创建元组列表 """
# 输入参数
item1_deviation = 1  # %
item1_steps = 10

item2_deviation = 10  # %
item2_N = 3  # 最小为2

# 生成第一个列表
item1_base = 360
item1_deviation_rounded = item1_base * (item1_deviation / 100)
list_item1s = np.arange(item1_base - item1_deviation_rounded,
                        item1_base + item1_deviation_rounded + item1_deviation,
                        item1_steps).tolist()
list_item1s = [int(math.ceil(item / item1_steps)) * item1_steps for item in list_item1s]

# 生成第二个列表
item2_base = 0.77
list_item2s = np.linspace(item2_base - (item2_base * (item2_deviation / 100)),
                         item2_base + (item2_base * (item2_deviation / 100)),
                         item2_N).tolist()
list_item2s = [round(item, 3) for item in list_item2s]

# 生成所有组合的元组列表
list_total_sequence = []
for i in range(6):
    list_sequence = [list_item1s, list_item2s]
    list_product = list(itertools.product(*list_sequence))
    list_total_sequence.append(list_product)
list_total_product = list(itertools.product(*list_total_sequence))

""" 单进程测试 """
t1 = time.perf_counter()
for i in range(len(list_total_product[0:20])):
    # 示例仅保留循环结构
    i
t2 = time.perf_counter()
print(f"单进程耗时: {t2-t1:.6f}秒")

""" 多进程测试 """
if __name__ == '__main__':
    t1 = time.perf_counter()
    p = Process(target=func, args=(list_total_product[0:20],))
    p.start()
    p.join()
    t2 = time.perf_counter()
    print(f"多进程耗时: {t2-t1:.6f}秒")

更新后的代码(仍有问题)

根据建议修改后,当num_items=700且仅用1个核心时,多进程耗时约0.183秒,单进程仅耗时3.68e-05秒:

if __name__ == '__main__':
    processes = []
    num_cores = 1
    num_items = 700
    num_items_per_core = int(num_items / num_cores)
    ipc = num_items_per_core
    
    t1 = time.perf_counter()
    
    for i in range(0, num_cores):
        p = Process(target=func, args=(list_total_product[ipc*i:ipc*i+ipc],))
        processes.append(p)
        p.start()
        
    for p in processes:
        p.join()
    
    t2 = time.perf_counter()
    print(f"多进程耗时: {t2-t1:.6f}秒")

问题根源与解决方法

1. 任务分配错误(重复执行相同任务)

原始代码中,你给每个进程传递的是完整的任务列表(比如list_total_product[0:20]),导致所有进程都处理全部任务。即使更新后的代码做了切片,若num_cores设置不合理或切片逻辑有问题,也可能出现任务重复。

修复方案:正确拆分任务切片
确保每个进程只处理任务列表的一部分,同时处理整除剩余的任务:

if __name__ == '__main__':
    import multiprocessing

    t1 = time.perf_counter()
    num_cores = multiprocessing.cpu_count()  # 使用实际CPU核心数
    target_tasks = list_total_product[:700]  # 取前700个任务
    total_tasks = len(target_tasks)
    tasks_per_core = total_tasks // num_cores
    remaining_tasks = total_tasks % num_cores

    processes = []
    start_idx = 0
    for i in range(num_cores):
        # 给前N个核心多分配1个任务,处理剩余量
        end_idx = start_idx + tasks_per_core + (1 if i < remaining_tasks else 0)
        sub_tasks = target_tasks[start_idx:end_idx]
        p = multiprocessing.Process(target=func, args=(sub_tasks,))
        processes.append(p)
        p.start()
        start_idx = end_idx

    for p in processes:
        p.join()

    t2 = time.perf_counter()
    print(f"多进程耗时: {t2-t1:.6f}秒")

2. 多进程开销超过并行收益

多进程存在固定开销:进程创建、内存复制、进程间通信等。如果func中的计算任务过于轻量(比如示例中仅循环变量i),这些开销会完全盖过并行带来的效率提升,导致总耗时更长。

修复方案:确保计算任务足够“重”
在实际业务代码中,若计算确实耗时,并行会有明显收益;若测试代码需要模拟,可添加耗时计算逻辑:

def func(list_products):
    result = []
    for item in list_products:
        # 模拟耗时计算:复杂数学运算
        temp = 0
        for _ in range(10000):
            temp += math.sqrt(item[0][0] ** 2 + item[0][1] ** 2)
        result.append(temp)
    return result

3. 额外优化建议

  • 使用进程池(multiprocessing.Pool):复用进程,减少进程创建的重复开销,更适合批量任务:
if __name__ == '__main__':
    import multiprocessing

    t1 = time.perf_counter()
    pool = multiprocessing.Pool(processes=multiprocessing.cpu_count())
    target_tasks = list_total_product[:700]
    # 自动分配任务到进程池
    pool.map(func, [ [task] for task in target_tasks ])
    pool.close()
    pool.join()
    t2 = time.perf_counter()
    print(f"进程池耗时: {t2-t1:.6f}秒")
  • 避免大数据复制:若list_total_product数据量极大,传递给进程时会复制内存,可改用共享内存或队列传递数据,降低开销。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.03 06:41:01