Ann*_*ana 4 csv parallel-processing r parallel.foreach doparallel
我有 10 个非常大的 CSV 文件(可能有也可能没有相同的标题),我正在使用“readr”包 read_csv_chunked() 连续读取和处理这些文件。目前,我可以使用 10 个内核并行读取 10 个文件。该过程仍需要一个小时。我有128个核心。我可以将每个 CSV 分成 10 个块,以便对每个文件并行处理,从而利用 100 个核心吗?这是我目前拥有的(创建两个示例文件仅用于测试):
library(doParallel)
library(foreach)
# Create a list of two sample CSV files and a filter by file
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")
# Create a function to read and filter each file in chunks
my_function <-
function(file) {
library(dplyr)
library(plyr)
library(readr)
filter_df <-
function(x, pos) {
subset(x, X1 %in% filter_by_df$X1 | X2 %in% filter_by_df$X1)
}
readr_df <-
read_csv_chunked(file,
callback = DataFrameCallback$new(filter_df),
progress = F,
chunk_size = 50) %>%
as.data.frame() %>%
distinct()
return(readr_df)
}
# Apply the custom function created above to all files in parallel and combine them
df_foreach <-
foreach(i = files, .combine = rbind.fill, .packages = c("plyr")) %dopar% my_function(i)
Run Code Online (Sandbox Code Playgroud)
我正在研究嵌套 foreach() 但不确定如何以嵌套方式传递不同的函数(read_csv_chunked 与 .combine = rbind 以及我的自定义函数 filter_df() 与 .combine = rbinf.fill())。我还研究了“未来”包。任何意见是极大的赞赏。谢谢。
正如 Sirius 所提到的,data.table::fread()它无疑是最快的 csv 阅读器R,并且内置的多线程应该充分利用您可以使用的资源。
但还有一个想法 - 由于您只需要最终结果中的行的子集,arrow因此该包是一个不错的选择。您可以使用 arrow 的“下推”功能来扫描多文件数据集,并仅将符合您条件的行读入内存,而不是将整个文件读入内存。
arrow 部分实现了 dplyr 的功能,因此根据您的经验,语法可能已经熟悉。
library(arrow)
library(dplyr)
DS <- arrow::open_dataset(sources = c("df_1.csv","df_1.csv"),
format = "csv")
DS |>
filter(X1 %in% filter_by_df[["X1"]]| X2 %in% filter_by_df[["X1"]] ) |>
distinct() |>
collect()
Run Code Online (Sandbox Code Playgroud)
请注意 - 该答案假设 R 版本为 4.1 或更高版本。%>%对于早期的 R 版本,可以使用magrittr 管道代替|>4.1 中介绍的基本 R 管道。