Enhancing Fault Tolerance based on Hadoop Cluster
Peng Hu, Wei Dai · International Journal of Database Theory and Application · 2014
Failures happen for large scale distributed systems such as Hadoop clusters. Native Hadoop provides basic support for failure tolerance. For example, data blocks are replicated over several HDFS nodes, and Map or Reduce tasks would be re-executed if they fail. However, simply re-processing the whole task decreases the efficiency of job execution, especially when the task is almost done. To this end, we propose a fault tolerance mechanism to detect and then recover from failures. Specifically, instead of simply using a timeout configuration, we design a trust based method to detect failures in a fast way. Then, a checkpoint based algorithm is applied to perform data recovery. Our experiments shows that our method exhibits good performance and is proved to be efficient.