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

如何在本地与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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.05 01:17:44