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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.09 14:01:00