Scalable and efficient graph traversal on high-throughput cluster

Scalable and efficient graph traversal on high-throughput cluster
复制标题

高吞吐量集群上可扩展且高效的图遍历

DOI:
10.1007/s42514-020-00056-3
复制
发表时间:
2021
影响因子:
0.9
通讯作者:
Sun Ninghui
Sun Ninghui
中科院分区:
--
文献类型:
--
作者:
Fan Dongrui;Cao Huawei;Wang Guobo;Nie Na;Ye Xiaochun;Sun Ninghui

文献摘要

相似文献

图是现代大数据应用中最重要的数据结构之一,广泛应用于各个领域。在众多图算法中,广度优先搜索(BFS)算法是解决图遍历问题的经典算法,也是Graph500基准测试的关键内核。在现代CPU架构上,图遍历在单节点系统上的实现已经取得了显着的改进。然而,由于资源利用率低和通信开销高,分布式集群上的图遍历性能较差且能源效率低下。高吞吐量集群(HTC)采用高吞吐量众核架构,具有高并发、实时性强、低功耗等特点。在这项工作中,我们提出了多种技术,包括异步虚拟环方法、线程缓存方案和顶点ID重新排序来解决上述问题并提高HTC上的BFS性能。我们系统地评估了优化的 BFS 算法,并在 72 个节点(2880 个核心)HTC 上实现了每秒 249.74 千兆遍历边(GTEPS)。与Graph500列表结果相比,优化后的算法在相同集群规模下实现了最高的节点效率,并且随着集群节点数量的增加,性能表现出弱线性可扩展性。在效率方面,HTC 上的平均性能为 3.47 GTEPS/节点,是 2019 年 11 月 Graph500 榜单上基于 CPU 的分布式系统中最好的。
Graph is one of the most important data structures in modern big data applications and is widely used in various fields. Among many graph algorithms, the Breadth-First Search (BFS) algorithm is a classic algorithm to solve the graph traversal problem and also the key kernel of Graph500 benchmark. On modern CPU architecture, the implementation of graph traversal on single-node systems has achieved significant improvement. However, due to the low resource utilization and high communications overhead, graph traversal on distributed clusters suffers from poor performance and energy inefficiency. High-throughput cluster (HTCs) adopt High-Throughput many-core architecture, which has the characteristics of high concurrency, strong real-time, and low-power consumption. In this work, we propose several techniques, including asynchronous virtual ring method, thread caching scheme and vertex ID reordering to solve above problems and improve BFS performance on HTCs. We systematically evaluate optimized BFS algorithm and achieve 249.74 giga-traversed edges per second (GTEPS) on 72 nodes (2880 cores) HTCs. Compared with results on Graph500 list, the optimized algorithm achieves the highest node efficiency under the same cluster scale and the performance shows weakly linear scalability as the number of cluster nodes increases. With regard to efficiency, the average performance on HTCs is 3.47 GTEPS/node, which is the best among CPU-based distributed systems on the November 2019 Graph500 list.