You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

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

问题原因分析

  1. 并行环境变量作用域:PSOCK集群的子进程有独立内存空间,主进程的BAGFDA无法被子进程修改,子进程内的rbind操作仅在自身内存生效,不会同步到主进程。
  2. foreach合并方式错误:嵌套循环使用cbind和c的合并逻辑不符合数据框需求,无法正确聚合结果。
  3. 数据与包未同步到子进程:子进程未加载所需包,也未获取Train、Valid数据,导致公式生成和模型训练异常。
  4. 错误处理输出无效:子进程的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

修正说明

  1. 扁平化循环逻辑:预先生成所有变量对,避免嵌套foreach的复杂合并,让并行逻辑更清晰。
  2. 同步资源到子进程:用clusterExport传递数据,clusterEvalQ加载包,确保子进程环境与主进程一致。
  3. 正确合并结果:使用rbind合并每个循环返回的数据框,直接生成最终结果集。
  4. 优化错误处理:让tryCatch返回包含错误信息的行,便于后续排查,同时保证结果结构统一。
  5. 移除无效全局操作:不再修改主进程的BAGFDA,直接通过foreach聚合子进程返回的结果。

内容的提问来源于stack exchange,提问作者Rodrigo Nehara

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.06.24 08:52:01