将R语言for循环正确转换为并行循环的技术问询
问题背景
现有一份考试结果数据集(学生多年多次参加考试,结果为通过/不通过,需研究前一次考试对后一次的影响),生成代码如下:
id = sample.int(10000, 100000, replace = TRUE) res = c(1,0) results = sample(res, 100000, replace = TRUE) date_exam_taken = sample(seq(as.Date('1999/01/01'), as.Date('2020/01/01'), by="day"), 100000, replace = TRUE) my_data = data.frame(id, results, date_exam_taken) my_data <- my_data[order(my_data$id, my_data$date_exam_taken),] my_data$general_id = 1:nrow(my_data) my_data$exam_number = ave(my_data$general_id, my_data$id, FUN = seq_along) my_data$general_id = NULL # 数据集示例 # id results date_exam_taken exam_number # 7992 1 1 2004-04-23 1 # 24837 1 0 2004-12-10 2 # 12331 1 1 2007-01-19 3 # 34396 1 0 2007-02-21 4 # 85250 1 0 2007-09-26 5 # 11254 1 1 2009-12-20 6
已编写的标准for循环运行正常:
my_list = list() for (i in 1:length(unique(my_data$id))) { tryCatch({ start_i = my_data[my_data$id == i,] pairs_i = data.frame(first = head(start_i$results, -1), second = tail(start_i$results, -1)) frame_i = as.data.frame(table(pairs_i)) frame_i$id = i print(frame_i) my_list[[i]] = frame_i }, error = function(e){}) } final_a = do.call(rbind.data.frame, my_list)
尝试用doParallel库优化循环,转换后的代码如下:
# detectCores()返回8,是否应该设置makeCluster(8)? > detectCores() [1] 8 my_list = list() max = length(unique(my_data$id)) library(doParallel) registerDoParallel(cl <- makeCluster(3)) # 注:该循环没有打印输出? test = foreach(i = 1:max, .combine = "rbind") %dopar% { tryCatch({ start_i = my_data[my_data$id == i,] pairs_i = data.frame(first = head(start_i$results, -1), second = tail(start_i$results, -1)) frame_i = as.data.frame(table(pairs_i)) frame_i$id = i print(frame_i) my_list[[i]] = frame_i }, error = function(e){}) } final_b = do.call(rbind.data.frame, test)
问题
- 我是否正确使用了doParallel的功能?
- 是否有更优的实现方式?
- 注:该代码将用于处理约1000万唯一ID的数据集。
解答
1. doParallel使用是否正确?
你的代码没有正确使用doParallel,存在以下核心问题:
- 全局变量复制风险:并行循环中直接引用
my_data,每个子进程都会复制整个数据集,针对1000万ID的超大数据集,会导致内存占用暴增,甚至集群崩溃。 - 无效的列表操作:
my_list[[i]] = frame_i在并行环境中完全无效,每个子进程有独立内存空间,无法修改主环境的列表。 - 无有效返回值:循环内没有明确返回
frame_i,最终test会是大量NULL,导致final_b生成失败。 - 打印无输出:并行子进程的标准输出默认不会传递到主进程,
print(frame_i)无法在控制台显示。 - 集群设置不合理:
detectCores()返回8时,建议设置makeCluster(detectCores() - 1)(留1个核心给系统),而非固定3个,以最大化CPU利用率。
修正后的正确并行写法:
library(doParallel) # 预留1个核心给系统进程 cl <- makeCluster(detectCores() - 1) registerDoParallel(cl) # 提前提取唯一ID,避免循环内重复计算 unique_ids <- unique(my_data$id) max_ids <- length(unique_ids) test <- foreach(i = 1:max_ids, .combine = rbind, .packages = "base") %dopar% { current_id <- unique_ids[i] start_i <- my_data[my_data$id == current_id, ] # 跳过只有1次考试的ID,避免head/tail报错 if(nrow(start_i) < 2) return(NULL) pairs_i <- data.frame(first = head(start_i$results, -1), second = tail(start_i$results, -1)) frame_i <- as.data.frame(table(pairs_i)) frame_i$id <- current_id # 明确返回结果 frame_i } # 关闭集群释放资源 stopCluster(cl) final_b <- test
2. 更优的实现方式
针对1000万唯一ID的数据集,并行循环仍不是最优解,推荐使用数据.table或dplyr的分组操作,底层基于C++实现,效率比R循环(包括并行)高几个数量级:
方式1:data.table实现(性能最优)
library(data.table) setDT(my_data) # 按ID分组,生成前一次考试结果 my_data[, prev_result := shift(results, n = 1, type = "lag"), by = id] # 过滤无前置结果的行,统计分组频次 final_dt <- my_data[!is.na(prev_result), .(Freq = .N), by = .(id, first = prev_result, second = results)]
方式2:dplyr实现(语法更简洁)
library(dplyr) final_dplyr <- my_data %>% group_by(id) %>% mutate(prev_result = lag(results)) %>% filter(!is.na(prev_result)) %>% group_by(id, prev_result, results) %>% summarise(Freq = n(), .groups = "drop") %>% rename(first = prev_result, second = results)
这两种方法无需手动循环,避免了数据重复复制,内存占用更可控,处理超大数据集的速度远快于并行循环。
3. 针对1000万ID数据集的额外建议
- 预处理过滤:提前过滤掉只有1次考试的ID,这类数据对研究“前一次考试的影响”无意义,可大幅减少计算量。
- 内存优化:用
data.table::fread()读取数据,比read.csv效率更高;将id、results等字段设为整数类型,降低内存占用。 - 避免全局复制:若必须用并行,可先按ID拆分数据集为列表,再分配给子进程,但分组操作已完全规避此问题。
内容的提问来源于stack exchange,提问作者stats_noob
相关产品推荐
相关产品推荐

