G-thinker: a general distributed framework for finding qualified subgraphs in a big graph with load balancing

G-thinker: a general distributed framework for finding qualified subgraphs in a big graph with load balancing
复制标题

DOI:
10.1007/s00778-021-00688-z
复制
发表时间:
2021-08
期刊:
The VLDB Journal
影响因子:
--
通讯作者:
Da Yan;Guimu Guo;J. Khalil;M. Tamer Özsu;Wei-Shinn Ku;John C.S. Lui
Da Yan;Guimu Guo;J. Khalil;M. Tamer Özsu;Wei-Shinn Ku;John C.S. Lui
中科院分区:
其他
文献类型:
--
作者:
Da Yan;Guimu Guo;J. Khalil;M. Tamer Özsu;Wei-Shinn Ku;John C.S. Lui

文献摘要

相似文献

Finding from a big graph those subgraphs that satisfy certain conditions is useful in many applications such as community detection and subgraph matching. These problems have a high time complexity, but existing systems that attempt to scale them are all IO-bound in execution. We propose the first truly CPU-bound distributed framework calledG-thinkerfor subgraph finding algorithms, which adopts a task-based computation model, and which also provides a user-friendly subgraph-centric vertex-pulling API for writing distributed subgraph finding algorithms that can be easily adapted from existing serial algorithms. To utilize all CPU cores of a cluster, G-thinker features (1) a highly concurrent vertex cache for parallel task access and (2) a lightweight task scheduling approach that ensures high task throughput. These designs well overlap communication with computation to minimize the idle time of CPU cores. To further improve load balancing on graphs where the workloads of individual tasks can be drastically different due to biased graph density distribution, we propose to prioritize the scheduling of those tasks that tend to be long running for processing and decomposition, plus a timeout mechanism for task decomposition to prevent long-running straggler tasks. The idea has been integrated into a novelty algorithm for maximum clique finding (MCF) that adopts a hybrid task decomposition strategy, which significantly improves the running time of MCF on dense and large graphs: The algorithm finds a maximum clique of size 1,109 on a large and denseWikiLinksgraph dataset in 70 minutes. Extensive experiments demonstrate that G-thinker achieves orders of magnitude speedup compared even with the fastest existing subgraph-centric system, and it scales well to much larger and denser real network data. G-thinker is open-sourced at http://bit.ly/gthinker with detailed documentation.