parLapply у реальних задачах

Паралельне програмування в R

Nabeel Imam

Data Scientist

Познайомимось із воркерами

cluster <- makeCluster(4)
clusterEvalQ(cluster, {
  id <- Sys.getpid()
  print(
    paste("Hello, my worker ID is", id)
  )
})
[[1]]
[1] "Hello, my worker ID is 425108"

[[2]]
[1] "Hello, my worker ID is 425129"

[[3]]
[1] "Hello, my worker ID is 425150"

[[4]]
[1] "Hello, my worker ID is 425171"
Паралельне програмування в R

Фільтрація даних паралельно

print(file_list)
 [1] "./health/Afghanistan.csv"            
 [2] "./health/Albania.csv"            
 [3] "./health/Algeria.csv"            
 [4] "./health/American Samoa.csv"     
 [5] "./health/Andorra.csv"            
...

Стетоскоп лежить на стосі стодоларових купюр.

Паралельне програмування в R

Фільтрація даних паралельно

filterCSV <- function (csv) {
  read.csv(csv) %>% 
    dplyr::filter(!is.na(health_exp_pc))
}


cl <- makeCluster(4) ls_df <- parLapply(cl, file_list, filterCSV) stopCluster(cl)
Error in checkForRemoteErrors(val) :
  first error: could not find function "%>%"
Паралельне програмування в R

clusterEvalQ виручає

Завантажте пакет на кластер

cl <- makeCluster(4)
clusterEvalQ(cl, library(dplyr))


ls_df <- parLapply(cl, file_list, filterCSV) stopCluster(cl)

Завантажте кілька пакетів на кластер

clusterEvalQ(cl, {
  library(dplyr)
  library(stringr)
})
[[1]]
       Country health_exp_pc Year
1  Afghanistan      81.27103 2002
2  Afghanistan      82.45785 2003
3  Afghanistan      89.47005 2004
...

[[2]]
   Country health_exp_pc Year
1  Albania      300.2757 2001
2  Albania      314.3254 2002
3  Albania      343.9442 2003
...
Паралельне програмування в R

Фільтрація за умовами

# Function with an argument for starting year
filterCSV <- function (csv, min_year) { 

  read.csv(csv) %>% 
    dplyr::filter(!is.na(health_exp_pc),

                  # Filter data for min_year and onwards
                  Year >= min_year) 
}


selected_year <- 2010 # Value to be supplied to min_year
Паралельне програмування в R

Фільтрація за умовами

cl <- makeCluster(4)
clusterEvalQ(cl, library(dplyr))

clusterExport(cl, "selected_year",
envir = environment())
ls_df <- parLapply(cl, file_list, filterCSV,
min_year = selected_year)
stopCluster(cl)

   

  • Експортуйте selected_year до кластера
  • Експортуйте з поточного середовища

 

  • Передайте selected_year у min_year
Паралельне програмування в R

Фільтрація за умовами

[[1]]
       Country health_exp_pc Year
1  Afghanistan      143.6695 2010
2  Afghanistan      143.0915 2011
3  Afghanistan      151.9180 2012
...
  Country health_exp_pc Year
1 Albania      451.8820 2010
2 Albania      485.5835 2011
3 Albania      529.6322 2012
...
Паралельне програмування в R

Чекліст гігієни кластера

  • Визначте кількість ядер
  • Створіть відповідний кластер
    • PSOCK для сумісності на всіх системах
    • FORK для Linux або Mac (і швидкості!)
  • Завантажте потрібні бібліотеки
  • Експортуйте потрібні змінні
  • Передайте експортовану змінну в іменований аргумент
  • Зупиніть кластер, коли завершите
n_cores <- detectCores() - 2


cluster <- makeCluster(n_cores)
clusterEvalQ(cluster, library(crucial_package))
clusterExport(cluster, "variable_we_need")
parLapply(cluster, ls_inputs, our_function, named_argument = variable_we_need)
stopCluster(cluster)
Паралельне програмування в R

Давайте потренуємось!

Паралельне програмування в R

Preparing Video For Download...