【问题标题】:R parallel - functions gives only intialmatrix backR 并行 - 函数只返回初始矩阵
【发布时间】:2020-10-20 19:39:20
【问题描述】:

我尝试从 data.frame 中逐行检查一行中的元素是否相同。

我的真实数据集包含超过 100 万行和几列负数和正数,以及 0 和 NA。总共有 15 个数据集我要检查,因此是并行变体。

不幸的是,我当前的代码只给了我初始矩阵。由于我从未使用过“并行”包,所以我不太了解它。我尝试使用“clusterExport”参数,但到目前为止没有任何帮助。

因此我的错误在哪里的问题以及寻求帮助的请求。

非常感谢。

x_x <- data.frame("x"=rep(c(1,2,3,4,5,6,7,8,9,10),10),"y"=rep(c(1,2,3,2,1,NA,7,8,9,10),10),"z"=rep(c(1,2,3,4,5,6,7,8,9,10),10))

library(foreach)
library(doParallel)
no_cores <- detectCores() - 1
cl <- makeCluster(no_cores) 
registerDoParallel(cl)

Test_parallel2 <- function(data_01)
{
  # data_01 <- x_x
  return_data <- data.frame("V1"=matrix(FALSE,nrow = nrow(data_01),ncol = 1))
  data_01 <- as.data.frame(t(data_01))

  is_true_eigen <- function(data_vec)
  {
    # data_vec <- c(TRUE,FALSE,TRUE)
    return_data_is_true <- TRUE
    for(i in 1:length(data_vec))
    {
      if(data_vec[i] == FALSE)
      {
        return_data_is_true <- FALSE
        break(i)
      }
    }
    return(return_data_is_true)
  }
  
  y <- foreach(i=1:ncol(data_01)) %dopar% {
    
    if( (data_01[1,i] == data_01[2,i]) == TRUE & ( mean(as.numeric(data_01[,i]),na.rm = TRUE) == data_01[1,i]) & ( is_true_eigen(colSums(!is.na(t(data_01[,i]))) > 0) == TRUE ) ) 
    {
      return_data[i,1] <- TRUE  # i=1
    }
    
  }
  
  #parallel::clusterExport(cl = cl,varlist = c("y"),envir=environment())
  
  return(return_data)
  
}

Test_parallel2(x_x)

编辑:

输出应该是每行的向量,有真或假(如果行元素相同)

示例:

第 1 行(来自 x_x):

1 | 1 | 1

应该返回一个 TRUE

第 4 行(来自 x_x):

4 | 2 | 4

应该返回一个 FALSE

第 6 行(来自 x_x):

6 |不适用 | 6

应该返回一个 FALSE

【问题讨论】:

  • 我不太确定你想做什么。您能否给我们一个示例,说明您期望从 x_x 样本数据输入中获得的输出?
  • 感谢您的建议。我已经添加了它。希望现在清楚了。

标签: r parallel-processing


【解决方案1】:

由于没有其他人插话,所以我尝试了这个。我在日常工作或业余爱好项目中不使用 foreachdoParallel 包,因此我恢复了并行化的首选方法,这是 parallel 包。

我使用三种并行化方案完成了这项工作:串行(无并行化)、fork 线程(仅限类 Unix 系统)和 PSOCK(Windows 和类 Unix 系统)。

要在没有任何并行化的情况下做到这一点,我们可以定义以下函数:

#####################
### Serial Method ###
#####################

### row_equal() ###
row_equal <- function(data){
  
  ## Transpose Data ##
  data <- as.data.frame(t(data))
  
  ## Apply Function Internals ##
  n <- lapply(
    X = data,
    FUN = unique
  )
  
  ## Number of Unique Values ##
  n <- lengths(n, use.names = FALSE)
  
  ## Convert to Output to Logical ##
  eq <- ifelse(n == 1L, TRUE, FALSE)
  
  ## Output ##
  return(eq)
  
}

要使用 fork 线程并行化 row_equal() 函数,我们可以将其修改为以下内容:

##############################
### Parallel (FORK) Method ###
##############################

### row_equal_fork() ###
row_equal_fork <- function(data){
  
  ## Transpose Data ##
  data <- as.data.frame(t(data))
  
  ## Cores ##
  n_cores <- max(parallel::detectCores() - 1L, 1L)
  
  ## Apply Function Internals ##
  n <- parallel::mclapply(
    X = data,
    FUN = unique,
    mc.cores = n_cores
  )
  
  ## Number of Unique Values ##
  n <- lengths(n, use.names = FALSE)
  
  ## Convert to Output to Logical ##
  eq <- ifelse(n == 1L, TRUE, FALSE)
  
  ## Output ##
  return(eq)
  
}

