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

如何用R并行分块读取多个大CSV,充分利用多核资源提效

利用128核并行处理单文件内的分块任务

完全可以通过文件级并行 + 单文件内分块并行的方式利用100核资源,下面提供两种可行的实现方案,同时解决表头不一致的合并问题。

方案一:嵌套foreach实现双层并行

思路:外层并行处理10个文件,内层将单个文件拆分为多个块并行读取过滤,最后按文件合并块结果,再合并所有文件的结果。

关键调整点

  1. 手动拆分文件行范围(替代read_csv_chunked的串行分块),实现单文件内的并行读取
  2. 处理表头:仅第一个块读取表头,后续块跳过表头避免重复
  3. 内层用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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.29 13:08:15