Replication-Based Fault-Tolerance for Large-Scale Graph Processing

Replication-Based Fault-Tolerance for Large-Scale Graph Processing
复制标题

DOI:
10.1109/tpds.2017.2703904
复制
发表时间:
2018-07
影响因子:
5.3
通讯作者:
Rong Chen;Youyang Yao;Peng Wang;Kaiyuan Zhang;Zhaoguo Wang;Haibing Guan;B. Zang;Haibo Chen
Rong Chen;Youyang Yao;Peng Wang;Kaiyuan Zhang;Zhaoguo Wang;Haibing Guan;B. Zang;Haibo Chen
中科院分区:
计算机科学2区
文献类型:
--
作者:
Rong Chen;Youyang Yao;Peng Wang;Kaiyuan Zhang;Zhaoguo Wang;Haibing Guan;B. Zang;Haibo Chen

文献摘要

被引文献

相似文献

算法复杂性和数据集大小的增加使得许多图并行算法需要使用网络机器,这也使得容错成为机器规模增加的必要条件。然而,现有的大规模图并行系统通常采用分布式检查点机制来容错,这不仅会带来显著的性能开销,而且恢复时间也很长.本文观察到,为分布式图计算创建的顶点副本可以自然扩展,以快速在内存中恢复图状态。本文介绍了模仿者,一种新的容错机制,它支持廉价的维护顶点状态复制到他们的副本在正常的消息交换,并提供快速的内存中重建失败的顶点从副本在其他机器。模仿者已经在Cyclops上实现了边缘切割和PowerLyra上实现了顶点切割。在50个节点的EC-2类集群上进行的评估显示,对于Cyclops和PowerLyra,Imitator平均分别产生1.37%和2.32%的性能开销(范围从-0.6%至3.7%不等),并且可以在不到3.4秒的时间内从超过100万个顶点的故障中恢复。
The increasing algorithmic complexity and dataset sizes necessitate the use of networked machines for many graph-parallel algorithms, which also makes fault tolerance a must due to the increasing scale of machines. Unfortunately, existing large-scale graph-parallel systems usually adopt a distributed checkpoint mechanism for fault tolerance, which incurs not only notable performance overhead but also lengthy recovery time. This paper observes that the vertex replicas created for distributed graph computation can be naturally extended for fast in-memory recovery of graph states. This paper describes Imitator, a new fault tolerance mechanism, which supports cheap maintenance of vertex states by replicating them to their replicas during normal message exchanges, and provides fast in-memory reconstruction of failed vertices from replicas in other machines. Imitator has been implemented on Cyclops with edge-cut and PowerLyra with vertex-cut. Evaluation on a 50-node EC-2 like cluster shows that Imitator incurs an average of 1.37 and 2.32 percent performance overhead (ranging from −0.6 to 3.7 percent) for Cyclops and PowerLyra respectively, and can recover from failures of more than one million of vertices with less than 3.4 seconds.