SAND Join — A skew handling join algorithm for Google's MapReduce framework

SAND Join — A skew handling join algorithm for Google's MapReduce framework
复制标题

DOI:
10.1109/inmic.2011.6151466
复制
发表时间:
2011-12
期刊:
2011 IEEE 14th International Multitopic Conference
影响因子:
--
通讯作者:
F. Atta;Stratis Viglas;Salman Niazi
F. Atta;Stratis Viglas;Salman Niazi
中科院分区:
其他
文献类型:
--
作者:
F. Atta;Stratis Viglas;Salman Niazi

文献摘要

被引文献

相似文献

MapReduce框架的简单性和灵活性促使大型分布式数据处理应用程序的程序员使用该框架开发他们的应用程序。然而,该框架的实现,包括Hadoop,并不能有效地处理输入数据中的偏差。输入数据的偏差会导致较差的负载平衡,这可能会淹没在这种并行处理框架上并行应用程序所能获得的好处。连接操作是代价最高、执行频率最高的操作,在待连接的输入数据集中存在严重的偏差时,连接操作的性能会严重下降。Hadoop的联接操作实现不能有效地处理这种不对称联接,这归因于使用散列分区进行负载分配。在这项工作中,我们引入了“倾斜处理连接”(SAND JOIN),它使用范围划分而不是哈希划分来进行负载分配。实验表明,SAND JOIN算法能够有效地对倾斜程度较大的数据集进行连接。并将该算法与Hadoop的连接算法进行了性能比较。
The simplicity and flexibility of the MapReduce framework have motivated programmers of large scale distributed data processing applications to develop their applications using this framework. However, the implementations of this framework, including Hadoop, do not handle skew in the input data effectively. Skew in the input data results in poor load balancing which can swamp the benefits achievable by parallelization of applications on such parallel processing frameworks. The performance of join operation, which is the most expensive and most frequently executed operation, is severely degraded in the presence of heavy skew in the input datasets to be joined. Hadoop's implementation of the join operation cannot effectively handle such skewed joins, attributed to the use of hash partitioning for load distribution. In this work, we introduce “Skew hANDling Join” (SAND Join) that employs range partitioning instead of hash partitioning for load distribution. Experiments show that SAND Join algorithm can efficiently perform joins on the datasets that are sufficiently skewed. We also compare the performance of this algorithm with that of Hadoop's join algorithms.