Distributed MapReduce engine with fault tolerance
Lixing Song, Shaoen Wu, Honggang Wang, Qing Yang · 2014
Hadoop is the de facto engine that drives current cloud computing practice. Current Hadoop architecture suffers from single point of failure problems: its job management lacks of fault tolerance. If a job management fails, even if its tasks remains still active on cloud nodes, this job loses all state information and has to restart from scratch. In this work, we propose a distributed MapReduce engine for Hadoop with the Distributed Hash Table (DHT) algorithm that drives the scalable peer-to-peer networks today. The distributed Hadoop engine provides the fault-tolerance capability necessary to support efficient job computation required in the cloud computing with numerous jobs running at a moment. We have implemented the proposed distributed solution into Hadoop and evaluated its performance in job failures under various network deployments.