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

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. 只要任务的依赖已完成,空闲资源就会贪婪执行就绪任务
  2. 有多个就绪任务时,高优先级任务优先执行

实现方案

方案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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.10.02 13:27:02