【问题标题】:How to run another Rscript after several R jobs running in parallel are done?完成多个并行运行的 R 作业后如何运行另一个 Rscript?
【发布时间】:2020-12-17 10:56:11
【问题描述】:

我需要如何运行我的脚本的安排是首先使用rstudioapi::jobRunScript() 函数并行运行 4 个 R 脚本。并行运行的每个脚本不会从任何环境导入任何内容,而是将创建的数据帧导出到全局环境。 我的第 5 个 R 脚本建立在由并行运行的 4 个 R 脚本创建的数据帧之上,而且这个第 5 个脚本也在控制台中运行。如果有办法在后台运行第 5 个脚本,而不是在前 4 个 R 脚本并行运行后在控制台中运行,那会好很多。我也在尝试减少整个过程的总运行时间。

虽然我能够弄清楚如何并行运行前 4 个 R 脚本,但我的任务还没有完全完成,因为我找不到触发我的第 5 个 R 脚本运行的方法。希望大家能帮帮我

【问题讨论】:

  • 我认为使用jobRunScript 运行的作业不会影响(例如,将结果存储在)调用环境。前 4 个工作是否在外部保存结果? (例如文件或数据库)。
  • ~~~jobRunScript()是什么?~~~哦,我明白了,是rstudioapi::jobRunScript()
  • 嗨@r2evans,结果只是在全球环境中浮动。我是这样做的: 'jobRunScript("C:\\Desktop\\CY alloc\\Scripts\\Prep 1.R", importEnv=FALSE, exportEnv="R_GlobalEnv")' 'jobRunScript("C:\\Desktop \\CY alloc\\Scripts\\Prep 2.R", importEnv=FALSE, exportEnv="R_GlobalEnv")'

标签: r windows parallel-processing rscript rstudioapi


【解决方案1】:

这对我来说有点太开放了。虽然rstudioapi 绝对可以用于运行并行任务,但它不是很通用,也不会为您提供非常有用的输出。 parallel Universe 在 R 中得到了很好的实现,其中有几个包提供了更简单和更好的接口来执行此操作。这里有 3 个选项,它们还允许从不同文件中“输出”某些内容。

包 = 并行

使用并行包,我们可以非常简单地实现这一点。只需创建要获取的文件向量并在每个线程中执行source。主进程会在它们运行时锁定,但如果你必须等待它们完成,这并不重要。

library(parallel)
ncpu <- detectCores()
cl <- makeCluster(ncpu)
# full path to file that should execute 
files <- c(...) 
# use an lapply in parallel.
result <- parLapply(cl, files, source)
# Remember to close the cluster
stopCluster(cl)
# If anything is returned this can now be used.

附带说明一下,有几个包具有与 parallel 包类似的接口,后者是在 snow 包的基础上构建的,因此了解它是一个很好的基准。

包 = foreach

parallel 包的替代方案是foreach 包,它提供了类似于for-loop 接口的功能,简化了接口,同时提供了更大的灵活性并自动导入必要的库和变量(尽管这样更安全)手动执行此操作)。
foreach 包确实依赖于 paralleldoParallel 包来设置集群

library(parallel)
library(doParallel)
library(foreach)
ncpu <- detectCores()
cl <- makeCluster(ncpu)
files <- c(...) 
registerDoParallel(cl)
# Run parallel using foreach
# remember %dopar% for parallel. %do% for sequential.
result <- foreach(file = files, .combine = list, .multicombine = TRUE) %dopar% { 
  source(file)
  # Add any code before or after source.
}
# Stop cluster
stopCluster(cl)
# Do more stuff. Result holds any result returned by foreach.

虽然它确实添加了几行代码,但 .combine.packages.export 提供了一个非常简单的接口,可以在 R 中使用并行计算。

包装 = 未来

现在这是比较少用的软件包之一。 future 提供了一个比parallelforeach 更灵活的并行接口,允许异步并行编程。然而,实现看起来有点令人生畏,而我在下面提供的示例只是触及了可能的表面。
另外值得一提的是,虽然future 包确实提供了运行代码所需的函数和包的自动导入,但经验让我意识到这仅限于任何调用中的第一级深度(有时更少),因此仍然需要导出。
虽然foreach 依赖parallel(或类似的)来启动集群,但foreach 将使用所有可用内核自行启动。对plan(multiprocess) 的简单调用将启动一个多核会话。

library(future)
files <- c(...) 
# Start multiprocess session
plan(multiprocess)
# Simple wrapper function, so we can iterate over the files variable easier
source_future <- function(file)
  future(file)
