R中parallel与foreach嵌套循环并行处理失效求助
问题:使用
foreach实现嵌套并行循环时无法得到预期结果 我尝试用R的parallel和foreach包实现嵌套循环并行,提升运行速度。数据框首列为分类变量,其余为连续型变量,目标是遍历所有无重复的两个连续型变量组合,以分类变量为响应变量执行caret::bagFDA函数。
常规嵌套循环可正常运行(代码如下),但并行版本运行无报错却无法填充结果数据框,results始终为NULL,公式显示为(unknown)。
常规串行代码(可正常运行)
#Install packages #install.packages("caret") #install.packages("dplyr") start.time <- Sys.time() #Libraries library(caret) library(dplyr) #Read iris data data(iris) iris <- iris[c(ncol(iris), 1:(ncol(iris) - 1))] #Sampling 50/50 for train and valid set.seed(2000) Train <- iris %>% group_by(Species) %>% sample_frac(.5, replace = FALSE) Valid <- anti_join(iris, Train) # bagFDA BAGFDA <- data.frame(Variable_Name1 = character(0), Variable_Name2 = character(0), Accuracy = numeric(0), Kappa = numeric(0)) set.seed(3000) for(i in 2:(ncol(Train) - 1)) { for (j in (i + 1):ncol(Train)) { tryCatch({ formula <- as.formula(paste("as.factor(Species) ~", names(Train)[i], "+", names(Train)[j])) bag <- caret::bagFDA(formula, data = Train) bag_predict <- predict(bag, newdata = Valid) bag_CM <- confusionMatrix(bag_predict, Valid$Species) iteration_results <- data.frame( Variable_Name1 = names(Train)[i], Variable_Name2 = names(Train)[j], Accuracy = bag_CM$overall["Accuracy"], Kappa = bag_CM$overall["Kappa"] ) BAGFDA <- rbind(BAGFDA, iteration_results) print("Good") }, error = function(e) { cat("ERROR:", conditionMessage(e), "\n") }) } } print(BAGFDA) end.time <- Sys.time() time.taken <- round(end.time - start.time,2) time.taken
并行版代码(存在问题)
#Install packages #install.packages("foreach") #install.packages("doParallel") #install.packages("caret") #install.packages("dplyr") start.time <- Sys.time() #Libraries library(foreach) library(doParallel) library(caret) library(dplyr) #Read iris data data(iris) iris <- iris[c(ncol(iris), 1:(ncol(iris) - 1))] #Creating clusters cores = parallel::detectCores() - 1 cluster = parallel::makeCluster(cores, type = "PSOCK") doParallel::registerDoParallel(cluster) if (!foreach::getDoParRegistered()) { print("ERROR") } print(foreach::getDoParWorkers()) #Sampling 50/50 for train and valid set.seed(2000) Train <- iris %>% group_by(Species) %>% sample_frac(.5, replace = FALSE) Valid <- anti_join(iris, Train) # bagFDA BAGFDA <- data.frame(Variable_Name1 = character(0), Variable_Name2 = character(0), Accuracy = numeric(0), Kappa = numeric(0)) set.seed(1001) results <- foreach(i = 2:(ncol(Train) - 1), .combine='cbind') %:% foreach (j = (i + 1):ncol(Train), .combine='c') %dopar% tryCatch({ formula <- as.formula(paste("as.factor(Species) ~", names(Train)[i], "+", names(Train)[j])) bag <- caret::bagFDA(formula, data = Train) bag_predict <- predict(bag, newdata = Valid) bag_CM <- confusionMatrix(bag_predict, Valid$Species) iteration_results <- data.frame( Variable_Name1 = names(Train)[i], Variable_Name2 = names(Train)[j], Accuracy = bag_CM$overall["Accuracy"], Kappa = bag_CM$overall["Kappa"] ) BAGFDA <- rbind(BAGFDA, iteration_results) print("Good") }, error = function(e) { cat("ERROR:", conditionMessage(e), "\n") }) print(BAGFDA) stopCluster(cluster) end.time <- Sys.time() time.taken <- round(end.time - start.time,2) time.taken
问题原因分析
- 并行环境变量作用域:PSOCK集群的子进程有独立内存空间,主进程的
BAGFDA无法被子进程修改,子进程内的rbind操作仅在自身内存生效,不会同步到主进程。 foreach合并方式错误:嵌套循环使用cbind和c的合并逻辑不符合数据框需求,无法正确聚合结果。- 数据与包未同步到子进程:子进程未加载所需包,也未获取
Train、Valid数据,导致公式生成和模型训练异常。 - 错误处理输出无效:子进程的
cat输出无法传递到主进程,错误信息丢失,也未返回有效结果。
修正后的并行代码
#Install packages #install.packages("foreach") #install.packages("doParallel") #install.packages("caret") #install.packages("dplyr") start.time <- Sys.time() #Libraries library(foreach) library(doParallel) library(caret) library(dplyr) #Read iris data data(iris) iris <- iris[c(ncol(iris), 1:(ncol(iris) - 1))] #Sampling 50/50 for train and valid set.seed(2000) Train <- iris %>% group_by(Species) %>% sample_frac(.5, replace = FALSE) Valid <- anti_join(iris, Train) #Creating clusters and sync resources cores = parallel::detectCores() - 1 cluster = parallel::makeCluster(cores, type = "PSOCK") doParallel::registerDoParallel(cluster) # 导出数据到子进程 parallel::clusterExport(cluster, c("Train", "Valid")) # 子进程加载必要包 parallel::clusterEvalQ(cluster, { library(caret) library(dplyr) }) # 预先生成所有无重复变量对,简化并行逻辑 var_pairs <- expand.grid( i = 2:(ncol(Train)-1), j = (2+1):ncol(Train), stringsAsFactors = FALSE ) %>% filter(i < j) set.seed(1001) results <- foreach(k = 1:nrow(var_pairs), .combine='rbind', .packages=c("caret", "dplyr")) %dopar% { i <- var_pairs$i[k] j <- var_pairs$j[k] tryCatch({ var1 <- names(Train)[i] var2 <- names(Train)[j] formula <- as.formula(paste("as.factor(Species) ~", var1, "+", var2)) bag <- caret::bagFDA(formula, data = Train) bag_predict <- predict(bag, newdata = Valid) bag_CM <- confusionMatrix(bag_predict, Valid$Species) data.frame( Variable_Name1 = var1, Variable_Name2 = var2, Accuracy = bag_CM$overall["Accuracy"], Kappa = bag_CM$overall["Kappa"], Status = "Success" ) }, error = function(e) { data.frame( Variable_Name1 = names(Train)[i], Variable_Name2 = names(Train)[j], Accuracy = NA_real_, Kappa = NA_real_, Status = paste("Error:", conditionMessage(e)) ) }) } # 输出结果 print(results) stopCluster(cluster) end.time <- Sys.time() time.taken <- round(end.time - start.time,2) time.taken
修正说明
- 扁平化循环逻辑:预先生成所有变量对,避免嵌套
foreach的复杂合并,让并行逻辑更清晰。 - 同步资源到子进程:用
clusterExport传递数据,clusterEvalQ加载包,确保子进程环境与主进程一致。 - 正确合并结果:使用
rbind合并每个循环返回的数据框,直接生成最终结果集。 - 优化错误处理:让
tryCatch返回包含错误信息的行,便于后续排查,同时保证结果结构统一。 - 移除无效全局操作:不再修改主进程的
BAGFDA,直接通过foreach聚合子进程返回的结果。
内容的提问来源于stack exchange,提问作者Rodrigo Nehara
相关产品推荐
相关产品推荐

