如何在本地与AWS实例及核心间用future_map实现并行计算
问题:本地机器与远程AWS EC2实例间的并行代码搭建
已成功通过SSH连接AWS实例,但无法正确搭建多级并行代码。此前在本地用plan(multisession)启动8个worker,将.x = bestVars和自定义函数传入future_map的.f = run_sim_in_par()运行正常。现在需要实现本地实例(保持multisession的运行方式)与EC2实例及其核心间的并行计算,相关代码如下:
library(future) library(furrr) library(data.table) library(kit) library(tictoc) library(collapse) plan(multisession) ## 示例数据 vars <- paste0(letters,1:10) bestVars <- combn(vars, 4, simplify = F) df <- data.frame( matrix(data = rnorm(10000*length(vars),200,500), nrow = 10000, ncol = length(vars)) ) names(df) <- vars df$value <- rnorm(n = nrow(df), 350, 300) df <- df %>% dplyr::select(value,everything(.)) df <- lapply(split.default(x = df, names(df)), function(x) x[[1]]) ## 拆分变量列表为块,避免真实数据(>1100万行)出现内存错误 chunks_run <- collapse::rsplit(1:length(bestVars), ceiling(seq_along(bestVars)/1000)) chunks_list <- vector("list", length = length(chunks_run)) ## 本地并行运行变量组合的函数 run_sim_in_par <- function(df, var_to_sim) { sampled_rows <- sample(x = 1:length(df[["value"]]), size = 50, replace = F) varname <- paste(var_to_sim, collapse = "*") best <- Reduce(df[var_to_sim], f = '*')[sampled_rows] row_idx <- kit::topn(best, n = 5, decreasing = T, hasna = FALSE, index = TRUE) best_row_value <- df[["value"]][sampled_rows][row_idx] sim <- data.table(var = varname, mean_value = mean(best_row_value)) return(sim) } chunks_idx <- 1 for (chunks_idx in seq_along(chunks_run)) { tic() simulated_roi <- future_map( .x = bestVars[ chunks_run[[chunks_idx]] ], .f = function(x) run_sim_in_par(df = df[c("value", x)], var_to_sim = x)) toc() chunks_list[[chunks_idx]] <- rbindlist(simulated_roi) } simulated_roi <- rbindlist(chunks_list) base::closeAllConnections() ## 启动远程集群并手动连接(remote=FALSE会挂起无法连接) public_ip <- "x.xx.aa.bc" ssh_private_key_file <- "path/r.pem" cl <- makeClusterPSOCK( # EC2实例的公网IP workers = public_ip, # 用户名(固定为'ubuntu') user = "ubuntu", ## 使用AWS注册的私钥 rshopts = c( "-o", "StrictHostKeyChecking=no", "-o", "IdentitiesOnly=yes", "-i", ssh_private_key_file ), rscript_args = c( ## 为'ubuntu'用户设置.libPaths() "-e", shQuote( paste0( "local({", "p <- Sys.getenv('R_LIBS_USER'); ", "dir.create(p, recursive = TRUE, showWarnings = FALSE); ", ".libPaths(p)", "})" ) ), ## 安装所需包 "-e", shQuote("install.packages(c('furrr','purrr','kit','dqrng','tictoc','data.table'))") ), # 设为TRUE可查看在worker上运行的代码,无需建立连接 dryrun = FALSE, manual = TRUE ) cl ## socket cluster with 1 nodes on host ‘xx.xxx.xx.xxx’ ## 尝试在本地和AWS集群上运行代码 ## 似乎本地和远程集群都没有处理任务 local_workers <- makeClusterPSOCK(8) plan(list(tweak(cluster, workers = c(cl, local_workers)), multisession)) chunks_idx <- 1 for (chunks_idx in seq_along(chunks_run)) { tic() simulated_roi <- future_map( .x = c(1:12), # 8个本地,4个远程 .f = ~{ future_map( .x = bestVars[ chunks_run[[chunks_idx]] ], .f = function(x) run_sim_in_par(df = df[c("value", x)], var_to_sim = x) ) } ) toc() chunks_list[[chunks_idx]] <- rbindlist(simulated_roi) }
内容的提问来源于stack exchange,提问作者On_an_island
相关产品推荐
相关产品推荐

