如何在R中使用doParallel替代apply实现data.frame逐行运算并行处理
Windows/Unix 通用逐行并行处理解决方案
你不需要为每一行生成独立R文件,直接基于你已有的doParallel使用经验,调整foreach迭代逻辑即可适配需求,完整可运行代码如下:
# 加载依赖包 library(doParallel) # ---------------------- # 1. 准备你的原数据和自定义函数(和你现有代码完全一致) # ---------------------- master.iter <- read.table(text = ' scenario aaa bbb ccc ddd eee 1 1 5 0 20 10 2 1 10 0 2000 1000 ', header = TRUE, stringsAsFactors = FALSE) master.function <- function(scenario, aaa, bbb, ccc, ddd, eee) { scenario <- as.numeric(c(scenario)) aaa <- as.numeric(c(aaa)) bbb <- as.numeric(c(bbb)) ccc <- as.numeric(c(ccc)) ddd <- as.numeric(c(ddd)) eee <- as.numeric(c(eee)) AAA <- seq(aaa,bbb,1) BBB <- AAA * ddd CCC <- AAA * eee my.table <- data.frame(AAA = AAA, BBB = BBB, CCC = CCC) output.list <- list(scenario = scenario, aaa = aaa, bbb = bbb, ccc = ccc, ddd = ddd, eee = eee, my.table = my.table) master_output <- do.call(cbind, output.list) return = list(master_output = master_output) } # 你预期的对照结果,用于验证 desired.result <- read.table(text = ' scenario aaa bbb ccc ddd eee my.table.AAA my.table.BBB my.table.CCC 1 1 5 0 20 10 1 20 10 1 1 5 0 20 10 2 40 20 1 1 5 0 20 10 3 60 30 1 1 5 0 20 10 4 80 40 1 1 5 0 20 10 5 100 50 2 1 10 0 2000 1000 1 2000 1000 2 1 10 0 2000 1000 2 4000 2000 2 1 10 0 2000 1000 3 6000 3000 2 1 10 0 2000 1000 4 8000 4000 2 1 10 0 2000 1000 5 10000 5000 2 1 10 0 2000 1000 6 12000 6000 2 1 10 0 2000 1000 7 14000 7000 2 1 10 0 2000 1000 8 16000 8000 2 1 10 0 2000 1000 9 18000 9000 2 1 10 0 2000 1000 10 20000 10000 ', header = TRUE) # ---------------------- # 2. 配置并行集群并执行计算 # ---------------------- # 自动检测可用核心,也可手动指定数值 n_cores <- detectCores() # 创建集群,Windows默认用PSOCK模式,Unix可加type="FORK"提升效率,不加也兼容 my.cluster <- makeCluster(n_cores) registerDoParallel(my.cluster) start.time <- Sys.time() # foreach逐行迭代计算 master.df_parallel <- foreach( i = 1:nrow(master.iter), .combine = 'rbind', # 自动将所有迭代结果按行合并 .errorhandling = 'remove', # 出错跳过对应行,也可设为"stop"直接终止 .export = 'master.function' # 必须导出自定义函数到worker进程,Windows下必填 ) %dopar% { # 取第i行的参数传入自定义函数 row_data <- master.iter[i,] res <- master.function( scenario = row_data$scenario, aaa = row_data$aaa, bbb = row_data$bbb, ccc = row_data$ccc, ddd = row_data$ddd, eee = row_data$eee ) # 直接返回需要合并的结果表 res$master_output } # 关闭集群释放资源 stopCluster(my.cluster) # 计算耗时 end.time <- Sys.time() total.time <- end.time - start.time print(total.time) # 验证结果和预期完全一致,执行后会返回TRUE all.equal(master.df_parallel, desired.result)
核心说明
- 跨系统兼容:该代码在Windows10和Unix集群都可以直接运行,无需修改。如果是Unix环境,创建集群时可以指定
makeCluster(n_cores, type = "FORK"),FORK模式会共享父进程内存,不需要额外导出变量和函数,运行效率更高。 - 扩展适配:如果后续你的自定义函数依赖其他全局变量、第三方包,只需要在foreach的参数里补充:
- 依赖包加
.packages = c("包名1", "包名2") - 依赖的全局变量加进
.export = c("master.function", "你的全局变量名")
- 依赖包加
- 性能优化:如果你的数据集行数非常多,可以通过调整
chunksize参数调整单次分配给worker的任务量,减少进程间通信开销。
内容的提问来源于stack exchange,提问作者Mark Miller
相关产品推荐
相关产品推荐

