【发布时间】:2020-11-18 06:29:48
【问题描述】:
我有一个 4500 行的输入 csv 文件。每行都有一个唯一的 ID,对于每一行,我必须读取一些数据,进行一些计算,然后将输出写入 csv 文件,这样我的输出目录中就有 4500 个 csv 文件。单个输出 csv 文件包含 8 列的单行数据
由于我必须对输入 csv 的每一行执行相同的计算,我想我可以使用 foreach 并行化这个任务。以下是整体逻辑结构
library(doSNOW)
library(foreach)
library(data.table)
input_csv <- fread('inputFile.csv'))
# to track the progres of the loop
iterations <- nrow(input_csv)
pb <- txtProgressBar(max = iterations, style = 3)
progress <- function(n) setTxtProgressBar(pb, n)
opts <- list(progress = progress)
myClusters <- makeCluster(6)
registerDoSNOW(myClusters)
results <-
foreach(i = 1:nrow(input_csv),
.packages = c("myCustomPkg","dplyr","arrow","zoo","data.table","rlist","stringr"),
.errorhandling = 'remove',
.options.snow = opts) %dopar%
{
rowRef <- input_csv[i, ]
# read data for the unique location in `rowRef`
weather.path <- arrow(paste0(rowRef$locationID'_weather.parquet')))
# do some calculations
# save the results as csv
fwrite(temp_result, file.path(paste0('output_iter_',i,'.csv')))
return(temp_result)
}
上述代码运行良好,但在完成input_csv 中 25% 或 30% 的行后总是卡住/不活动/不执行任何操作。我一直在查看我的输出目录,在 N% 的迭代之后,没有文件被写入。我怀疑 foreach 循环是否进入某种睡眠模式?我发现更令人困惑的是,如果我终止工作,重新运行上述代码,它确实会显示 16% 或 30%,然后再次进入非活动状态,即每次重新运行时,它都会以不同的进度级别“休眠”。
在这种情况下,我不知道如何给出一个最小的可重现示例,但我想如果有人知道我应该检查的任何清单或导致这种情况的潜在问题,那将非常有帮助。谢谢
编辑我仍在努力解决这个问题。如果我可以提供更多信息,请告诉我。
EDIT2
我原来的 inputFile 包含 213164 行。所以我拆分了我的大文件
分成 46 个小文件,每个文件有 4634 行
library(foreach)
library(data.table)
library(doParallel)
myLs <- split(mydat, (as.numeric(rownames(mydat))-1) %/% 46))
然后我这样做了:
for(pr in 1:46){
input_csv <- myLs[[pr]]
myClusters <- parallel::makeCluster(6)
doParallel::registerDoParallel(myClusters)
results <-
foreach(i = 1:nrow(input_csv),
.packages = c("myCustomPkg","dplyr","arrow","zoo","data.table","rlist","stringr"),
.errorhandling = 'remove',
.verbose = TRUE) %dopar%
{
rowRef <- input_csv[i, ]
# read data for the unique location in `rowRef`
weather.path <- arrow(paste0(rowRef$locationID'_weather.parquet')))
# do some calculations
# save the results as csv
fwrite(temp_result, file.path(paste0('output_iter_',i,'_',pr,'.csv')))
gc()
}
parallel::stopCluster(myClusters)
gc()
}
这也有效,直到说 pr = 7 或 pr = 8 迭代然后不继续 也不会生成任何错误消息。我很困惑。
编辑 这就是我的 CPU 使用率。我只使用了 4 个核心来生成这个图像。谁能解释这张图片中是否有任何东西可以解决我的问题。
【问题讨论】:
-
好像你正在返回
temp_result。是内存问题吗? -
是的,我正在返回 temp_result。有什么方法可以检查它是否确实是由内存问题引起的,因为没有产生错误。脚本仅在 25% 或 30% 或 10% 处停止并且不会移动。如果我终止作业,仍然不会产生错误。
-
你应该打开某种系统监视器。
-
几个月前,有人在导出大量文件时遇到问题,他们也使用了
fwrite(),但看起来他们删除了这个问题。如果我没记错的话,例如 50 个文件的速度更快,但例如 500 个文件的速度较慢。我无法记住差异的大小。综上所述,可能值得尝试将fwrite()换成readr::write_csv()。另一种可能性是,考虑到将文件全部保存到results,您可以尝试在另一个步骤中写入文件 -
好的。感谢您的评论。我将阅读 readr 函数并检查它是否有帮助
标签: r foreach doparallel