results <- lapply(files, source_future)
# Do some calculations in the meantime
print('hello world, I am running while waiting for the futures to finish')
# Force waiting for the futures to finish
resolve(results)
# Extract any result from the futures
results <- values(results)
# Clean up the process (close down clusters)
plan(sequential)
# Run some more code.

现在这可能看起来很重,但一般机制是:

  1. 致电plan(multiprocess)
  2. 使用future(或%&lt;-%,我不会介绍)执行一些函数
  3. 如果您有更多代码要运行,则执行其他操作,这不依赖于进程
  4. 使用 resolve 等待结果,它适用于列表(或环境)中的单个未来或多个未来
  5. 使用value 收集结果,用于单个期货或values 用于列表(或环境)中的多个期货
  6. 使用plan(sequential) 清除在future 环境中运行的任何集群
  7. 继续使用取决于未来结果的代码。

我相信这 3 个包为任何用户需要与之交互的多处理的每个必要元素(至少在 CPU 上)提供了接口。其他包提供替代接口,而对于异步我只知道futurepromises。一般来说,我建议大多数用户在转向异步编程时要非常小心,因为与同步并行编程相比,这可能会导致一整套不太常见的问题。

我希望这可能有助于为(非常有限的)rstudioapi 接口提供一个替代方案,我相当肯定它从来没有打算由用户自己用于并行编程,但更有可能用于执行任务,例如由接口本身并行构建一个包。

【讨论】:

  • 关于“package = future: ...自动导入运行代码所需的函数和包,经验让我意识到这仅限于任何调用的第一级深度(有时更少),因为这样的出口仍然是必要的。”请考虑将可重现的示例发布到未来的问题跟踪器,因为目标是涵盖尽可能多的用例(模块棘手的极端案例)。
  • @HenrikB,我会记住的。它在我的工作中出现过几次,但目前我手头没有可重现的例子。谁知道他们可能在此期间已经修复了。
  • 嗨奥利弗,谢谢。关于 package = parallel,我是否需要在每个并行运行的脚本末尾使用 return 函数,以便“结果”对象可以捕获生成的数据帧?
  • 嗨@Danelle。不,没有必要。如果没有要返回的内容,则默认返回 NULL。如果您也没有任何东西要返回,您可以删除 result &lt;-。这适用于所有 3 个解决方案。 :-)
  • 知道了,@Oliver!并行包对我有用,并且在我当前的代码结构中是最直接和最容易实现的。感谢您的回复!
【解决方案2】:

您可以将promisesfuture 结合使用:promises::promise_all 后跟promises::then 允许等待前一个futures 完成,然后再将最后一个futures 作为后台进程启动。
当作业在后台运行时,您可以控制控制台。

library(future)
library(promises)

plan(multisession)
# Job1
fJob1<- future({
  # Simulate a short duration job
  Sys.sleep(3)
  cat('Job1 done \n')
})

# Job2
fJob2<- future({
  # Simulate a medium duration job
  Sys.sleep(6)
  cat('Job2 done \n')
})


# Job3
fJob3<- future({
  # Simulate a long duration job
  Sys.sleep(10)
  cat('Job3 done \n')
})

# last Job
runLastJob <- function(res) {
  cat('Last job launched \n')
  # launch here script for last job
}

# Cancel last Job
cancelLastJob <- function(res) {
  cat('Last job not launched \n')
}

#  Wait for all jobs to be completed and launch last job
p_wait_all <- promises::promise_all(fJob1, fJob2, fJob3 )
promises::then(p_wait_all,  onFulfilled = runLastJob, onRejected = cancelLastJob)

Job1 done 
Job2 done 
Job3 done 
Last job launched 