不幸的是,简单的 fork-threading 版本仅适用于类 Unix 系统。 Windows 做不到。对于 Windows,我们需要设置一个 PSOCK 集群,将我们希望它执行的工作传递给它,然后在完成/失败时停止集群。在这种情况下,作业非常简单,但对于更复杂的作业,您可能需要使用parallel::clusterEvalQ() 传递集群所需的包,或使用parallel::clusterExport() 传递集群附加对象。

要使用 PSOCK 集群并行化 row_equal() 函数,我们可以将其修改为以下内容:

###############################
### Parallel (PSOCK) Method ###
###############################

### row_equal_psock() ###
row_equal_psock <- function(data){
  
  ## Transpose Data ##
  data <- as.data.frame(t(data))
  
  ## Cores ##
  n_cores <- max(parallel::detectCores() - 1L, 1L)
  cl <- parallel::makeCluster(n_cores, type = "PSOCK")
  on.exit(parallel::stopCluster(cl))
  
  ## Apply Function Internals ##
  n <- parallel::parLapply(
    cl = cl,
    X = data,
    fun = unique
  )
  
  ## Number of Unique Values ##
  n <- lengths(n, use.names = FALSE)
  
  ## Convert to Output to Logical ##
  eq <- ifelse(n == 1L, TRUE, FALSE)
  
  ## Output ##
  return(eq)
  
}

对于您的测试数据框 (x_x),我在使用函数时得到以下输出:

row_equal(x_x)
##  [1]  TRUE  TRUE  TRUE FALSE FALSE FALSE  TRUE  TRUE  TRUE  TRUE  TRUE  TRUE
## [13]  TRUE FALSE FALSE FALSE  TRUE  TRUE  TRUE  TRUE  TRUE  TRUE  TRUE FALSE
## [25] FALSE FALSE  TRUE  TRUE  TRUE  TRUE  TRUE  TRUE  TRUE FALSE FALSE FALSE
## [37]  TRUE  TRUE  TRUE  TRUE  TRUE  TRUE  TRUE FALSE FALSE FALSE  TRUE  TRUE
## [49]  TRUE  TRUE  TRUE  TRUE  TRUE FALSE FALSE FALSE  TRUE  TRUE  TRUE  TRUE
## [61]  TRUE  TRUE  TRUE FALSE FALSE FALSE  TRUE  TRUE  TRUE  TRUE  TRUE  TRUE
## [73]  TRUE FALSE FALSE FALSE  TRUE  TRUE  TRUE  TRUE  TRUE  TRUE  TRUE FALSE
## [85] FALSE FALSE  TRUE  TRUE  TRUE  TRUE  TRUE  TRUE  TRUE FALSE FALSE FALSE
## [97]  TRUE  TRUE  TRUE  TRUE

所有版本的函数都给出相同的输出:

identical(row_equal(x_x), row_equal_fork(x_x))
## [1] TRUE
identical(row_equal(x_x), row_equal_psock(x_x))
## [1] TRUE

但是,请注意,并行化函数不一定会使其运行得更快,因为两种并行化方法都有相关的开销;特别是使用 PSOCK 方法(无论如何,在我蹩脚的 ARM64 笔记本电脑上):

library(microbenchmark)
microbenchmark(
  serial = row_equal(x_x),
  parallel_fork = row_equal_fork(x_x),
  parallel_psock = row_equal_psock(x_x),
  times = 100L
)
## Unit: milliseconds
##            expr        min         lq        mean      median          uq        max neval
##          serial   2.425209   5.059251    6.211456    5.545606    6.358625   21.22458   100
##   parallel_fork  76.978126  92.248626  113.667318  117.062001  127.628751  166.58804   100
##  parallel_psock 949.944959 990.536689 1014.325810 1009.042293 1036.322251 1120.35083   100

如果您的数据集包含许多行,那么您可能会开始从并行方法中看到一些好处。我内心的纯粹主义者觉得必须有一种矢量化的方式来做这件事。

【讨论】:

  • 非常感谢您的努力。我已经立即尝试了,您的串行方法效果很好。随着时间的推移,这 3 种变体的表现令人印象深刻。老实说,我没想到会这样。谢谢和问候
猜你喜欢
  • 2014-09-18
  • 1970-01-01
  • 2021-05-11
  • 2017-12-14
  • 2018-12-04
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 2012-12-26
相关资源
最近更新 更多