如何用R并行分块读取多个大CSV,充分利用多核资源提效
利用128核并行处理单文件内的分块任务
完全可以通过文件级并行 + 单文件内分块并行的方式利用100核资源,下面提供两种可行的实现方案,同时解决表头不一致的合并问题。
方案一:嵌套foreach实现双层并行
思路:外层并行处理10个文件,内层将单个文件拆分为多个块并行读取过滤,最后按文件合并块结果,再合并所有文件的结果。
关键调整点
- 手动拆分文件行范围(替代
read_csv_chunked的串行分块),实现单文件内的并行读取 - 处理表头:仅第一个块读取表头,后续块跳过表头避免重复
- 内层用
rbind合并同文件的块结果,外层用rbind.fill合并不同表头的文件结果
完整代码示例
library(doParallel) library(foreach) library(readr) library(plyr) library(dplyr) # 生成测试数据 df_1 <- data.frame(matrix(sample(1:300), ncol = 3)) df_2 <- data.frame(matrix(sample(1:200), ncol = 4)) filter_by_df <- data.frame(X1 = 1:100) write.csv(df_1, "df_1.csv", row.names = FALSE) write.csv(df_2, "df_2.csv", row.names = FALSE) files <- c("df_1.csv", "df_2.csv") # 工具函数:获取文件总行数(不含表头) get_file_row_count <- function(file) { length(read_lines(file, skip = 1)) } # 工具函数:读取文件指定行范围的块并过滤 read_filter_chunk <- function(file, skip_rows, n_rows, has_header) { chunk <- read_csv(file, skip = skip_rows, n_max = n_rows, col_names = has_header) filter(chunk, X1 %in% filter_by_df$X1 | X2 %in% filter_by_df$X1) %>% distinct() } # 外层并行:处理每个文件 cl <- makeCluster(10) # 外层用10核对应10个文件 registerDoParallel(cl) df_result <- foreach(file = files, .combine = rbind.fill, .packages = c("readr", "dplyr", "plyr")) %dopar% { total_rows <- get_file_row_count(file) chunk_num <- 10 # 每个文件拆成10块 chunk_size <- ceiling(total_rows / chunk_num) # 内层并行:处理当前文件的所有块 inner_cl <- makeCluster(10) # 每个文件用10核,总共10*10=100核 registerDoParallel(inner_cl) chunk_results <- foreach(i = 1:chunk_num, .combine = rbind, .packages = c("readr", "dplyr")) %dopar% { if (i == 1) { # 第一个块:读取表头,从第0行开始(不跳过) read_filter_chunk(file, skip_rows = 0, n_rows = chunk_size, has_header = TRUE) } else { # 后续块:跳过表头+前面已读的行,不读取表头 skip <- 1 + (i-1)*chunk_size read_filter_chunk(file, skip_rows = skip, n_rows = chunk_size, has_header = FALSE) } } stopCluster(inner_cl) chunk_results } stopCluster(cl)
方案二:用future框架简化并行逻辑
future的并行模型更灵活,无需手动管理嵌套集群,代码更简洁:
完整代码示例
library(future) library(future.apply) library(readr) library(plyr) library(dplyr) # 初始化并行:启用多核模式,直接用100核 plan(multisession, workers = 100) # 测试数据和工具函数同方案一 df_1 <- data.frame(matrix(sample(1:300), ncol = 3)) df_2 <- data.frame(matrix(sample(1:200), ncol = 4)) filter_by_df <- data.frame(X1 = 1:100) write.csv(df_1, "df_1.csv", row.names = FALSE) write.csv(df_2, "df_2.csv", row.names = FALSE) files <- c("df_1.csv", "df_2.csv") get_file_row_count <- function(file) { length(read_lines(file, skip = 1)) } read_filter_chunk <- function(file, skip_rows, n_rows, has_header) { chunk <- read_csv(file, skip = skip_rows, n_max = n_rows, col_names = has_header) filter(chunk, X1 %in% filter_by_df$X1 | X2 %in% filter_by_df$X1) %>% distinct() } # 处理单个文件的所有块 process_file <- function(file) { total_rows <- get_file_row_count(file) chunk_num <- 10 chunk_size <- ceiling(total_rows / chunk_num) # 生成每个块的参数 chunk_params <- lapply(1:chunk_num, function(i) { if (i == 1) { list(skip_rows = 0, n_rows = chunk_size, has_header = TRUE) } else { list(skip_rows = 1 + (i-1)*chunk_size, n_rows = chunk_size, has_header = FALSE) } }) # 并行处理当前文件的所有块 chunk_results <- future_lapply(chunk_params, function(params) { read_filter_chunk(file, params$skip_rows, params$n_rows, params$has_header) }) # 合并同文件的块结果 do.call(rbind, chunk_results) } # 并行处理所有文件,合并结果 all_results <- future_lapply(files, process_file) df_result <- do.call(rbind.fill, all_results)
注意事项
- IO瓶颈:100核同时读取文件可能触发磁盘IO瓶颈,如果是机械硬盘可能提升有限,建议用SSD存储
- 内存控制:调整
chunk_size避免单个块过大导致内存溢出,可根据单块处理后的数据集大小估算 - 全局变量传递:确保
filter_by_df在并行进程中可访问,future会自动处理全局变量的传递,foreach需注意.export参数(如果变量不在全局环境) - 包加载:尽量在并行任务外加载依赖包,减少进程初始化开销
内容的提问来源于stack exchange,提问作者Annabanana
相关产品推荐
相关产品推荐