【讨论】:

    【解决方案3】:

    我不知道这对您当前的场景有多大的适应性,但这里有一种方法可以让四个东西并行运行,获取它们的返回值,然后触发第五个表达式/函数。

    前提是使用callr::r_bg 来运行单个文件。这实际上运行的是一个function,而不是一个文件,所以我将修改这些文件的预期外观一点

    我将编写一个辅助脚本,旨在模仿您的四个脚本之一。我猜你也希望能够正常获取它(直接运行它而不是作为函数运行),所以我将生成脚本文件,以便它“知道”它是被获取还是直接运行(基于Rscript detect if R script is being called/sourced from another script)。 (如果你知道python,这类似于python的if __name__ == "__main__"技巧。)

    名为somescript.R的辅助脚本。

    somefunc <- function(seconds) {
      # put the contents of a script file in this function, and have
      # it return() the data you need back in the calling environment
      Sys.sleep(seconds)
      return(mtcars[sample(nrow(mtcars),2),1:3])
    }
    
    if (sys.nframe() == 0L) {
      # if we're here, the script is being Rscript'd, not source'd
      somefunc(3)
    }
    

    作为演示,如果source'd 在控制台上,这个只是定义了函数(或多个,如果你愿意),它不会执行最后一个if内的代码块:

    system.time(source("~/StackOverflow/14182669/somescript.R"))
    #                                   # <--- no output, it did not return a sample from mtcars
    #    user  system elapsed 
    #       0       0       0           # <--- no time passed
    

    但如果我在终端中使用Rscript 运行它,

    $ time /c/R/R-4.0.2/bin/x64/Rscript somescript.R
                   mpg cyl  disp
    Merc 280C     17.8   6 167.6
    Mazda RX4 Wag 21.0   6 160.0
    
    real    0m3.394s                    # <--- 3 second sleep
    user    0m0.000s
    sys     0m0.015s
    

    回到前提。而不是四个“脚本”,重写你的脚本文件,就像我上面的somescript.R。如果操作正确,它们可以是 Rscripted 以及 sourced,具有不同的意图。

    我将使用这个脚本四次而不是四个脚本。这是我们想要自动化的手动运行:

    # library(callr)
    tasks <- list(
      callr::r_bg(somefunc, args = list(5)),
      callr::r_bg(somefunc, args = list(1)),
      callr::r_bg(somefunc, args = list(10)),
      callr::r_bg(somefunc, args = list(3))
    )
    sapply(tasks, function(tk) tk$is_alive())
    # [1]  TRUE FALSE  TRUE FALSE
    ### time passes
    sapply(tasks, function(tk) tk$is_alive())
    # [1] FALSE FALSE  TRUE FALSE
    sapply(tasks, function(tk) tk$is_alive())
    # [1] FALSE FALSE FALSE FALSE
    
    tasks[[1]]$get_result()
    #                    mpg cyl  disp  hp drat    wt  qsec vs am gear carb
    # Merc 280          19.2   6 167.6 123 3.92 3.440 18.30  1  0    4    4
    # Chrysler Imperial 14.7   8 440.0 230 3.23 5.345 17.42  0  0    3    4
    

    我们可以使用

    source("somescript.R")
    message(Sys.time(), " starting")
    # 2020-08-28 07:45:31 starting
    tasks <- list(
      callr::r_bg(somefunc, args = list(5)),
      callr::r_bg(somefunc, args = list(1)),
      callr::r_bg(somefunc, args = list(10)),
      callr::r_bg(somefunc, args = list(3))
    )
    # some reasonable time-between-checks
    while (any(sapply(tasks, function(tk) tk$is_alive()))) {
      message(Sys.time(), " still waiting")
      Sys.sleep(1)                      # <-- over to you for a reasonable poll interval
    }
    # 2020-08-28 07:45:32 still waiting
    # 2020-08-28 07:45:33 still waiting
    # 2020-08-28 07:45:34 still waiting
    # 2020-08-28 07:45:35 still waiting
    # 2020-08-28 07:45:36 still waiting
    # 2020-08-28 07:45:37 still waiting
    # 2020-08-28 07:45:38 still waiting
    # 2020-08-28 07:45:39 still waiting
    # 2020-08-28 07:45:40 still waiting
    # 2020-08-28 07:45:41 still waiting
    message(Sys.time(), " done!")
    # 2020-08-28 07:45:43 done!
    results <- lapply(tasks, function(tk) tk$get_result())
    str(results)
    # List of 4
    #  $ :'data.frame': 2 obs. of  3 variables:
    #   ..$ mpg : num [1:2] 24.4 32.4
    #   ..$ cyl : num [1:2] 4 4
    #   ..$ disp: num [1:2] 146.7 78.7
    #  $ :'data.frame': 2 obs. of  3 variables:
    #   ..$ mpg : num [1:2] 30.4 14.3
    #   ..$ cyl : num [1:2] 4 8
    #   ..$ disp: num [1:2] 95.1 360
    #  $ :'data.frame': 2 obs. of  3 variables:
    #   ..$ mpg : num [1:2] 15.2 15.8
    #   ..$ cyl : num [1:2] 8 8
    #   ..$ disp: num [1:2] 276 351
    #  $ :'data.frame': 2 obs. of  3 variables:
    #   ..$ mpg : num [1:2] 14.3 15.2
    #   ..$ cyl : num [1:2] 8 8
    #   ..$ disp: num [1:2] 360 304
    

    现在运行你的第五个函数/脚本。

    【讨论】:

      猜你喜欢
      • 2011-09-27
      • 2019-05-24
      • 1970-01-01
      • 1970-01-01
      • 2016-08-27
      • 1970-01-01
      • 1970-01-01
      • 2019-01-18
      • 2019-03-26
      相关资源
      最近更新 更多