On exploring efficient shuffle design for in-memory MapReduce

On exploring efficient shuffle design for in-memory MapReduce
复制标题

DOI:
10.1145/2926534.2926538
复制
发表时间:
2016-06
期刊:
Proceedings of the 3rd ACM SIGMOD Workshop on Algorithms and Systems for MapReduce and Beyond
影响因子:
--
通讯作者:
Harunobu Daikoku;H. Kawashima;O. Tatebe
Harunobu Daikoku;H. Kawashima;O. Tatebe
中科院分区:
其他
文献类型:
--
作者:
Harunobu Daikoku;H. Kawashima;O. Tatebe

文献摘要

相似文献

MapReduce作为一种大数据分析的方式,在很多领域都有广泛的应用。洗牌,MapReduce的节点间数据交换阶段,已被报告为该框架的主要瓶颈。文献中对洗牌的加速进行了研究,本文提出了两个问题。第一个问题涉及远程直接内存访问(RDMA)对洗牌性能的影响。RDMA使一台机器能够在另一台机器的本地存储器上读取和写入数据,并且已知是一种有效的数据传输机制。单纯使用RDMA会影响洗牌的性能吗?第二个问题是要使用的数据传输算法。传统的MapReduce实现有两种类型的洗牌算法:全连接和更复杂的算法,如Pairwise。数据传输算法是否影响洗牌的性能?为了回答这些问题,我们设计并实现了另一个MapReduce系统从头开始在C/C++获得最大的性能和保留设计的灵活性。对于第一个问题,我们比较了基于rsocket的RDMA洗牌和基于IPoIB的RDMA洗牌。GroupBy的实验结果表明,RDMA将map+shuffle阶段加速了约50%。对于第二个问题,我们首先将我们的内存系统与Apache Spark进行了比较,以调查我们的系统是否比现有系统更有效。与Spark相比,我们的系统在Word Count上的性能提高了3.04倍,在BiGram Count上提高了2.64倍。然后,我们比较了两种数据交换算法,全连接和成对。使用BiGram Count的实验结果表明,不使用RDMA的全连接比使用RDMA的Pairwise效率高13%。我们的结论是,它是必要的重叠地图和洗牌阶段,以获得性能的改善。改进百分比相对较小的原因可以归因于在映射阶段将键值对插入到哈希映射中的耗时。
MapReduce is commonly used as a way of big data analysis in many fields. Shuffling, the inter-node data exchange phase of MapReduce, has been reported as the major bottleneck of the framework. Acceleration of shuffling has been studied in literature, and we raise two questions in this paper. The first question pertains to the effect of Remote Direct Memory Access (RDMA) on the performance of shuffling. RDMA enables one machine to read and write data on the local memory of another and has been known to be an efficient data transfer mechanism. Does the pure use of RDMA affect the performance of shuffling? The second question is the data transfer algorithm to use. There are two types of shuffling algorithms for the conventional MapReduce implementations: Fully-Connected and more sophisticated algorithms such as Pairwise. Does the data transfer algorithm affect the performance of shuffling? To answer these questions, we designed and implemented yet another MapReduce system from scratch in C/C++ to gain the maximum performance and to reserve design flexibility. For the first question, we compared RDMA shuffling based on rsocket with the one based on IPoIB. The results of experiments with GroupBy showed that RDMA accelerates map+shuffle phase by around 50%. For the second question, we first compared our in-memory system with Apache Spark to investigate whether our system performed more efficiently than the existing system. Our system demonstrated performance improvement by a factor of 3.04 on Word Count, and by a factor of 2.64 on BiGram Count as compared to Spark. Then, we compared the two data exchange algorithms, Fully-Connected and Pairwise. The results of experiments with BiGram Count showed that Fully-Connected without RDMA was 13% more efficient than Pairwise with RDMA. We conclude that it is necessary to overlap map and shuffle phases to gain performance improvement. The reason of the relatively small percentage of improvement can be attributed to the time-consuming insertions of key-value pairs into the hash-map in the map phase.