As both modern supercomputers and new generation scientific computing applications grow in size and complexity, the probability of system failure rises commensurately. Making parallel computing fault tolerant has become an increasingly important issue.
At PPL, we have research (and resultant software) on multiple schemes that support fault tolerance:(1) automated traditional checkpoint with restart on a different number of processors (2) an in-remote-memory (or in-local-disk) double checkpointing scheme for online fault recovery (3) an ongoing research project that avoids sending all processors back to their checkpoints, and speeds up restart (4) a proactive fault avoidance scheme that evacuates a processor while letting the computation continue.
On-disk checkpoint/restart is the simplest approach. It requires synchronization in the system. In AMPI, this is invoked by a collective call. All the data is saved as disk files on a file server. When a processor crashes, the application is killed and then restarted from its checkpoint. This approach allows the scheduler more flexibility, since a job can be restarted on a different number of processors.
Double checkpoint-based fault tolerance for Charm++ and AMPI.
The above traditional disk-based method of dealing with faults is to checkpoint the state of the entire application periodically to reliable storage and restart from the recent checkpoint. The recovery of the application from faults involves (often manually) restarting applications on all processors and having it read the data from disks on all processors. The restart can therefore take minutes after it has been initiated. Traditional strategies often require that the failed processor can be replaced so that the number of processors at checkpoint-time and recovery-time are the same. Double in-memory checkpoint/restart is a scheme for fast and scalable in-memory checkpoint and restart. The scheme
- does not require any individual component to be fault-free because checkpoint data is stored in two copies on two different locations. (although not infalliable)
- provides fast and efficient checkpoint and restart using memory for checkpointing. It can take advantage of high performance network to speedup checkpointing process.
- at restart, it provides fault tolerance support for both cases with or without backup processors. If there are no backup processors, the program can continue to run on the remaining processors without performance penalty due to load imbalance.
- the memory-based method is useful for applications where the memory footprint is small at the checkpoint state, while a variation of this scheme --- in-disk checkpoint/restart can be applied to applications with large memory footprint. It uses local scratch disk for storing checkpoints, which is much more efficient than using centralized reliable file server.
- provides asynchronous checkpointing support for applications have moderate checkpoint size to relieve the network pressure. There are two phases in the support. In the first phase, after reaching a global synchronization point, each node will store its checkpoint in the main memory. After every node safely stores the newest copy of the checkpoint, they will resume to the normal execution. Then, the Charm++ runtime system will take the responsibility to send the checkpoints of each node to another node as a double copy at the appropriate time: when the applications are busy computing and leave the network free.
The checkpoint and restart time of the leanmd benchmark is shown below (on BG/P intrepid). It runs well and could scale to 32K cores.
For applications (wave2D and ChaNGa) uses asynchronous checkpointing, the checkpoint time is further reduced as shown below (on Trestles).FTL-Charm++:is a sender based pessimistic fault tolerant protocol for Charm++ and AMPI. The sender saves a copy of the message in its memory while sending a message to an object on another processor. It sends a ticket request to the receiver, which replies with a ticket. A ticket is the sequence in which the receiver will process this message. The sender saves the ticket in its memory and sends the message to the receiver. For messages sent to an object on the same processor, the ticket number is fetched through a function call and along with the sender and receiver's id is saved on a buddy processor. This allows message logging to work along with virtualization. While restarting after a crash, this protocol can spread the objects on the restarted processor among other processors. This provides for much faster restarts compared to other message logging protocols where all the work of the restarted processor is done on one processor. Checkpoint based protocols, which rollback all processors to their previous checkpoint when one processor crashes, obviously cannot redistribute work from the restarted processor as all processors are busy. So our protocol is unique in allowing restarts much faster than the time between the crash and the previous checkpoint. Work is going on to combine the message logging protocol with migration of objects to allow us to do dynamic runtime load balancing.
Proactive Fault Tolerance: is a scheme for reacting to fault warnings for Charm++ and AMPI. Modern hardware and fault prediction schemes can predict many faults with a good degree of accuracy. We developed a scheme for reacting to such faults. When a node is warned that it might crash, the charm++ objects on it are migrated away. The runtime system is also changed so that message delivery can continue seamlessly even if the warned node crashes. Reduction trees are also modified to remove the warned node from the tree so that its crash does not cause the tree to become disconnected. A node in the tree is replaced by one of its children. If the tree becomes too unbalanced, the whole reduction tree can be recreated from scratch. The proactive fault tolerance protocol does not require any extra nodes, it continues working on the remaining ones. It can deal with multiple simultaneous warnings, which is useful for real life cases where all the nodes in a rack are about to go down.
Below we show a run of NPB BT-MZ on 256 cores of Ranger with proactive fault tolerance support. The first load balancing at iteration 20 solves the initial load imbalancing of the application. At iteration 60, one core receives the signal of failure and objects on it migrate to other cores which creates a load imbalancing in the system. The problem is solved by applying load balancing once more.
Papers / Talks
-
16-202016
Phd Thesis -
16-152016
PaperFlipBack: Automatic Targeted Protection Against Silent Data Corruption
- Xiang Ni
- Laxmikant Vasudeo Kale
-
15-162015
PaperA Fault-Tolerance Protocol for Parallel Applications with Communication Imbalance
- Esteban Meneses
- Laxmikant Vasudeo Kale
-
14-312014
TalkScalable Replay with Partial-Order Dependencies for Message-Logging Fault Tolerance
- Jonathan Lifflander
- Esteban Meneses
- Harshitha Menon
- Phil Miller
- Sriram Krishnamoorthy
- Laxmikant Vasudeo Kale
-
14-212014
PaperScalable Replay with Partial-Order Dependencies for Message-Logging Fault Tolerance
- Jonathan Lifflander
- Esteban Meneses
- Harshitha Menon
- Phil Miller
- Sriram Krishnamoorthy
- Laxmikant Vasudeo Kale
-
14-202014
PaperUsing Migratable Objects to Enhance Fault Tolerance Schemes in Supercomputers
- Esteban Meneses
- Xiang Ni
- Gengbin Zheng
- Celso Mendes
- Laxmikant Vasudeo Kale
-
14-022014
PaperEnergy Profile of Rollback-Recovery Strategies in High Performance Computing
- Esteban Meneses
- Osman Sarood
- Laxmikant Vasudeo Kale
-
13-502013
TalkA ‘Cool’ Way of Improving the Reliability of HPC Machines
- Osman Sarood
- Esteban Meneses
- Laxmikant Vasudeo Kale
-
13-252013
PaperA ‘Cool’ Way of Improving the Reliability of HPC Machines
- Osman Sarood
- Esteban Meneses
- Laxmikant Vasudeo Kale
-
13-242013
PaperACR: Automatic Checkpoint/Restart for Soft and Hard Error Protection
- Xiang Ni
- Esteban Meneses
- Nikhil Jain
- Laxmikant Vasudeo Kale
-
13-172013
Phd Thesis -
13-062013
TalkAdoption Protocols for Fanout-Optimal Fault-Tolerant Termination Detection
- Jonathan Lifflander
- Phil Miller
- Laxmikant Vasudeo Kale
-
12-462012
PaperAdoption Protocols for Fanout-Optimal Fault-Tolerant Termination Detection
- Jonathan Lifflander
- Phil Miller
- Laxmikant Vasudeo Kale
-
12-432012
TalkAssessing Energy Efficiency of Fault Tolerance Protocols for HPC Systems
- Esteban Meneses
- Osman Sarood
- Laxmikant Vasudeo Kale
-
12-372012
PaperAssessing Energy Efficiency of Fault Tolerance Protocols for HPC Systems
- Esteban Meneses
- Osman Sarood
- Laxmikant Vasudeo Kale
-
12-322012
PaperHiding Checkpoint Overhead in HPC Applications with a Semi-Blocking Algorithm
- Xiang Ni
- Esteban Meneses
- Laxmikant Vasudeo Kale
-
12-302012
TalkA Message-Logging Protocol for Multicore Systems
- Esteban Meneses
- Xiang Ni
- Laxmikant Vasudeo Kale
-
12-142012
PaperA Message-Logging Protocol for Multicore Systems
- Esteban Meneses
- Xiang Ni
- Laxmikant Vasudeo Kale
-
12-122012
PaperA Scalable Double In-memory Checkpoint and Restart Scheme towards Exascale
- Gengbin Zheng
- Xiang Ni
- Laxmikant Vasudeo Kale
-
12-092012
TalkComposable and Modular Exascale Programming Models with Intelligent Runtime Systems
- Laxmikant Vasudeo Kale
-
12-042012
PaperA Scalable Double In-memory Checkpoint and Restart Scheme towards Exascale
- Gengbin Zheng
- Xiang Ni
- Esteban Meneses
- Laxmikant Vasudeo Kale
-
11-382011
TalkDynamic Load Balance for Optimized Message Logging in Fault Tolerant HPC Applications
- Esteban Meneses
- Greg Bronevetsky
- Laxmikant Vasudeo Kale
-
11-302011
PaperDesign and Analysis of a Message Logging Protocol for Fault Tolerant Multicore Systems
- Esteban Meneses
- Xiang Ni
- Laxmikant Vasudeo Kale
-
11-262011
PaperDynamic Load Balance for Optimized Message Logging in Fault Tolerant HPC Applications
- Esteban Meneses
- Greg Bronevetsky
- Laxmikant Vasudeo Kale
-
11-042011
PaperEvaluation of Simple Causal Message Logging for Large-Scale Fault Tolerant HPC Systems
- Esteban Meneses
- Greg Bronevetsky
- Laxmikant Vasudeo Kale
-
10-022010
PaperTeam-based Message Logging: Preliminary Results
- Esteban Meneses
- Celso Mendes
- Laxmikant Vasudeo Kale
-
09-092009
PaperTowards a Framework for Abstracting Accelerators in Parallel Applications: Experience with Cell
- David Kunzman
- Laxmikant Vasudeo Kale
-
06-122006
PaperA Fault Tolerance Protocol with Fast Fault Recovery
- Sayantan Chakravorty
- Laxmikant Vasudeo Kale
-
06-112006
PaperProactive Fault Tolerance in MPI Applications via Task Migration
- Sayantan Chakravorty
- Celso Mendes
- Laxmikant Vasudeo Kale
-
06-042006
PaperHPC-Colony: Services and Interfaces for Very Large Systems
- Sayantan Chakravorty
- Celso Mendes
- Laxmikant Vasudeo Kale
- Terry Jones
- Andrew Tauferner
- Todd Inglett
- Jose Moreira
-
06-032006
PaperPerformance Evaluation of Automatic Checkpoint-based Fault Tolerance for AMPI and Charm++
- Gengbin Zheng
- Chao Huang
- Laxmikant Vasudeo Kale
-
04-142004
PaperProactive Fault Tolerance in Large Systems
- Sayantan Chakravorty
- Celso Mendes
- Laxmikant Vasudeo Kale
-
04-072004
MS Thesis -
04-062004
PaperFTC-Charm++: An In-Memory Checkpoint-Based Fault Tolerant Runtime for Charm++ and MPI
- Gengbin Zheng
- Lixia Shi
- Laxmikant Vasudeo Kale
-
04-032004
PaperA Fault Tolerant Protocol for Massively Parallel Systems
- Sayantan Chakravorty
- Laxmikant Vasudeo Kale