Chi squared feature selection over Apache Spark
Mohamed Nassar, Haı̈dar Safa, Alaa Al Mutawa, Helal Ahmed, Iskander Gaba · 2019
We live in the age of big data and distributed computing. The current large scale computation frameworks are based on a scaling-out approach for distributing tasks over a cluster of commodity machines. Apache Spark is one of these frameworks that has excelled in many computational tasks. Implementation of statistical learning algorithms over Spark is a challenging task. A bad implementation may lead to a significant decrease in performance and a waste of cluster time and money. Poor performance is mostly due to a lack of understanding of the data in hand and Spark's underlying mechanisms more than it is due to a deficit in the framework itself. In this paper, we consider the use case of X2 feature selection which is very popular in supervised learning pipelines. Our implementation follows the algorithm of the Scikit-learn Python machine learning library which is different than the algorithm used by the Spark machine learning library. The Spark ML library implementation of X2 feature selection accepts only categorical features. Our alternative implementation is more suitable for numerical features. We experiment in particular with features of high sparsity such as n-gram counts. We study the best partitioning scheme of the data and the optimal number of partitions. Our experiments are run over the Databricks platform.