Applying user defined function to normalize all columns in sparklyr using spark_apply

Viewed 90

I have a spark dataframe that i manipulate using sparklyr that has > 100 columns. I would like to normalize each column in the following way (vector - mean(vector)) / sd(vector). To achieve that in R, I could use dplyr in the following way:

library(dplyr)

normalize <- function(vector){
  vector_norm = (vector - mean(vector)) / sd(vector)
  return(vector_norm)
}

iris %>%
  select(-Species) %>%
  mutate_all(funs(normalize(.))) %>% 
  view

Unfortunately, sparklyr is incapable of running user defined functions in R natively. There is an approach using spark_apply that allows this to be run (though inefficiently). My best attempt at that approach is the following:

# Connect to Spark and push iris dataset to Spark
library(sparklyr)
sc <- spark_connect(method = "databricks")
iris_sdf <- sdf_copy_to(sc, iris %>% head(4), overwrite = T)
 
schema <- as.list(colnames(iris))
 
results_sdf <- spark_apply(iris_sdf,
                           function(vector){
                                vector_norm = (vector - mean(vector)) / sd(vector)
                                    return(vector_norm)
                                  }, 
                           columns = schema)
head(results_sdf, 10)

But i got the following error:

Error : org.apache.spark.SparkException: Job aborted due to stage failure: Task 0 in stage 21.0 failed 4 times, most recent failure: Lost task 0.3 in stage 21.0 (TID 25, 10.19.216.60, executor 0): java.lang.Exception: sparklyr worker rscript failure with status 255, check worker logs for details. Error : org.apache.spark.SparkException: Job aborted due to stage failure: Task 0 in stage 21.0 failed 4 times, most recent failure: Lost task 0.3 in stage 21.0 (TID 25, 10.19.216.60, executor 0): java.lang.Exception: sparklyr worker rscript failure with status 255, check worker logs for details.
    at sparklyr.Rscript.init(rscript.scala:83)
    at sparklyr.WorkerApply$$anon$2.run(workerapply.scala:133)

I also tried:

iris_sdf %>%
  spark_apply(
    function(e) data.frame((e$Sepal.Length - mean(e$Sepal.Length)) / sd(e$Sepal.Length)),
    names = c("Sepal.Length")
  ) 

No error but the resulting output had zero rows.

I would be open to any solution in sparklyr, pyspark, or scala.

0 Answers
Related