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

Python中如何实现真正的并行编程并充分利用CPU资源?

如何让Python多进程充分利用CPU算力?

你的问题主要出在任务粒度太细和进程间数据传递/共享的方式错误这两个方面,导致CPU大部分时间消耗在进程调度和通信上,而不是实际计算,所以使用率上不去。另外你的代码还有一个隐藏的逻辑错误:子进程修改的score数组是主进程数组的副本,最后主进程的score根本不会被更新。

核心问题分析

  1. 任务粒度极小:你给每个i,j的组合都单独提交一个任务,总共10000个任务。每个任务只做一次加法运算,进程启动、参数传递的开销远大于计算本身,CPU大部分时间在等待任务调度,而非执行计算。
  2. 大对象重复传递:每次调用apply_async都把整个df传给子进程,df是一个大的DataFrame,重复传递会产生大量内存复制和通信开销,拖慢整体速度。
  3. 共享数据错误:score数组在子进程中是独立副本,子进程对它的修改不会同步回主进程,最后你得到的score还是初始的全0数组。

改进方案

下面给出两种可行的优化方案,从调整任务粒度和优化数据传递两个方向解决问题:

方案一:调整任务粒度,批量处理任务

把原来的单i,j任务改成按行批量处理,每个任务负责计算一整行的所有j值,这样任务数量从10000降到100,极大减少进程调度和通信开销。同时通过返回结果的方式收集计算值,避免共享内存的复杂操作。

import numpy as np
import pandas as pd
import multiprocessing

def calc_score_row(df, i):
    # 一次性计算第i行所有j的score值
    row_data = df.loc[i, 'data']
    # 利用numpy向量运算加速,比循环快很多
    row_score = row_data + df['data'].values
    return i, row_score

if __name__ == '__main__':
    df = pd.read_csv('data.csv')
    n = 100
    score = np.zeros([n, n])
    
    # 初始化进程池,使用全部CPU核心
    pool = multiprocessing.Pool(multiprocessing.cpu_count())
    
    # 提交所有行的计算任务
    results = []
    for i in range(n):
        res = pool.apply_async(calc_score_row, args=(df, i))
        results.append(res)
    
    # 逐个获取结果并填充到score数组
    for res in results:
        i, row_score = res.get()
        score[i, :] = row_score
    
    pool.close()
    pool.join()

方案二:使用共享内存传递大数组(适合超大数据集)

如果你的df非常大,重复传递df会消耗大量内存,这时可以用共享内存让所有子进程直接访问同一份数据,避免内存复制。

import numpy as np
import pandas as pd
import multiprocessing
from multiprocessing import shared_memory

def calc_score_row_shared(shm_name, array_shape, array_dtype, i):
    # 连接到主进程创建的共享内存
    existing_shm = shared_memory.SharedMemory(name=shm_name)
    # 从共享内存中加载数据数组
    data_array = np.ndarray(array_shape, dtype=array_dtype, buffer=existing_shm.buf)
    # 计算当前行的所有score值
    row_score = data_array[i] + data_array
    existing_shm.close()
    return i, row_score

if __name__ == '__main__':
    df = pd.read_csv('data.csv')
    data_array = df['data'].values
    n = len(data_array)
    score = np.zeros([n, n])
    
    # 创建共享内存,存储data列的数据
    shm = shared_memory.SharedMemory(create=True, size=data_array.nbytes)
    # 将数据写入共享内存
    shared_data = np.ndarray(data_array.shape, dtype=data_array.dtype, buffer=shm.buf)
    shared_data[:] = data_array[:]
    
    pool = multiprocessing.Pool(multiprocessing.cpu_count())
    results = []
    for i in range(n):
        res = pool.apply_async(
            calc_score_row_shared,
            args=(shm.name, data_array.shape, data_array.dtype, i)
        )
        results.append(res)
    
    # 收集结果填充score数组
    for res in results:
        i, row_score = res.get()
        score[i, :] = row_score
    
    pool.close()
    pool.join()
    # 释放共享内存资源
    shm.close()
    shm.unlink()

额外优化建议

  • 优先用numpy向量运算:如果你的计算可以用numpy的向量/矩阵运算实现(比如score = df['data'].values[:, np.newaxis] + df['data'].values[np.newaxis, :]),那完全不需要多进程——numpy的底层是C实现的,速度快且能充分利用CPU,开销比多进程小很多。
  • 任务粒度适中:不要把任务拆得太细,也不要让单个任务过大导致负载不均,尽量让每个进程的任务量相近。
  • 用imap_unordered收集结果:如果结果的顺序不重要,可以用pool.imap_unordered代替apply_async,能更快地获取计算完成的结果,减少等待时间。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.28 06:27:50