Optimizing Multiway Joins in a Map-Reduce Environment

Optimizing Multiway Joins in a Map-Reduce Environment
复制标题

DOI:
10.1109/tkde.2011.47
复制
发表时间:
2011-09
影响因子:
8.9
通讯作者:
F. Afrati;J. Ullman
F. Afrati;J. Ullman
中科院分区:
计算机科学2区
文献类型:
--
作者:
F. Afrati;J. Ullman

文献摘要

被引文献

相似文献

Map-Reduce的实现正被用于对非常大的数据执行许多操作。我们研究了在Map-Reduced环境中连接几个关系的策略。我们的新方法从标识“map-key”开始,这是一组属性,用于标识Map进程必须向其发送特定元组的Reduced进程。映射键的每个属性都有一个“共享”,这是它的值被散列到其中的存储桶的数量,以形成Reduce进程的标识符的一个组件。关系以有限的方式复制其元组,复制的程度取决于其模式中缺失的映射键属性的共享。在给定一定数量的归约过程的情况下,我们研究了份额的优化问题。给出了一种检测和修复变量错误地包含在映射键中的问题的算法。然后,我们考虑了两种重要的特殊情况:链连接和星形连接。在每种情况下,我们都能够确定映射键,并确定产生最少复制的共享。虽然我们提出的方法并不总是优于使用MAP-REDUE实现连接的传统方法,但在一些涉及大规模数据的重要情况下,我们的方法是成功的,包括:1)分析查询,其中非常大的事实表与较小的维表连接;2)查询涉及通过出局程度较高的图的路径的查询,例如Web或社交网络。
Implementations of map-reduce are being used to perform many operations on very large data. We examine strategies for joining several relations in the map-reduce environment. Our new approach begins by identifying the “map-key,” the set of attributes that identify the Reduce process to which a Map process must send a particular tuple. Each attribute of the map-key gets a “share,” which is the number of buckets into which its values are hashed, to form a component of the identifier of a Reduce process. Relations have their tuples replicated in limited fashion, the degree of replication depending on the shares for those map-key attributes that are missing from their schema. We study the problem of optimizing the shares, given a fixed number of Reduce processes. An algorithm for detecting and fixing problems where a variable is mistakenly included in the map-key is given. Then, we consider two important special cases: chain joins and star joins. In each case, we are able to determine the map-key and determine the shares that yield the least replication. While the method we propose is not always superior to the conventional way of using map-reduce to implement joins, there are some important cases involving large-scale data where our method wins, including: 1) analytic queries in which a very large fact table is joined with smaller dimension tables, and 2) queries involving paths through graphs with high out-degree, such as the Web or a social network.