AWS EC2服务器上R的parallel foreach多核提速失效问题排查
我在博士阶段开发了一个模型并开展模拟研究,一共设置了64种由不同参数与样本量(1000或10000)组合的模拟场景,每种场景需要重复执行500次。办公电脑为8核,处理样本量10000的场景需1-15小时,但每晚系统超时导致无法完成,因此搭建了48核的AWS EC2服务器,预期速度能达到办公电脑的6倍。但实际运行中,小数据集场景提速明显,样本量10000的场景反而慢很多。附上R代码,想请教:并行模拟是否存在错误?为什么多核不仅没提速反而变慢?
library(Rlab) # for rbern function library(matrixStats) # for colSD and colMean functions library(dplyr) library(foreach) # for parallelisation to speed up simulations library(doSNOW) library(progress) # for progress bar source("/home/ec2-user/functions_fit.R") source("/home/ec2-user/functions_simulator.R") # Run simulation study for this combination of parameters simulation_study <- function(reps, i){ # Define simulation settings for this iteration simulation_setting <- sim_params[i,] # Initialize cluster with 32 cores cl <- makeCluster(48) registerDoSNOW(cl) #registerDoParallel(cl) # Make progress bar for this simulation setting pb <- progress_bar$new(format = "Iteration = :letter [:bar] :elapsed | eta: :eta", total = reps, width = 100) progress <- function(n){ pb$tick(tokens = list(letter = n)) } opts <- list(progress = progress) # Start parallel simulation for this simulation setting res <- foreach(k=1:reps, .options.snow = opts, .packages = c('Rlab', 'dplyr'), .export = c("model.fit", "model.simulator" ), .errorhandling = 'remove') %dopar% { sim.dat <- model.simulator(simulation.setting) time <- system.time(result <- model.fit(sim.dat)) lower <- result[["summary"]][["lower"]] # lower confidence interval upper <- result[["summary"]][["upper"]] # upper confidence interval c(result$theta.hat, lower, upper, time[1]) } stopCluster(cl) results <- data.frame(matrix(unlist(res), nrow=reps, byrow=TRUE)) length.theta <- length(true.theta) avg.time <- mean(results[, ncol(results)], na.rm=T) # average time for model to converge convergence <- sum(sapply(1:reps, function(x) ifelse(any(is.na(results[x, 1:(length.theta)])), 0, 1)))/reps*100 # convergence rate estimates <- results[, 1:(length.theta)] # columns with estimates lowers <- results[, (length.theta+1):(2*length.theta)] # columns with lower CI uppers <- results[, (2*length.theta+1):(3*length.theta)] # columns with upper CI cps <- rep(0, length.theta) # coverage probability for (i in 1:length.theta){ cps[i] <- mean(ifelse((lowers[,i] <= true.theta[i]) & (uppers[,i] >= true.theta[i]), 1, 0), na.rm=T) # calculate coverage prob } theta.med <- colMedians(matrix(unlist(estimates), nrow=reps, ncol = length.theta, byrow = F), na.rm = T) # median estiamtes # return results sim.results <- round(data.frame(Actual = true.theta, Estimate = theta.med, Bias = (theta.med - true.theta)/true.theta*100, Std.Dev = colSds(as.matrix(estimates, nrow=reps), na.rm=T), CP = cps, Avg.Time = avg.time, Convergence = convergence),4) return(sim.results) } # Run each of the 64 simulation settings sim_results <- list() for (i in 1:64){ sim_params[i,] sim_results[[i]] <- simulation_study(500, i) print(sim_results[[i]]) }
注:model.fit和model.simulator定义在R文件functions_fit.R和functions_simulator.R中
1. 集群反复创建销毁的巨大开销
你在simulation_study函数内,每处理一个场景就创建一次48核集群,处理完立即销毁。集群的初始化和销毁本身存在不小的系统开销,小样本场景单次任务耗时短,这个开销占比低,所以仍能体现提速效果;但大样本场景单次任务耗时久,反复创建销毁集群的额外成本会被放大,直接拖慢整体运行速度。
2. 变量传递错误与冗余开销
- 代码中定义的变量是
simulation_setting(下划线),但foreach循环内用的是simulation.setting(点),这会导致每个并行任务都要去全局环境查找变量,增加了不必要的跨进程通信开销。大样本场景下数据体积大,这个错误带来的通信成本会被进一步放大。 .packages和.export参数在每次调用foreach时都会重复加载包和传递函数,没有复用集群资源,冗余开销在多核大负载下会被放大。
3. 进度条的同步开销
使用doSNOW的进度条,每个并行任务完成后都要和主进程通信更新进度,大样本场景下虽然任务完成频率低,但多核高负载时,这种同步通信会抢占计算资源,反而拖慢运行效率。
4. 内存资源瓶颈
大样本(10000)的模拟和模型拟合本身内存占用高,48核同时运行时,每个核都要加载大样本数据、进行模型计算,如果EC2实例的内存不足以支撑48个大样本任务同时运行,就会触发磁盘交换(swap),磁盘读写速度远低于内存,导致整体运行速度急剧下降。小样本内存占用低,不会触发这个问题,因此提速正常。
1. 全局初始化集群,避免反复创建销毁
把集群创建移到外层循环外,所有场景共用一个集群,处理完所有场景再销毁:
# 全局初始化一次集群,建议用detectCores()自动识别可用核心,避免超配 cl <- makeCluster(min(48, detectCores())) registerDoSNOW(cl) # Run each of the 64 simulation settings sim_results <- list() for (i in 1:64){ sim_results[[i]] <- simulation_study(500, i) print(sim_results[[i]]) } # 所有场景处理完再销毁集群 stopCluster(cl)
同时修改simulation_study函数,删除其中的makeCluster和stopCluster代码。
2. 修正变量名,优化资源传递
- 把
foreach内的simulation.setting改为simulation_setting,修正变量名错误。 - 在外层集群初始化后,提前导出全局变量和加载依赖包,避免
foreach重复传递:
# 提前导出全局变量和函数到集群节点 clusterExport(cl, c("sim_params", "model.fit", "model.simulator")) # 提前在集群节点加载依赖包 clusterEvalQ(cl, { library(Rlab) library(dplyr) })
修改后foreach可以去掉.packages和.export参数,减少冗余开销。
3. 匹配核心数与内存资源
如果EC2实例内存有限,不要用满48核,比如先尝试24核,观察内存使用情况,找到核心数与内存的平衡点。大样本场景下,核心数过多反而会因内存不足导致减速。
4. 移除或优化进度条
如果进度条不是必需的,直接去掉.options.snow = opts参数,减少主从进程的通信开销;如果需要进度条,可以改用主进程定期输出进度的方式,而非每个任务同步更新。
5. 优化模型拟合函数的内存效率
检查model.fit函数,避免不必要的大对象复制,尝试用更高效的数据结构(比如data.table代替data.frame),降低大样本场景下的内存占用,缓解多核运行时的内存压力。
内容的提问来源于stack exchange,提问作者kenny

