Chaos: scale-out graph processing from secondary storage

Chaos: scale-out graph processing from secondary storage
复制标题

DOI:
10.1145/2815400.2815408
复制
发表时间:
2015-10
期刊:
Proceedings of the 25th Symposium on Operating Systems Principles
影响因子:
--
通讯作者:
Amitabha Roy;Laurent Bindschaedler;Jasmina Malicevic;W. Zwaenepoel
Amitabha Roy;Laurent Bindschaedler;Jasmina Malicevic;W. Zwaenepoel
中科院分区:
其他
文献类型:
--
作者:
Amitabha Roy;Laurent Bindschaedler;Jasmina Malicevic;W. Zwaenepoel

文献摘要

被引文献

相似文献

混乱量表图形处理从辅助存储到集群中的多台计算机。从辅助存储中进行处理图的早期系统仅限于一台机器,因此受单台计算机上存储系统的带宽和容量的限制。混乱仅受整个集群中所有存储设备的总带宽和容量的限制。混乱构建在X-Stream引入的流隔板上,以实现对存储的顺序访问,但并行地将流隔板的执行执行。混乱以三种方式是新颖的。首先,用于顺序存储访问的混乱分区,而不是用于局部和负载平衡,从而导致预处理时间较低。其次,混乱在整个群集上均匀地分发图形数据,并且没有试图实现局部性,这是基于小型群集网络带宽远远超过存储带宽的观察结果。第三,混乱使用偷窃工作允许多台机器在单个分区上工作,从而在运行时实现了负载余额。在性能缩放方面,在32台机器上,混乱平均只需长1.61倍以比单个机器上的32倍处理图32倍。在容量缩放方面,混乱能够处理1万亿个边缘的图形,代表16 TB的输入数据,这是一个小商品集群上图形处理能力的新里程碑。
Chaos scales graph processing from secondary storage to multiple machines in a cluster. Earlier systems that process graphs from secondary storage are restricted to a single machine, and therefore limited by the bandwidth and capacity of the storage system on a single machine. Chaos is limited only by the aggregate bandwidth and capacity of all storage devices in the entire cluster. Chaos builds on the streaming partitions introduced by X-Stream in order to achieve sequential access to storage, but parallelizes the execution of streaming partitions. Chaos is novel in three ways. First, Chaos partitions for sequential storage access, rather than for locality and load balance, resulting in much lower pre-processing times. Second, Chaos distributes graph data uniformly randomly across the cluster and does not attempt to achieve locality, based on the observation that in a small cluster network bandwidth far outstrips storage bandwidth. Third, Chaos uses work stealing to allow multiple machines to work on a single partition, thereby achieving load balance at runtime. In terms of performance scaling, on 32 machines Chaos takes on average only 1.61 times longer to process a graph 32 times larger than on a single machine. In terms of capacity scaling, Chaos is capable of handling a graph with 1 trillion edges representing 16 TB of input data, a new milestone for graph processing capacity on a small commodity cluster.