Python中如何实现真正的并行编程并充分利用CPU资源?
如何让Python多进程充分利用CPU算力?
你的问题主要出在任务粒度太细和进程间数据传递/共享的方式错误这两个方面,导致CPU大部分时间消耗在进程调度和通信上,而不是实际计算,所以使用率上不去。另外你的代码还有一个隐藏的逻辑错误:子进程修改的score数组是主进程数组的副本,最后主进程的score根本不会被更新。
核心问题分析
- 任务粒度极小:你给每个
i,j的组合都单独提交一个任务,总共10000个任务。每个任务只做一次加法运算,进程启动、参数传递的开销远大于计算本身,CPU大部分时间在等待任务调度,而非执行计算。 - 大对象重复传递:每次调用
apply_async都把整个df传给子进程,df是一个大的DataFrame,重复传递会产生大量内存复制和通信开销,拖慢整体速度。 - 共享数据错误:
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
相关产品推荐
相关产品推荐

