Recovery in Distributed Systems Using Optimistic Message Logging and Checkpointing
Willy Zwaenepoel, Johnson, D.B. · 1988
.... Abstract a process is logged on stable storage [5], and each process is occasionally checkpointed to stable stor-In a distributed system using message logging and age, but no coordination is required between the checkpointing to provide fault tolerance, there is checkpoints of different processes. Between received always a unique maximum recoverable system state, messages, the execution of each process is assumed to (, regardless of the message logging protocol used. The be deterministic. L proof of this relies on the observation that the set of The protocols used for message logging are typi-system states that have occurred during any single cally pessimistic. With these protocols, each message execution of a system forms a lattice, with the sets is synchronously logged as it is received, either by of consistent and recoverable system states as sublat- blocking the receiver until the message is logged [1, 6], tices. The maximum recoverable system state never or by blocking the receiver if it attempts to send a new decreases, and if all messages are eventually logged, message before this received message is logged [3]. the domino effect cannot occur. This paper presents Recovery based on pessimistic message logging is a general model for reasoning about recovery in such straightforward. A failed process is restarted from its a system and, based on this model, an efficient algo- last checkpoint, and all messages originally received rithm for determining the maximum recoverable sys- by this process since the checkpoint are replayed to it tem state at any time. This work unifies existing ap- f,'om the log in the same order as they were received Iuf proaches to fault tolerance based on message logging before the failure. The process reexecutes based on and checkpointing, and improves on existing methods these messages to its state at the time of the failure. for optimistic recovery in distributed systems. Messages sent by the process during recovery are ig-nored since they are duplicates of those sent before