11. Fault Tolerance in Large-Scale Scientific Computing
Patricia Hough, Victoria E. Howle · Society for Industrial and Applied Mathematics eBooks · 2006
Large-scale simulation is becoming an increasingly valuable tool in exploring and understanding complex scientific phenomena. Furthermore, the availability of parallel machines with tens or hundreds of thousands of processors makes it possible to conduct computational studies at a level of detail that was never before tractable. Key to the successful use of these machines is the ability of enabling technologies to keep pace with the challenges presented by the massive scale, including scalability, latency, and regular failures. It is on this last category that we focus. Birman [6] defines fault tolerance as the ability of a distributed system to run correctly through a failure. The system consists not only of the hardware platform, but also of the application software and any middleware upon which it depends. Fault tolerance has long been an active research area in the computer science community. The problem domain for this work has typically consisted of business-, mission-, or life-critical applications, which are characterized by their loosely coupled nature and their need to withstand arbitrary failures. The techniques developed for fault tolerance are usually based on replication of some sort and require little or no knowledge of the application. Therefore, they can be implemented in middleware and used by an application with minimal effort. The size and complexity of large-scale parallel computers and scientific applications continue to grow at a rapid rate. With that growth comes a corresponding increase in the frequency of failures and thus, a greater concern for fault tolerance. While scientific applications are important, they are notably less critical than those that fall in the traditional problem domain. Coupled with that characterization is the ever-present desire for high performance. The overarching goal, then, is to tolerate the most common failures while introducing only minimal overhead during failure-free execution. The traditional replication-based approaches, while transparent and easy to use, introduce an unacceptable amount of overhead in this setting.