R中使用future_map处理大对象时,如何优化并行调度效率?
当并行任务依赖超大对象时,你观察到的手动分块比直接拆分大量小任务更快的现象,核心原因是减少了大对象在进程间的重复序列化/反序列化传输次数:100个独立小任务会触发100次大对象传输,而分拆为10个块后仅需传输10次,自然能节省时间。针对你的需求,以下几种方法可以让内置调度器更高效,无需手动嵌套分块:
1. 显式预加载全局共享对象
future框架支持提前将大对象一次性传递给所有worker进程,避免每个任务重复传输。你可以通过globals参数指定要共享的对象,让调度器在启动worker时完成一次传输:
library(future) library(furrr) library(microbenchmark) big_file <- matrix(1:1e8, nrow=1e3) # 显式指定共享大对象,提前传递给所有worker plan("multisession", workers = 10, globals = list(big_file = big_file)) microbenchmark::microbenchmark( "preloaded_globals" = future_map(1:100, function(a){dim(big_file)}), "manual_chunks" = future_map(1:10, function(a){map(1:10, function(b){dim(big_file)})}), times = 10 )
这种方式下,preloaded_globals的耗时会接近手动分块的结果,因为大对象仅在worker启动时传输一次。
2. 切换到FORK模式(Linux/macOS专属)
你提到最终会用multicore(FORK)模式,这是解决大对象传输问题的最优方案:FORK模式下子进程会直接继承父进程的内存空间,大对象无需序列化传输,完全共享内存。此时分块和不分块的耗时差异会几乎消失:
big_file <- matrix(1:1e8, nrow=1e3) plan("multicore", workers = 10) # Linux/macOS支持FORK microbenchmark::microbenchmark( "fork_one_run" = future_map(1:100, function(a){dim(big_file)}), "fork_chunks" = future_map(1:10, function(a){map(1:10, function(b){dim(big_file)})}), times = 10 )
测试结果会显示两种方式耗时几乎一致,且整体耗时远低于multisession模式。
3. 调整调度器的批量任务分配参数
furrr的future_map支持通过.options参数设置scheduling值,控制每个worker一次性接收的任务数量,实现自动分块:
big_file <- matrix(1:1e8, nrow=1e3) plan("multisession", workers = 10) microbenchmark::microbenchmark( "auto_chunk" = future_map(1:100, function(a){dim(big_file)}, .options = furrr_options(scheduling = 10)), # 每个worker处理10个任务 "manual_chunks" = future_map(1:10, function(a){map(1:10, function(b){dim(big_file)})}), times = 10 )
scheduling值越大,每个worker拿到的任务块越大,传输大对象的次数越少,效果和手动分块一致,无需手动嵌套map。
4. 用共享内存存储大对象
如果无法使用FORK模式(比如Windows系统),可以用bigmemory包将大对象存入共享内存,所有worker直接读取,彻底避免传输:
library(future) library(furrr) library(microbenchmark) library(bigmemory) # 创建共享内存矩阵 big_shared <- filebacked.big.matrix(nrow = 1e3, ncol = 1e5, type = "integer") big_shared[] <- 1:1e8 plan("multisession", workers = 10) microbenchmark::microbenchmark( "shared_memory" = future_map(1:100, function(a){dim(big_shared)}), times = 10 )
这种方式下,无论多少任务,大对象都无需在进程间传输,耗时会显著降低。
内容的提问来源于stack exchange,提问作者m.evans

