R语言线性回归并行化:解决美国新冠县域数据furrr提速差问题
可行优化方案如下:
1. 优化并行调用逻辑,减少调度开销
你当前代码分4次调用future_map(),每次都要完成全量数据的进程分发、结果回收的完整流程,额外开销极高,且前一次调用未完成时后一次不会启动,也会造成CPU资源闲置。之前单核满载(CPU使用率26%对应4核设备单核心跑满)的情况,大多是因为多次future_map调用的调度开销远大于实际拟合开销,导致CPU大部分时间在等待调度,没有实际执行计算任务。
你可以将4个模型的拟合逻辑合并到同一个自定义函数中,单次调用future_map()即可完成所有模型的拟合,仅产生1次调度开销,示例代码如下:
# 合并四个模型拟合到同一个函数 fit_all_models <- function(tbl) { list( deathmodel = deathsmodel(tbl), casemodel = casesmodel(tbl), absdeathmodel = absdeathsmodel(tbl), abscasemodel = abscasesmodel(tbl) ) } # 仅调用一次future_map uscases_twoweeks <- casesdeaths %>% filter(date >= twoweeksago, !is.na(population), population > min_country_population) %>% mutate(countyid = paste(county, state, sep = ", ")) %>% arrange(countyid, date) %>% group_by(countyid) %>% nest() %>% mutate(models = future_map(data, fit_all_models)) %>% unnest_wider(models) # 将列表拆分为独立的模型列
2. 调整并行策略参数
- 如果你使用的是Linux/macOS系统,优先选
multicore模式,相比multisession不需要重复复制进程数据,开销低很多。如果是Windows系统只能用multisession,可以先调大全局变量传输上限,避免数据传输被限制:
# 调整并行参数示例 options(future.globals.maxSize = 2 * 1024^3) # 设为2G,可根据你的数据量调整 future::plan(multicore, workers = future::availableCores() - 1) # 留1核给系统调度
- 要是每个县的数据集非常小(比如只有十几行),可以设置
future_map的.options = furrr_options(scheduling = 10),将任务拆分为更小的块分配给进程,减少资源浪费。
3. 替换低效率的拟合函数
base R的lm()本身会输出大量冗余的诊断信息,针对小样本的分组拟合场景额外开销占比极高。你可以换成RcppEigen::fastLm(),拟合速度是base lm()的3-5倍,输出结果和lm()基本兼容:
library(RcppEigen) # 替换原有模型函数的拟合逻辑即可 casesmodel <- function(tbl) { fastLm(casesper100k ~ time, data = tbl) }
如果只需要提取回归系数、p值等核心指标,也可以直接用矩阵运算手动计算,速度还能再提升一倍以上。
4. 替换数据处理框架进一步提速
如果你的原始数据量超过100万行,tidyverse的nest()分组操作本身就会占用不少时间,可以换成data.table做预处理+parallel包做并行拟合,整体速度能提升2-10倍:
library(data.table) library(parallel) # 预处理改用data.table setDT(casesdeaths) prep_dt <- casesdeaths[date >= twoweeksago & !is.na(population) & population > min_country_population, ] prep_dt[, countyid := paste(county, state, sep = ", ")] setorder(prep_dt, countyid, date) # 按县拆分数据 split_dt <- split(prep_dt, by = "countyid") # 初始化并行集群 cl <- makeCluster(future::availableCores() - 1) # 导出需要的函数到子进程 clusterExport(cl, c("fit_all_models", "fastLm", "deathsmodel", "casesmodel", "absdeathsmodel", "abscasesmodel")) # 并行拟合 result_list <- parLapply(cl, split_dt, fit_all_models) stopCluster(cl) # 结果合并可按需自行处理
内容的提问来源于stack exchange,提问作者Rob Hanssen
相关产品推荐
相关产品推荐

