A Comparative Analysis of Join Algorithms Using the Hadoop Map/Reduce Framework

Konstantina Palla · 2009

The Map/Reduce framework is a programming model recently introduced by Google Inc. to support distributed computing on very large datasets across a large number of machines. It provides a simple but yet powerful way to implement distributed applications without having deeper knowledge of parallel programming. Each participating node executes Map and/or Reduce tasks which involve reading and writing large datasets. In this work, we exploited the open source Hadoop implementation of the Map/Reduce framework and developed a theoretical cost model to evaluate the I/O cost induced on each node during the execution of a task. Furthermore, we used our model to evaluate and compare two basic join algorithms already provided by the framework — the reduce side and the map side join — and a third one we implemented and which makes use of a Bloom filter on the map tasks. The experimental results proved the validity of our cost model and, furthermore, stressed out that our proposed algorithm claims the least I/O cost in comparison to the others. This work is expected to provide a good insight into the Map/Reduce framework in terms of I/O cost on the one hand, and a thorough analysis of three join implementations under the Map/Reduce dataflow on the other hand.

Read the paper · More papers on PaperTik