Multi-SPMD Programming Model with YML and XcalableMP
233
7 Fault-Tolerance Features in the mSPMD Programming
Model
7.1 Overview and Implementation
As well as scalability and programmability, reliability is an important issue in
exascale computing. Since the number of components of an exascale supercomputer
should be tremendously large, it is evident that the mean time between failure
(MTBF) of a system decreases as the number of the system’s components increases.
Therefore, fault tolerance becomes essential for systems and applications. Here, we
develop a fault-tolerance mechanism in an mSPMD programming model, and its
development and execution environment. The fault tolerance in the mSPMD programming model can be realized without modifying applications’ source codes[13].
Figure 14 illustrates the fault-tolerant mechanism in the mSPMD programming
model. If the workflow scheduler can find an error in a task and execute the
task again on different nodes, then we can realize a fault-tolerance and resilience
mechanism automatically.
We have extended the OmniRPC-MPI described in Sect. 3.3 to detect errors
in remote programs and notify the errors to the YML workflow scheduler. For
these purposes, heartbeat messages between master and remote programs have been
introduced in the OmniRPC-MPI library. If an error is detected in a remote program,
then it is reported to the YML workflow scheduler as a return value of existing
APIs. The OmniRpcProbe(Request r) API has been designed to listen to the
status of a requested task in a remote program. This returns success if the remote
program sends a signal to indicate the requested task r has successfully finished. On
the other hand, if heartbeat messages from the remote program executing the task r
have stopped, OmniRpcProbe(Request r) returns fail.
The YML scheduler re-schedules the failed task if it receives fail signal
from the OmniRPC-MPI library. The re-scheduling method is simple; The YML
scheduler puts the failed task at the head of the “ready” task queue.
7.2 Experiments
We have performed some experiments to investigate the overhead of the fault
detection and the elapsed time when errors occur on a cluster shown in Table 5. The
BGJ method shown in Sect. 5 had been used. The size of a matrix is 20,480 × 20,480
and divided into
233
7 Fault-Tolerance Features in the mSPMD Programming
Model
7.1 Overview and Implementation
As well as scalability and programmability, reliability is an important issue in
exascale computing. Since the number of components of an exascale supercomputer
should be tremendously large, it is evident that the mean time between failure
(MTBF) of a system decreases as the number of the system’s components increases.
Therefore, fault tolerance becomes essential for systems and applications. Here, we
develop a fault-tolerance mechanism in an mSPMD programming model, and its
development and execution environment. The fault tolerance in the mSPMD programming model can be realized without modifying applications’ source codes[13].
Figure 14 illustrates the fault-tolerant mechanism in the mSPMD programming
model. If the workflow scheduler can find an error in a task and execute the
task again on different nodes, then we can realize a fault-tolerance and resilience
mechanism automatically.
We have extended the OmniRPC-MPI described in Sect. 3.3 to detect errors
in remote programs and notify the errors to the YML workflow scheduler. For
these purposes, heartbeat messages between master and remote programs have been
introduced in the OmniRPC-MPI library. If an error is detected in a remote program,
then it is reported to the YML workflow scheduler as a return value of existing
APIs. The OmniRpcProbe(Request r) API has been designed to listen to the
status of a requested task in a remote program. This returns success if the remote
program sends a signal to indicate the requested task r has successfully finished. On
the other hand, if heartbeat messages from the remote program executing the task r
have stopped, OmniRpcProbe(Request r) returns fail.
The YML scheduler re-schedules the failed task if it receives fail signal
from the OmniRPC-MPI library. The re-scheduling method is simple; The YML
scheduler puts the failed task at the head of the “ready” task queue.
7.2 Experiments
We have performed some experiments to investigate the overhead of the fault
detection and the elapsed time when errors occur on a cluster shown in Table 5. The
BGJ method shown in Sect. 5 had been used. The size of a matrix is 20,480 × 20,480
and divided into
