Massively Parallel Sort-Merge Joins in Main Memory Multi-Core Database Systems

Massively Parallel Sort-Merge Joins in Main Memory Multi-Core Database Systems
复制标题

DOI:
10.14778/2336664.2336678
复制
发表时间:
2012-06-01
影响因子:
2.5
通讯作者:
Neumann, Thomas
Neumann, Thomas
中科院分区:
计算机科学2区
文献类型:
--
作者:
Albutiu, Martina-Cezara;Kemper, Alfons;Neumann, Thomas

文献摘要

被引文献

相似文献

在不久的将来,两个新兴的硬件趋势将主导数据库系统技术:每台服务器增加几TB的主存容量和大规模并行多核处理。当前数据库技术中的许多算法和控制技术都是为基于磁盘的系统设计的,其中I/O占主导地位。在这项工作中,我们重新审视了著名的排序合并连接,到目前为止,还没有在可扩展的大规模并行多核数据处理的研究重点,因为它被认为不如散列连接。我们设计了一套新的大规模并行排序合并(MPSM)加入算法的基础上部分分区排序。与经典的排序合并连接相反,我们的MPSM算法不依赖于难以并行化的最终合并步骤来创建一个完整的排序顺序。相反,它们在独立创建的运行上并行工作。这样,我们的MPSM算法就是NUMA仿射的,因为所有排序都是在本地内存分区上进行的。一个现代的32核机器与一TB的主内存上的广泛的实验评估证明了MPSM的大型主内存数据库与数十亿个对象的竞争力的性能。它在使用的核心数量上(几乎)线性扩展,并且明显优于竞争的哈希连接建议-特别是它比“尖端”Vectorwise并行查询引擎的性能高出四倍。
Two emerging hardware trends will dominate the database system technology in the near future: increasing main memory capacities of several TB per server and massively parallel multi-core processing. Many algorithmic and control techniques in current database technology were devised for disk-based systems where I/O dominated the performance. In this work we take a new look at the well-known sort-merge join which, so far, has not been in the focus of research in scalable massively parallel multi-core data processing as it was deemed inferior to hash joins. We devise a suite of new massively parallel sort-merge (MPSM) join algorithms that are based on partial partition-based sorting. Contrary to classical sort-merge joins, our MPSM algorithms do not rely on a hard to parallelize final merge step to create one complete sort order. Rather they work on the independently created runs in parallel. This way our MPSM algorithms are NUMA- affine as all the sorting is carried out on local memory partitions. An extensive experimental evaluation on a modern 32-core machine with one TB of main memory proves the competitive performance of MPSM on large main memory databases with billions of objects. It scales (almost) linearly in the number of employed cores and clearly outperforms competing hash join proposals - in particular it outperforms the "cutting-edge" Vectorwise parallel query engine by a factor of four.