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.