为何该Python代码并行化后未提升运行速度?
我在做蒙特卡洛练习时遇到了并行化性能瓶颈:针对模型(由抽样次数draws和随机种子seed参数化)做并行处理时,串行与并行运行耗时相近;但加入强制sleep后并行化有明显性能提升,说明并行逻辑本身是有效的。核心问题出在列表推导式[linreg(shuffle_treat(df)) for k in range(draws)]的并行化上,甚至针对单抽样并行时效果更差。
核心代码实现
import time import pandas as pd import numpy as np import multiprocessing as mp def linreg(df): y = df[['y']].values x = np.hstack([np.ones((df.shape[0], 1)), df[['treat']].values]) xx_inv = np.linalg.inv(x.T @ x) beta_hat = xx_inv @ (x.T @ y) return pd.Series(beta_hat.flat, index=['intercept', 'coef']) def shuffle_treat(df): df['treat'] = df['treat'].sample(frac=1, replace=False).values return df def run_analysis(draws, seed, sleep=0): N = 5000 df = pd.DataFrame({'treat':np.random.choice([0,1], size=N, replace=True)}) df['u'] = np.random.normal(size=N) df['y'] = df.eval('10 + 5*treat + u') np.random.seed(seed) time.sleep(sleep) est = [linreg(shuffle_treat(df)) for k in range(draws)] est = pd.concat(est, axis=0, sort=False, keys=range(draws), names=['k', 'param']) return est
测试代码与结果
param_list = [dict(draws=500, seed=1029), dict(draws=500, seed=1029)] param_list_sleep = [dict(draws=500, seed=1029, sleep=5), dict(draws=500, seed=1029, sleep=5)] def run_analysis_wrapper(params): run_analysis(**params) # 串行测试 start = time.time() for params in param_list: run_analysis_wrapper(params) end = time.time() print(f'double run 1 process: {(end - start):.2f} sec') # 并行测试 start = time.time() with mp.Pool(processes=2) as pool: pool.map(run_analysis_wrapper, param_list) end = time.time() print(f'double run 2 processes: {(end - start):.2f} sec') # 串行带sleep测试 start = time.time() for params in param_list_sleep: run_analysis_wrapper(params) end = time.time() print(f'double run 1 process w/ sleep: {(end - start):.2f} sec') # 并行带sleep测试 start = time.time() with mp.Pool(processes=2) as pool: pool.map(run_analysis_wrapper, param_list_sleep) end = time.time() print(f'double run 2 processes w/ sleep: {(end - start):.2f} sec')
输出结果:
double run 1 process: 2.52 sec double run 2 processes: 2.94 sec double run 1 process w/ sleep: 12.30 sec double run 2 processes w/ sleep: 7.71 sec
环境:Linux EC2实例(48核CPU),Python 3.9.16 conda环境
问题原因分析
进程通信开销抵消并行收益
计算任务单次耗时极短(线性回归+数据洗牌),但multiprocessing创建进程、序列化/反序列化数据(比如传递大DataFrame)的固定开销远超过并行计算节省的时间,导致整体并行效率反而低于串行。而sleep属于IO等待任务,并行时进程可以利用等待时间处理其他任务,因此能体现出性能优势。任务粒度太小
如果尝试针对单个抽样(循环内的单次linreg调用)做并行,任务粒度会进一步缩小,进程调度和通信的开销占比会更高,导致性能更差。计算逻辑可优化空间大
手动实现的线性回归(求逆矩阵)效率较低,且pandas的循环操作本身也存在一定的性能损耗,进一步压缩了并行收益的空间。
解决方案
1. 增大任务粒度
减少进程调度次数,将多个抽样打包为一个并行任务。例如,把原本的2个500抽样任务合并为几个更大的任务,或者在子进程内批量处理更多抽样:
# 修改run_analysis,让每个进程处理多组draws def run_analysis_batch(draws_per_batch, seed, batch_idx): # 每个batch使用独立种子保证随机性 np.random.seed(seed + batch_idx) N = 5000 df = pd.DataFrame({'treat':np.random.choice([0,1], size=N, replace=True)}) df['u'] = np.random.normal(size=N) df['y'] = df.eval('10 + 5*treat + u') est = [linreg(shuffle_treat(df)) for k in range(draws_per_batch)] return pd.concat(est, axis=0, sort=False, keys=range(batch_idx*draws_per_batch, (batch_idx+1)*draws_per_batch), names=['k', 'param']) # 并行处理4个batch,每个batch处理250次抽样 param_list_batch = [dict(draws_per_batch=250, seed=1029, batch_idx=i) for i in range(4)] with mp.Pool(processes=4) as pool: results = pool.map(run_analysis_batch, param_list_batch) final_est = pd.concat(results)
2. 优化计算逻辑
用更高效的实现替代手动线性回归,减少单次计算耗时:
def linreg(df): y = df['y'].values x = np.column_stack([np.ones(df.shape[0]), df['treat'].values]) # 使用np.linalg.lstsq替代手动求逆,更快且数值稳定性更好 beta_hat, *_ = np.linalg.lstsq(x, y, rcond=None) return pd.Series(beta_hat, index=['intercept', 'coef'])
3. 减少进程间数据传递
避免在进程间传递大DataFrame,改为在子进程内生成数据:
def run_single_draw(seed, N=5000): np.random.seed(seed) df = pd.DataFrame({'treat':np.random.choice([0,1], size=N, replace=True)}) df['u'] = np.random.normal(size=N) df['y'] = df.eval('10 + 5*treat + u') df['treat'] = df['treat'].sample(frac=1, replace=False).values return linreg(df) # 并行处理所有抽样,仅传递种子而非DataFrame with mp.Pool(processes=4) as pool: seeds = [1029 + k for k in range(1000)] results = pool.map(run_single_draw, seeds) final_est = pd.concat(results, axis=0, sort=False, keys=range(1000), names=['k', 'param'])
4. 选择合适的并行工具
可以尝试用concurrent.futures.ProcessPoolExecutor简化并行逻辑,或用joblib库(针对科学计算优化的并行工具),它能更高效地处理numpy/pandas数据的传递:
from joblib import Parallel, delayed def run_analysis_joblib(draws, seed): np.random.seed(seed) N = 5000 df = pd.DataFrame({'treat':np.random.choice([0,1], size=N, replace=True)}) df['u'] = np.random.normal(size=N) df['y'] = df.eval('10 + 5*treat + u') est = Parallel(n_jobs=4)(delayed(lambda: linreg(shuffle_treat(df)))() for k in range(draws)) return pd.concat(est, axis=0, sort=False, keys=range(draws), names=['k', 'param'])
内容的提问来源于stack exchange,提问作者jtorca

