Complete Process Recovery: Using Vector Time to Handle Multiple Failures in Distributed Systems
Golden G. Richard, Mukesh Kumar Singhal · 1997
Abstract: Distributed applications are generally structured as a set of communicating processes executing on multiple processors. As the number of processors in a distributed system and the running time of applications increases, the likelihood of processor failure increases. Without mechanisms for recovery, the failure of even a single processor can mandate restarting an entire application from scratch. Handling multiple failures is important because power failures, user error, etc. can often cause several machines to fail almost simultaneously. With appropriate recovery mechanisms, distributed applications can survive failures and complete restarts can be avoided. In this paper we present a set of distributed recovery techniques called CPR (Complete Process Recovery) which utilize vector time to handle failures and address both consistent state restoration and the associated message handling issues. The latter is important, since some recovery techniques delegate the handling of lost or duplicate messages to the message transport mechanism. This requires the ability to checkpoint the network layer (a serious restriction), since the network is generally unaware of the anomalous messages induced by process failure and recovery. CPR provides a comprehensive approach to process recovery which melds consistent state restoration and proper message handling. Our technique requires non-failed processes to roll back at most once in response to a single failure and has low message complexity. Processes required to rollback after a failure do so concurrently, which substantially decreases recovery delay after a failure has occurred.