R语言中如何实现带优先级与依赖关系的异步future任务调度
R 多任务依赖调度需求实现方案
背景
以下是典型的分析工作流:
该流程的顺序执行代码如下:
## some simple functions to mimic longer running models m1 <- \() { Sys.sleep(4*15) "model 1 results" } m2 <- \() { Sys.sleep(3*15) "model 2 results" } m3 <- \() { Sys.sleep(2*15) "model 3 results" } process <- \(text) { Sys.sleep(0.5 * 10) sprintf("processed %s", text) } ###### Approach 1 ~155.2 seconds --- sequential ###### start <- proc.time() ## model 1 m1res <- m1() p1res <- process(m1res) cat(p1res, fill = TRUE) ## model 2 m2res <- m2() p2res <- process(m2res) cat(p2res, fill = TRUE) ## model 3 m3res <- m3() p3res <- process(m3res) cat(p3res, fill = TRUE) stop <- proc.time() ## total run time stop - start
顺序执行总耗时约155.2秒,我们可以使用future包提速,假设可用CPU为2核:
###### Approach 2 ~102.7 seconds --- models not blocked but not optimal run ###### ## load future package and setup a multisession with 2 workers library(future) plan("multisession", workers = 2) start <- proc.time() ## fit models p1res <- future({ model1 <- m1() process(model1) }) p2res <- future({ model2 <- m2() process(model2) }) p3res <- future({ model3 <- m3() process(model3) }) ## present results cat(value(p1res), fill = TRUE) cat(value(p2res), fill = TRUE) cat(value(p3res), fill = TRUE) stop <- proc.time() ## total run time stop - start
该方案耗时降至约102.7秒,但仍未达到最优调度。理想的调度逻辑为:
- CPU 1:运行模型1、处理模型1结果、处理模型2结果
- CPU 2:运行模型2、运行模型3、处理模型3结果
###### Approach 3 ~82.7 seconds --- models not blocked but not optimal run ###### start <- proc.time() ## fit models m1res %<-% m1() m2res %<-% m2() m3res %<-% m3() ## grab results p1res %<-% process(m1res) p2res %<-% process(m2res) p3res %<-% process(m3res) ## present results cat(p1res, fill = TRUE) cat(p2res, fill = TRUE) cat(p3res, fill = TRUE) stop <- proc.time() ## total run time stop - start
上述方案实现了最优调度,耗时仅约82.7秒,但代码组织清晰度差,同一条分析链路的模型和后处理逻辑没有聚合在一起。
核心诉求是实现类似以下伪代码的效果,将同链路的模型训练与后处理逻辑聚合,同时满足依赖调度、优先级调度规则:
m1res <- future({ m1() }, priority = 1) p1res <- future({ process(m1res) }, priority = 2, dependson = "m1res") m2res <- future({ m2() }, priority = 1) p2res <- future({ process(m2res) }, priority = 2, dependson = "m2res")
调度规则要求:
- 只要任务的依赖已完成,空闲资源就会贪婪执行就绪任务
- 有多个就绪任务时,高优先级任务优先执行
实现方案
方案1:搭配promises包实现显式依赖绑定
原生future生态下结合promises包即可实现需求,代码结构和你期望的伪代码高度一致:
library(future) library(promises) plan(multisession, workers = 2) start <- proc.time() # 链路1:模型1+后处理,逻辑聚合编写 m1res <- future(m1()) p1res <- m1res %...>% process() # 链路2:模型2+后处理,逻辑聚合编写 m2res <- future(m2()) p2res <- m2res %...>% process() # 链路3:模型3+后处理,逻辑聚合编写 m3res <- future(m3()) p3res <- m3res %...>% process() # 输出结果 cat(value(p1res), fill = TRUE) cat(value(p2res), fill = TRUE) cat(value(p3res), fill = TRUE) stop <- proc.time() stop - start
%...>%是promises的管道运算符,会自动等待上游future完成后再执行下游逻辑,同时不会提前阻塞主线程,调度逻辑和最优方案完全一致,耗时稳定在82秒左右,同链路逻辑完全聚合,可读性符合要求。
方案2:使用targets包实现生产级工作流调度
如果流程后续会持续迭代扩展,更推荐使用专门的R工作流调度包targets,原生支持依赖检测、优先级调度、缓存复用,代码结构更清晰:
library(targets) # 配置并行参数 tar_option_set( controller = crew::crew_controller_local(workers = 2), garbage_collection = TRUE ) # 定义工作流,同链路逻辑聚合 list( # 链路1 tar_target(m1_res, m1(), priority = 1), tar_target(p1_res, process(m1_res), priority = 2), # 链路2 tar_target(m2_res, m2(), priority = 1), tar_target(p2_res, process(m2_res), priority = 2), # 链路3 tar_target(m3_res, m3(), priority = 1), tar_target(p3_res, process(m3_res), priority = 2) )
执行tar_make()即可自动按照最优调度逻辑运行所有任务,支持断点续跑,修改某个步骤后不需要重新运行全流程。
内容的提问来源于stack exchange,提问作者Joshua
相关产品推荐
相关产品推荐

