如何并行化R中的lapply()函数?网格搜索遇性能异常
网格搜索并行化性能暴跌问题排查与解决
问题描述
正在执行网格搜索,寻找能最大化x的线性组合与y之间相关性的系数。核心计算函数corr_grid_search接收一组theta参数和数据集modeling_df,代码如下:
corr_grid_search <- function(thetas, modeling_df) { # thetas = as.vector(thetas) coeff1 = modeling_df$penalty1 / thetas[1] coeff2 = modeling_df$penalty2 / thetas[2] coeff3 = modeling_df$penalty3 / thetas[3] coeff4 = modeling_df$penalty4 / thetas[4] coeff5 = modeling_df$penalty5 / thetas[5] coeff6 = modeling_df$penalty6 / thetas[6] coeff7 = modeling_df$penalty7 / thetas[7] coeff8 = modeling_df$penalty8 / thetas[8] coeff9 = modeling_df$penalty9 / thetas[9] coeff10 = modeling_df$penalty10 / thetas[10] df = data.frame(coeff1, coeff2, coeff3, coeff4, coeff5, coeff6, coeff7, coeff8, coeff9, coeff10) pp_1 = modeling_df$x1 / df$coeff1 pp_2 = modeling_df$x2 / df$coeff2 pp_3 = modeling_df$x3 / df$coeff3 pp_4 = modeling_df$x4 / df$coeff4 pp_5 = modeling_df$x5 / df$coeff5 pp_6 = modeling_df$x6 / df$coeff6 pp_7 = modeling_df$x7 / df$coeff7 pp_8 = modeling_df$x8 / df$coeff8 pp_9 = modeling_df$x9 / df$coeff9 pp_10 = modeling_df$x10 / df$coeff10 recip = 1/df[, c('coeff1', 'coeff2', 'coeff3', 'coeff4', 'coeff5', 'coeff6', 'coeff7', 'coeff8', 'coeff9', 'coeff10')] recip = as.data.frame(lapply(recip, function(x) replace(x, is.infinite(x), NA))) df = data.frame(pp_1, pp_2, pp_3, pp_4, pp_5, pp_6, pp_7, pp_8, pp_9, pp_10) weighted_x = rowSums(df, na.rm=T) / rowSums(recip, na.rm=T) cor(weighted_x[!is.na(weighted_x)], modeling_df[!is.na(weighted_x),]$y) }
使用lapply串行执行时正常:
lapply(blah, corr_grid_search, modeling_df)
但尝试用future.apply和parallel并行化时,运行时间比串行慢2-3个数量级:
library(future.apply) plan(multisession) cors = future_lapply(blah, corr_grid_search, modeling_df)
library(parallel) cl = makeCluster(32) clusterExport(cl=cl, varlist=c("modeling_df")) cors = parLapply(cl, blah, corr_grid_search, modeling_df)
核心问题分析
- 数据传输开销过大:
modeling_df会被重复传输到每个并行进程,若数据集较大,序列化、传输和复制的成本远超过并行计算节省的时间。 - 单任务计算量过小:单个theta迭代的计算逻辑简单,并行调度、进程间通信的开销远大于计算本身,导致整体效率暴跌。
- 函数内部冗余计算:原函数采用逐列手动赋值的方式,没有利用R的向量化特性,单任务计算效率本身就低,并行后问题被放大。
- Worker数量过多:设置32个Worker远超常规CPU核心数,会导致进程调度、上下文切换的开销剧增。
优化方案
1. 减少数据重复传输
- 针对parallel包:确保
modeling_df只被导出到集群节点一次,且函数不再重复接收该参数:library(parallel) # 根据CPU核心数设置Worker数量,比如8 cl = makeCluster(8) # 一次性导出数据到所有节点 clusterExport(cl=cl, varlist=c("modeling_df")) # 修改函数,不再接收modeling_df参数 corr_grid_search <- function(thetas) { # 函数逻辑直接使用全局环境的modeling_df # ...(后续替换为向量化版本) } cors = parLapply(cl, blah, corr_grid_search) stopCluster(cl) - 针对future.apply包:指定Worker数量为CPU核心数,并让
modeling_df在全局环境中共享,避免重复传输:library(future.apply) plan(multisession, workers = 8) # 修改函数,不再接收modeling_df参数 cors = future_lapply(blah, corr_grid_search)
2. 向量化改造函数,提升单任务效率
将原函数的逐列操作改为向量化计算,减少冗余代码,大幅提升单任务速度:
corr_grid_search <- function(thetas) { # 提取penalty和x列 penalties <- modeling_df[, paste0("penalty", 1:10)] xs <- modeling_df[, paste0("x", 1:10)] # 向量化计算coeff矩阵 coeff <- penalties / matrix(thetas, nrow = nrow(penalties), ncol = 10, byrow = TRUE) # 计算pp矩阵 pp <- xs / coeff # 处理无穷值为NA recip <- 1 / coeff recip[is.infinite(recip)] <- NA # 计算加权x weighted_x <- rowSums(pp, na.rm = TRUE) / rowSums(recip, na.rm = TRUE) # 计算相关性 valid_idx <- !is.na(weighted_x) cor(weighted_x[valid_idx], modeling_df$y[valid_idx]) }
3. 调整Worker数量
Worker数量建议设置为CPU物理核心数(一般为8-16,而非32),避免过度调度导致的开销。
4. 批量处理任务
若单个theta计算量仍过小,可将blah分成若干批次,每个批次处理一组theta,减少并行调度次数:
# 把blah分成8组,对应8个Worker blah_batches <- split(blah, cut(1:length(blah), 8, labels = FALSE)) # 并行处理每个批次,批次内部串行计算 cors_batches <- future_lapply(blah_batches, function(batch) { lapply(batch, corr_grid_search) }) # 合并结果 cors <- unlist(cors_batches)
内容的提问来源于stack exchange,提问作者GBPU
相关产品推荐
相关产品推荐

