A method of recovering application data
19 claims: 7 independent, 12 dependent
- 1インターコネクトによって接続された複数のノードを有するコンピュータ・システムにおいて障害を起こしたノードのメモリからアプリケーション・データを回復し、該アプリケーション・データを置換ノードに書き込む方法であって、 前記コンピュータ・システムのノードがアプリケーションを実行し、該アプリケーションはアプリケーション・データを生成し、該アプリケーションの最も最近の状態をノード・メモリに記憶し;前記ノードが障害を起こし;前記障害を起こしたノードのノード・メモリはその後、フェイルオーバー・メモリ・コントローラを使って制御され;前記フェイルオーバー・メモリ・コントローラは前記アプリケーション・データを前記障害を起こしたノードのノード・メモリから前記置換ノードのノード・メモリに前記インターコネクトを通じてコピー し 、 前記アプリケーションが走っている間であって前記ノード障害の前に、前記アプリケーションは、前記障害を起こしたノードおよび/または前記置換ノードのノード・メモリにおいて前記アプリケーションによって利用可能または利用不能であるノード・メモリの一部を登録する、 方法。
- 2前記アプリケーションが、前記利用可能または利用不能であるノード・メモリの一部を登録することを、前記障害を起こしたノードが内部的に前記アプリケーションにメモリを割り当てた後に行なう、 請求項1記載の方法。
- 3前記アプリケーションが前記一部を前記フェイルオーバー・メモリ・コントローラに登録する、請求項2記載の方法。
- 4前記障害を起こしたノードのノード・メモリが前記障害後に補助電力を供給され、補助電源は前記フェイルオーバー・メモリ・コントローラによって制御される、請求項1または2記載の方法。
- 5前記補助電源が、前記フェイルオーバー・メモリ・コントローラと一緒に設けられたバッテリーの形である、請求項1ないし4のうちいずれか一項記載の方法。
- 6前記補助電源が、前記障害の前およびあとに前記フェイルオーバー・メモリ・コントローラに給電する、請求項1ないし5のうちいずれか一項記載の方法。
- 7前記補助電源が、前記フェイルオーバー・メモリ・コントローラのインターコネクト接続に給電する、請求項1ないし6のうちいずれか一項記載の方法。
- 8前記補助電力は、前記障害を起こしたノードのプロセッサまたはノード・メモリ・コントローラのような別のコンポーネントを介してではなく、前記ノード・メモリに直接供給される、請求項1ないし7のうちいずれか一項記載の方法。
- 9前記フェイルオーバー・メモリ・コントローラは、前記障害を起こしたノードのプロセッサまたはノード・メモリ・コントローラのような別のコンポーネントを介してではなく、前記インターコネクトに直接結合される、請求項1ないし8のうちいずれか一項記載の方法。
- 10前記コンピュータ・システムの管理プロセスがノードをモニタリングし、前記ノード障害を検出し、前記置換ノードを同定し、前記フェイルオーバー・メモリ・コントローラに前記アプリケーション・データを前記障害を起こしたノードのノード・メモリから前記置換ノードのノード・メモリにコピーするよう命令する、請求項1ないし9のうちいずれか一項記載の方法。
- 11前記管理プロセスがまた、前記置換ノード上の前記アプリケーションを再開する、請求項10記載の方法。
- 12補助電源および/またはフェイルオーバー・メモリ・コントローラが単一のノードのためにまたは一群のノードのために設けられる、請求項1ないし11のうちいずれか一項記載の方法。
- 13それぞれが独自のノード・メモリおよびノード・メモリ・コントローラを有する、インターコネクトによって接続された複数のノードを有するコンピュータ・システム上でアプリケーションを走らせているときに、障害を起こしたノードからの回復において使うためのフェイルオーバー・メモリ・コントローラであって、 当該フェイルオーバー・メモリ・コントローラは前記障害を起こしたノードの前記メモリ・コントローラおよび/またはメモリに接続するよう動作可能であり;当該フェイルオーバー・メモリ・コントローラは、前記障害を起こしたノードのノード・メモリに記憶されているアプリケーション・データの、置換ノードのノード・メモリへの、前記インターコネクトを通じた転送を制御するよう構成されて おり 、 前記アプリケーションが走っている間であって前記ノード障害の前に、前記アプリケーションは、前記障害を起こしたノードおよび/または前記置換ノードのノード・メモリにおいて前記アプリケーションによって利用可能または利用不能であるノード・メモリの一部を登録する、 フェイルオーバー・メモリ・コントローラ。
- 14前記フェイルオーバー・メモリ・コントローラは、前記コンピュータ・システム内のノードから自律的である、請求項13記載のフェイルオーバー・メモリ・コントローラ。
- 15それぞれが独自のノード・メモリおよびノード・メモリ・コントローラを有する複数のノードと;それらのノード上でアプリケーションを実行するときに障害を起こしたノードからの回復において使うためのフェイルオーバー・メモリ・コントローラと;それらのノードおよび前記フェイルオーバー・メモリ・コントローラを接続するインターコネクトとを有するコンピュータ・システムであって、 前記フェイルオーバー・メモリ・コントローラは障害を起こしたノードの前記メモリ・コントローラおよび/またはメモリに接続するよう動作可能であり;前記フェイルオーバー・メモリ・コントローラは、前記障害を起こしたノードのノード・メモリに記憶されているアプリケーション・データの、置換ノードのノード・メモリへの、前記インターコネクトを通じた転送を制御するよう構成されて おり 、 前記アプリケーションが走っている間であって前記ノード障害の前に、前記アプリケーションは、前記障害を起こしたノードおよび/または前記置換ノードのノード・メモリにおいて前記アプリケーションによって利用可能または利用不能であるノード・メモリの一部を登録する、 コンピュータ・システム。
- 16前記フェイルオーバー・メモリ・コントローラが各ノードのために設けられる、請求項15記載のコンピュータ・システム。
- 17前記フェイルオーバー・メモリ・コントローラは、ノードの電力およびインターコネクト接続とは別個の電力接続およびインターコネクト接続を有する、請求項15または16記載のコンピュータ・システム。
- 18フェイルオーバー・メモリ・コントローラおよびそれぞれノード・メモリを含む複数のノードを有するコンピュータ・システム上で走るデーモンであって、前記複数のノードおよびフェイルオーバー・メモリ・コントローラはみなインターコネクトによって接続されており、 当該デーモンは前記ノード上でのアプリケーションの実行をモニタリングし、 当該デーモンはノード障害を検出し、置換ノードを同定し、前記フェイルオーバー・メモリ・コントローラに、前記障害を起こしたノードのノード・メモリから前記置換ノードのノード・メモリに前記インターコネクトを通じてアプリケーション・データをコピーするよう命令 し 、 前記アプリケーションが走っている間であって前記ノード障害の前に、前記アプリケーションは、前記障害を起こしたノードおよび/または前記置換ノードのノード・メモリにおいて前記アプリケーションによって利用可能または利用不能であるノード・メモリの一部を登録する、 デーモン。
- 19前記アプリケーションが前記複数のノードによって分散式に実行される、請求項1ないし12のうちいずれか一項記載の方法。
Independent claims19
33 paragraphs, as filed
The present invention relates to recovering data from a failed compute node, but not only in parallel computing environments and high performance computing (HPC) applications. The present invention finds particular use in the field of fault-tolerant distributed computing and focuses on exascale computers.
HPC systems typically run compute-intensive and other large-scale applications. Such HPC systems often provide a distributed environment with multiple processing units or "cores" in which independent sequences of events, such as executable processing threads or processes, can run autonomously in parallel.
Many different hardware configurations and programming models are applicable to HPC. A common approach to HPC today is a cluster system in which multiple nodes, each with one or more multi-core processors (or "chips"), are interconnected by a high-speed network. Each node is assumed to have its own memory area, which is accessible to all cores within that node. Cluster systems can be programmed by human programmers who write source code, leveraging existing code libraries to perform common functions. The source code is then compiled and can be executed by lower level executable code, such as a processor type with a particular instruction set, ISA (Instruction Set). Architecture (instruction set architecture)) level code, or assembly language specific to a particular processor. Often there is a final step in assembling or (in the case of virtual machines) the assembly code into executable machine code. The executable form of an application (sometimes simply referred to as the "executable") runs under the supervision of an operating system (OS) and uses the OS and libraries to control the hardware. The various layers of software used may be collectively referred to as the software stack.
As used in this article, the term "software stack" includes all the software needed to run an application and is basic level software (operating system or OS); eg, interconnects between nodes, disks or others. Includes libraries (also a type of system software) with hardware components and interfaces such as memory and the application itself. The currently running application can be seen as the top layer of the software stack above the system software.
Applications for multi-core computer systems are written in a regular computer language (such as C / C ++ or Fortran) with augmented libraries to allow programmers to take advantage of the parallel processing power of multiple cores. May be good. In this regard, it is customary to refer to the "process" that runs in the core. A (multithreaded) process may run across several cores in a multicore CPU, and each node may contain one or more CPUs. One such library is the Message Passing Interface (MPI), which uses a distributed memory model (each process is assumed to have its own memory area) and is interprocessed. Make communication easier. MPI allows groups of processes to be defined and distinguished, and includes routines for so-called "barrier synchronization". This is an important feature that allows multiple processes or processing elements to work together.
Alternatively, shared memory parallel programming allows all processes or cores to access the same memory or memory area. The shared memory model does not require you to explicitly specify the communication of data between processes (because changes made by one process are transparent to all other processes). However, it may be necessary to use a library that controls access to shared memory to ensure that only one process modifies the data at a time.
Exascale computer (ie 1 exaflop (10)<sup>18</sup>Floating-point arithmetic (HPC system) capable of sustaining performance per second) is expected to be deployed by 2020. Several national projects have been announced to develop exascale systems during this time frame. Petascale (currently the latest technology, about 10<sup>15</sup>The transition from (flops) to exascale is expected to require fundamental changes in hardware technology. There is no longer any further increase in processor clock frequency, so performance improvements result in parallelism or increased concurrency (possibly up to about 1 billion cores). The requirement to keep the power usage of an exascale system in an acceptable window means that low power (and low cost) components are likely to be used, resulting in mean time between failures for each component. The interval becomes shorter. In this way, exascale systems will contain far more components than today's state-of-the-art systems. And each component is more likely to fail than its counterpart today. Mean time between failures for exascale systems is likely to be measured in minutes (rather than daily in the current system).
Therefore, exascale software will in particular need improved tolerance to these failures and will need to be able to overcome component failures and continue to run. HPC applications are generally load-balanced carefully to ensure that work is distributed across all available compute cores, so that they can perform the work assigned to the failed node. It can be important that the replacement node be made available to the application (assigning this work to one or more of the remaining nodes that are already loaded can disrupt load balance and lead to significant performance degradation. Is high).
Figure 1 shows the process of replacing a failed node with a replacement node. This figure shows six nodes (nodes 0-5) in the system, and node 5 has failed and can no longer contribute to the execution of the application. Of course, in reality, much more nodes make up a computer system. A replacement node (node 6) is made available to the system and inserted to replace the failed node. Once the replacement node is assigned to the application, it also needs to be initialized with the data needed to continue execution (for example, the value of a variable calculated by the application).
<p num="0010"> The need to initialize replacement nodes is not new, and known initialization techniques include:</p><p num="0011"> -Resume the node from the checkpoint file. This method ensures that the data is initialized to the correct values. However, generating checkpoints is time consuming (because it involves copying large amounts of data to memory on another node or to a file on disk). Therefore, the data is generally checkedpointed on a regular basis and the intervals between checkpoints are relatively large. Therefore, when restoring data from a checkpoint, all calculations since the last checkpoint must be repeated (at least on the failed node, potentially globally). In this way, checkpointing has a double overhead in terms of time. First when you generate the checkpoint, and then again when you read and recalculate to restore from the checkpoint.</p><p num="0012"> -Interpolate the value on the failed node from the value of the equivalent data on the surviving node. This method is not always possible (for example, if each node is responsible for a separate region of space in the modeling algorithm, it is possible to simply interpolate the solution over that region from the values on that boundary. It is unlikely). Even if it is possible to interpolate the data on the failed node from the data on other nodes, doing so results in loss of accuracy.</p><p num="0013"> All of these prior art techniques have drawbacks and therefore it is desirable to provide an alternative way to initialize the replacement node.</p>
<p num="0014"> According to an embodiment of one aspect of the invention, in a computer system with multiple nodes connected by an interconnect, recovering application data from the memory of a failed node and writing the application data to a replacement node. In a method, a node of the computer system runs an application, which generates application data and stores the most recent state of the application in node memory; the node of the failed node. The memory is then controlled using a failover memory controller; the failover memory controller passes the application data from the node memory of the failing node to the node memory of the replacement node through the interconnect. A way to copy is provided.</p><p num="0015"> The inventors need to actually copy the data off the failed node without having to repeatedly copy the data out of the node (because such copying is time consuming and often unnecessary). We have come to realize that this is a method of obtaining the data (not the interpolated one that loses accuracy), preferably from the state immediately before the failure (not the one before a certain time).</p><p num="0016"> Embodiments of the invention propose a method of recovering data directly from the memory of a failed node in order to overcome the limitations of current state-of-the-art technology.</p><p num="0017"> Thus, according to an embodiment of the invention, if there are a plurality of nodes and a node fails, that node is replaced by a replacement node. The failover memory controller (in effect, the standby memory controller used when the node fails) takes control of the node memory of the failed node, causing the application data to fail. Copy from the node memory of the new node to the node memory of the replacement node. Therefore, this method can restore the application data, which is the most recent state of the application, and copy it to the replacement node, including the entire contents of the node memory used by the application from the failed node. can do. The replacement node may be a node that was not previously used in the execution of the application. The replacement node may be a reserved node in case of failure of the node running the application. The replacement node may be treated in the same way as any other node after it has been initialized.</p><p num="0018"> In parallel execution of an algorithm in a system with multiple nodes, application data can include, for example, the latest version of a portion of the algorithm calculated on that node.</p><p num="0019"> Not all node memory may be available for application data. Node memory may also be used by other processes, such as those in an operating system. Some parts of node memory may be reserved for other uses and therefore unavailable for application data. In such situations, the application may register a portion of the node memory, preferably with a failover memory controller. This registration may be that of the currently running node, and thus may allow copying only the correct section of failed node memory. Alternatively or additionally, this registration may be that of a replacement node and may ensure that the data is copied to the available section of the replacement node. Registration can be done at any time the application is running, but before a node failure. In a preferred embodiment, registration is done as soon as possible after the application has been allocated memory within the node.</p><p num="0020"> The portion of node memory that is registered can be either a portion of memory that is available for use by the application or a portion of memory that is not available (these are equivalent and therefore neither. The option will also allow you to determine the correct section and / or available replacement section of the failed node). The said portion of node memory can be a failed node, a replacement node, or a portion of the node memory of both the failed node and the replacement node. In situations where not all of the node memory is available for application data on one or both of those nodes, this embodiment will only copy the application data of the failed node to the replacement node and / Or you can guarantee that the data will be copied to the available portion of the node memory of the replacement node. Registration is preferably node-by-node, rather than as a pattern for all nodes used in the application, so each node can internally determine how memory is allocated to the application.</p><p num="0021"> In a preferred embodiment of the invention, the failed node is likely to be completely non-functional. This can be due to an electrical failure in the node itself, or in the communication fabric that connects the node to the rest of the system, or due to some other flaw. The power connection to memory may not be available (for example, an interconnect or memory controller failure may leave the power supply to the memory intact, but the contents may become inaccessible), but in many node failure scenarios. Power is likely to be lost.</p><p num="0022"> According to some embodiments of the invention, node failure can thus mean that there is no longer power supply to the node, especially to node memory. In such situations, it may be appropriate to provide an auxiliary power source that is independent of the rest of the system. Auxiliary power can power any part or all of the system, but preferably powers the node memory of the failed node. Preferably, the auxiliary power is automatically switched on in the event of a node failure. Auxiliary power can act to continuously maintain power to memory so that it is not erased by a power failure. The auxiliary power supply may be controlled by any suitable means, but is preferably controlled by a failover memory controller.</p><p num="0023"> The auxiliary power source according to the embodiment of the invention may be any suitable power source of any kind, but is preferably a battery (s). Auxiliary power supplies may be located anywhere in the system, but are preferably with a failover memory controller (which may be attached to or separately from the node). If there are multiple failover memory controllers, the embodiments may provide the same number of auxiliary power supplies, each linked to its own failover memory controller, or each may have several failures. A smaller number of auxiliary power supplies may be provided linked to the overmemory controller.</p><p num="0024"> The auxiliary power supply may be able to power the failover memory controller at any time. This can be before or after the node failure, preferably both. When the auxiliary power source is the power source for the failover memory controller, it is also preferred to be the power source for the failover memory controller's interconnect connection and any other such interconnect connection that may require power.</p><p num="0025"> In embodiments of the invention, auxiliary power may be delivered to the node's memory of the node either directly or through another component of the node, such as a processor or node memory controller. Preferably, the auxiliary power is supplied directly to the node memory of the failed node.</p><p num="0026"> The failover memory controller's connection to the node's interconnect can be direct or through another component, such as the failed node's processor or node memory controller. Good. Preferably, the failover memory controller is directly coupled to the interconnect.</p><p num="0027"> The computer system may have a (central) management process that monitors the nodes (eg, at the operating system level). In such a situation, the management process can detect the node failure, identify the replacement node, and send the application data to the failover memory controller from the node memory of the failed node to the node memory of the replacement node. It is desirable to instruct to copy to. It is also preferable that the management process restarts the application.</p><p num="0028"> In an embodiment of the invention, the system may have one auxiliary power supply or a plurality of auxiliary power supplies. Depending on the number of nodes and the number of auxiliary power supplies in the system, each auxiliary power supply may be provided for a single node or for a group of nodes.</p><p num="0029"> Similarly, the system may have one failover memory controller or multiple failover memory controllers. Depending on the number of nodes in the system and the number of failover memory controllers, each failover memory controller may be provided for a single node or for a group of nodes.</p><p num="0030"> According to an embodiment of another aspect of the invention, when run on a computer system with multiple nodes connected by an interconnect, each with its own node memory and node memory controller. A failover memory controller for use in recovery from a failed node, which can operate to connect to said memory controller and / or memory on the failed node. The failover memory controller is configured to control the transfer of application data stored in the node memory of the failed node to the node memory of the replacement node through said interconnect. A failover memory controller that has been provided is provided.</p><p num="0031"> This aspect relates to a failover memory controller that may be provided as part of said computer system. As outlined above, there may be multiple failover memory controllers, one per node. In an embodiment of the invention, the only (or each) failover memory controller is autonomous with respect to some or all of the nodes in the computer system, including the failed node, preferably all of the nodes. It may be. However, a failover memory controller can have the ability to control parts of a node when needed.</p><p num="0032"> According to an embodiment of another aspect of the invention, with multiple nodes, each with its own node memory and node memory controller; from a failed node when running an application on those nodes. A computer system that has a failover memory controller for use in recovery; an interconnect that connects those nodes and the failover memory controller, and the failover memory controller fails. Can operate to connect to the memory controller and / or memory of the node; the failover memory controller is a replacement node for application data stored in the node memory of the failed node. A computer system is provided that is configured to control a transfer through said interconnect to a node memory of.</p><p num="0033"> The failover memory controller is preferably directly connected to the memory of the failed node and to the memory controller of the failed node. For example, the node memory controller leaves the node memory in a consistent state (this is possible, as it is possible that the node memory controller has failed), and then This is to ensure that control is passed to the failover controller.</p><p num="0034"> The computer system may be an HPC system or any other computer system with distributed memory and nodes connected by an interconnect.</p><p num="0035"> According to computer system embodiments, the failover memory controller may have power and interconnect connections that may or may not be separate from the node power and interconnect connections. Preferably, the failover memory controller has a power connection and / or an interconnect connection separate from the node power and interconnect connections. This allows more autonomy for the failover memory controller. In a preferred embodiment of the invention, the power supply to the failover memory controller, which may be the same as, but not necessarily the same as the auxiliary power supply, is independent of the operating state of all nodes. Therefore, it does not depend on whether any particular node is still operating to supply power thereafter. In most cases, the failover memory controller may use the same power supply as the rest of the computer system in the end. The embodiments of the invention are not specifically designed to be used in situations where the entire system loses power.</p><p num="0036"> According to an embodiment of another aspect of the invention, a daemon running on a computer system having a failover memory controller and multiple nodes, each containing node memory (in the background, not under user control). A computer program that runs as a process), the plurality of nodes and the failover memory controller are all connected by interconnects, and the daemon monitors the execution of applications on the nodes. Detects a node failure, identifies the replacement node, and orders the failover memory controller to copy application data from the failed node's node memory to the replacement node's node memory through the interconnect. , A daemon is provided.</p><p num="0037"> In effect, this daemon can provide the management process described above. Therefore, this daemon can be run automatically as part of the operating system.</p><p num="0038"> According to a more general programming aspect, when loaded into a computing device such as a distributed computer system, the computing device performs a method step based on any or any combination of the above method definitions. A program is provided that is configured to do so.</p><p num="0039"> The features and sub-features of any of the various aspects of the invention may be freely combined. For example, a preferred embodiment of a failover memory controller and / or computer system may be configured to incorporate functionality corresponding to one or more preferred features of the method.</p><p num="0040"> The present invention can be implemented in computer hardware, firmware, software or a combination thereof. Embodiments can be implemented as a computer program or computer program product, i.e. a computer program tangibly embodied in an information carrier, eg, in a non-transient machine-readable storage device or in a propagating signal. ..</p><p num="0041"> A computer program can be in the form of a computer program part or two or more computer programs, and can be written in any form of programming language, including compiled or interpreted languages, libraries, stand-alone. It can be deployed in any form, including as a program or as a module, component, subroutine or other unit suitable for use in a data processing environment.</p><p num="0042"> The method steps of the present invention can be performed by a programmable processor that executes a computer program to perform the functions of the present invention by acting on input data to produce output.</p><p num="0043"> The present invention has been described using individual embodiments. Other embodiments are within the appended claims. For example, the steps of the invention can be performed in different orders and still achieve the desired results.</p><p num="0044"> Devices based on preferred embodiments are described as operational or arranged to perform certain functions. This configuration or arrangement can be by the use of hardware or middleware or any other suitable system. In a preferred embodiment, the configuration or placement is software.</p><p num="0045"> The present invention will now be described with reference to the individual, non-limiting embodiments shown in the drawings.</p>
<figref num="1">It is a schematic diagram which shows the replacement of the failed node.</figref><figref num="2">It is a flowchart which draws the general embodiment of this invention.</figref><figref num="3">It is an overview of the restoration method based on the embodiment of the invention.</figref><figref num="4">a is a flowchart of node failure recovery according to the prior art, and b is a flowchart of recovery of a failed node based on the embodiment of the invention.</figref><figref num="5">It is a schematic diagram which shows the movement of data to a replacement node based on the embodiment of the invention.</figref><figref num="6">It is a figure which shows the schematic layout of the component in a node and the link which comes out from the node based on the embodiment of the invention.</figref>
FIG. 2 is a flowchart showing a general embodiment of the invention. In step S10, the compute node runs the application. In step S20, application data is generated and the most recent state of the application is stored in node memory. If the node fails in step S30, the failover memory controller controls the node memory of the failed node in step S40. Finally, in step S50, the failover memory controller copies data from the memory of the failed node to the memory of the replacement node.
Since the file-over memory controller is only used after a failure, application data is recovered from the failed node without performing redundant tasks such as data storage. Thus, the method is post-reactive. In addition, a single failover memory controller can interact with memory on multiple nodes if desired. Embodiments of the invention can address unexpected failures in applications running in parallel by copying the entire application data from one node to another.
The method proposed in the embodiments of the invention involves adding a failover memory controller to a computer system. This part can overlap with the functionality of the memory controller within one or more nodes or within each node. It powers the failover memory controller and maintains power to the memory of the failed node, allowing the data in that memory to be recovered and then transferred to the replacement node. May be attached to an auxiliary / backup power supply. A node failure management process may be provided. For example, following a failure on one node, the management process detects the failure, identifies the replacement node, and connects the failover memory controller to the memory of the failing node and replaces its contents. Instruct to copy directly to the node's memory. This reduces the time required to reinitialize the application after a failure (compared to standard checkpointing) and minimizes the amount of computation that must be repeated (in the processor registers). The data is not recovered, so a small amount of calculations still need to be repeated).
There are two main ways in which a replacement node can be assigned. First, the application can be launched with more nodes (possibly an extra 10% "spare" node) allocated than it actually needs. At that time, if a failure is detected in one of the nodes running the application (either by the application itself or through some monitoring software framework, such as the management process), the spare node is reserved and waits. Will be there.
Alternatively (possibly preferably) the system job scheduler can maintain a pool of spare nodes that can be assigned to any running application. In that case, following the detection of the node failure, the application (or monitoring framework) contacts the job scheduler and requests access to one of the spare nodes. In this scenario, the job scheduler is responsible for ensuring that a sufficiently large pool of spare nodes remains available when a new job begins to run.
FIG. 3 shows a specific embodiment of a method for recovery from a spoiled node. Figure 3 shows two of the many nodes 20 in the computer system 10. The failed (or failed) node is shown as node 20A on the left and the new / replacement node is shown as node 20B on the right. Each node contains memory 60 and memory controller 50. Other node parts are known to those of skill in the art and are not shown. The battery 40 and failover memory controller 30 are located separately from both nodes, possibly on the same physical support. The battery is shown linked to the failed node 20A memory and failover memory controller 30. Battery 40 does not need to be linked to memory 60 on replacement node 20B. That node is still connected to the power supply. The failover memory controller 30 is also linked to memory 60 on replacement node 20B through the network / interconnect.
Therefore, the failover memory controller and battery are additional parts of the computer system. In the event of a node failure, the battery powers the memory on the failed node, and the failover memory controller (which is also powered by the battery) recovers the contents of that memory and the failed node. Transfer to a new node that is intended to replace. Manager 70 controls this process (identifying the failed node, pointing the failover memory controller to the failed node, and specifying where memory should be copied).
The flowcharts a and b in FIG. 4 show the recovery of a failed node of the prior art without recovery and the recovery of a node using the failover memory controller of the invention embodiment that transfers the contents of memory to a new host node. Is shown.
If no recovery is provided, the application is started in S110, there is a node failure in step S120, and in step S130 the contents of memory on the failed node are lost to the application, as shown in a in FIG.
In Figure 4b, the file over memory controller is used to recover the data. Immediately after starting in step S210, in step S220, what part of node memory is being used (or equivalent, which part is reserved for other uses and for the application). It is necessary to register (is it not available). This is because some part of the total memory on the CPU is used by other processes, such as the operating system (OS). Equivalent processes are likely to be running on the replacement node, and overwriting the memory used by those processes would not be desirable (it is likely to cause problems on the new node). , Possibly even a new node can fail). Thus, only the memory directly used by the application should be transferred after a failure, so it needs to be registered at an early stage.
Following this registration, the application will continue to run until there is a node failure S230. In the meantime, the manager daemon monitors the health of each node (a daemon already exists to do this, and further discussion is omitted here). Following the failure detection in step S240, the manager assigns a new host node to the application in step S250 (to replace the failed node) and initiates the process of launching the application on that node. (S270), notify the failover memory controller (S260). Notifications can be in parallel with other actions. In the meantime, the application can continue on the remaining nodes (although execution is likely to eventually be held at the sync point, waiting for data from the failed node). In restore step S280, power is maintained by the auxiliary power supply to the memory of the failed node, and the failover memory controller copies the data from the previously registered section of memory and then on the replacement node. Transfer directly to memory (including memory controller notifications on the new node). Once the memory has been successfully copied to the replacement node, the management daemon can resume execution on this node.
The failover controller can be responsible for one or more nodes. The above process requires that the failover memory controller be directly connected to the memory and memory controller on each node in charge and have access to the network. The battery must be connected to the memory on the node (but not necessarily to the memory controller) and may power the failover memory controller's network connection.
FIG. 5 shows the movement of data from memory 60 on the failed node 20A to memory 60 on the replacement node 20B. Only application data is recovered (because the replacement node has its own instance of other data, such as the OS). This configuration requires the application to register the memory used by the application with the failover memory controller at startup. However, this is not the case if there are other methodologies or if there is a set portion of memory available for the application.
Failover memory controllers can be added to each node on a one-to-one basis (possibly on the same circuit board, but act autonomously, including their own network connectivity and power supply), or each with a group of processors. It can be implemented as separate components (s) within the system that are responsible for memory recovery (possibly the entire system). Embodiments of the invention work equally well in any of these possibilities.
Figure 6 depicts the physical relationships between the components of the system (differences from prior art are shown boldly). This figure shows the failed node, the management daemon, and the new node.
The failed (or old) node has a memory controller 50, a memory 60, a CPU 80, and a hard disk drive (HDD) 90 with an external connection via the network connector 100. This node also contains a failover memory controller 30 and a battery 40. The failover memory controller 30 is directly connected to the memory controller 50, and the battery 40 is directly connected to the memory 60 and the failover memory controller 30. Network connection 100 allows connections to CPU 80 and failover memory controller 30 management daemon 70.
The battery and failover memory controllers in the thick line are shown to be located inside the old node, but they do not have to be. The failover memory controller and battery can also be optionally housed elsewhere in the system (connected to the network and bypassing the CPU and memory controller to access the node's memory. As long as one failover memory controller may be responsible for memory on more than one node.
Embodiments of the invention may have some or all of the following advantages: · The most recent possible state of the application (rather than a version from some distance in the past when standard checkpoints are used) is restored, so the amount of computation that needs to be repeated after a node failure. Is greatly reduced. -The accuracy of the solution is not lost. On the other hand, interpolating the solution from the surviving solution on the failed node can result in increased error.
In summary, according to preferred embodiments of the invention, a failover memory controller that functions autonomously from the processing units in the system and has access to the memory of those processing units is the most important distinguishing technical feature. A technical feature that also distinguishes the use of batteries or other auxiliary power sources to power the memory on the failed node.
The following additional notes will be further disclosed with respect to the embodiments including the above embodiments. (Appendix 1) A method of recovering application data from the memory of a failed node and writing the application data to a replacement node in a computer system with multiple nodes connected by an interconnect. A node in the computer system runs an application, which generates application data and stores the most recent state of the application in node memory; The node has failed; The node memory of the failed node is then controlled using the failover memory controller; The failover memory controller copies the application data from the node memory of the failed node to the node memory of the replacement node through the interconnect. Method. (Appendix 2) The application is running and prior to the node failure, the application has a portion of the node memory available or unavailable by the application in the node memory of the failed node and / or the replacement node. Register as The method described in Appendix 1. (Appendix 3) The method according to Appendix 2, wherein the application registers a part of the failover memory controller. (Appendix 4) The method of Appendix 1 or 2, wherein the node memory of the failed node is supplied with auxiliary power after the failure, and the auxiliary power is controlled by the failover memory controller. (Appendix 5) The method according to any one of Supplementary note 1 to 4, wherein the auxiliary power supply is in the form of a battery provided together with the failover memory controller. (Appendix 6) The method according to any one of Supplementary note 1 to 5, wherein the auxiliary power supply supplies power to the failover memory controller before and after the failure. (Appendix 7) The method according to any one of Supplementary note 1 to 6, wherein the auxiliary power supply supplies power to the interconnect connection of the failover memory controller. (Appendix 8) The auxiliary power is supplied directly to the node memory, not through another component such as the processor or node memory controller of the failed node, any one of Appendix 1-7. Item description method. (Appendix 9) The failover memory controller is any of Appendix 1-8, which is directly coupled to the interconnect, not through another component such as the processor or node memory controller of the failed node. The method described in item 1. (Appendix 10) The management process of the computer system monitors the node, detects the node failure, identifies the replacement node, and transfers the application data to the failover memory controller to the node memory of the failed node. The method according to any one of Supplementary note 1 to 9, wherein the replacement node is instructed to be copied to the node memory of the replacement node. (Appendix 11) 10. The method of Appendix 10, wherein the management process also restarts the application on the replacement node. (Appendix 12) The method of any one of Appendix 1 to 11, wherein the auxiliary power supply and / or failover memory controller is provided for a single node or for a group of nodes. (Appendix 13) To be used in recovery from a failed node when running an application on a computer system with multiple nodes connected by an interconnect, each with its own node memory and node memory controller. Failover memory controller of The failover memory controller can operate to connect to the memory controller and / or memory of the failed node; The failover memory controller is configured to control the transfer of application data stored in the node memory of the failed node to the node memory of the replacement node through the interconnect. Yes, Failover memory controller. (Appendix 14) The failover memory controller according to Appendix 13, wherein the failover memory controller is autonomous from a node in the computer system. (Appendix 15) With multiple nodes, each with its own node memory and node memory controller; With a failover memory controller for use in recovery from a failed node when running an application on those nodes; A computer system having those nodes and the interconnect connecting the failover memory controller. The failover memory controller can operate to connect to the memory controller and / or memory of the failed node; The failover memory controller is configured to control the transfer of application data stored in the node memory of the failed node to the node memory of the replacement node through the interconnect. Yes, Computer system. (Appendix 16) The failover memory controller is provided for each node. The computer system described in Appendix 15. (Appendix 17) The failover memory controller has a power connection and an interconnect connection separate from the power and interconnect connections of the node. The computer system described in Appendix 15 or 16. (Appendix 18) A daemon that runs on a computer system that has a failover memory controller and multiple nodes, each containing node memory, all of which are connected by an interconnect. The daemon monitors the execution of applications on the node and The daemon detects a node failure, identifies the replacement node, and sends application data to the failover memory controller through the interconnect from the node memory of the failed node to the node memory of the replacement node. Order to copy, daemon.
S10 node runs application Generates S20 application data and stores the latest state of the application in node memory S30 node failure? S40 failover memory controller controls node memory on the failed node S50 Failover Memory Controller copies data from the memory of the failed node to the node memory of the replacement node S110 application start S120 node failure S130 The contents of memory on the failed node are lost to the application S210 application start S220 application registers memory with failover memory controller S230 node failure S240 node failure detected by daemon S250 Starts a new host node Recover memory through S260 failover memory controller S270 Start application on new host S280 recovery 20A node 20B new node 30 Failover memory controller 40 battery 50 memory controller 60 memory 70 Manager (administrative daemon) 80 CPU 90 HDD 100 networks
6 sheets
Sheet 1 Sheet 2 Sheet 3 Sheet 4 Sheet 5 Sheet 6
Every citation, both ways
| Document | Relation | Office |
|---|---|---|
| JP2006309292A | Cites | Japan |
| JP62005759A | Cites | Japan |
6 members in 3 offices
Priority claims5
| Document | Office | Kind | Date |
|---|---|---|---|
| 14165974 | European Patent Office (EPO) | A | |
| 14165974 | European Patent Office (EPO) | A | |
| 141659748 | European Patent Office (EPO) | – | |
| 141659748 | – | – | – |
| EP20140165974 | – | – | – |
Members6
| Document | Office | Kind | |
|---|---|---|---|
| EP2937785A1 | European Patent Office (EPO) | A1 | |
| US2015309893A1 | United States of America | A1 | |
| JP2015210812A | Japan | A | |
| EP2937785B1 | European Patent Office (EPO) | B1 | |
| US9852033B2 | United States of America | B2 | |
| JP6398658B2This record | Japan | B2 |
9 legal events, as the office reported them to INPADOC
Over the term
Point at a mark for the eventEvents
| Event | Code | |
|---|---|---|
| Cancellation because of no payment of annual feesLAPS | LAPS | |
| Certificate of patent or registration of utility modelJAPANESE INTERMEDIATE CODE: R150R150 | R150 | |
| First payment of annual fees (during grant procedure)JAPANESE INTERMEDIATE CODE: A61A61 | A61 | |
| Written decision to grant a patent or to grant a registration (utility model)JAPANESE INTERMEDIATE CODE: A01A01 | A01 | |
| Decision of grant or rejection writtenTRDD | TRDD | |
| Request for written amendment filedJAPANESE INTERMEDIATE CODE: A523A521 | A521 | |
| Notification of reasons for refusalJAPANESE INTERMEDIATE CODE: A131A131 | A131 | |
| Report on retrievalJAPANESE INTERMEDIATE CODE: A971007A977 | A977 | |
| Written request for application examinationJAPANESE INTERMEDIATE CODE: A621A621 | A621 |
Numbers
- Publication
- 6398658
- Publication, DOCDB
- 6398658
- Publication, EPODOC
- JP6398658B
- Application
- 240195
- Application, DOCDB
- 2014240195
- Application, EPODOC
- JP20140240195
Titles2
- Japanese
- アプリケーション・データを回復する方法
- English
- How to recover application data
Classification
- CPC, 4
- G06F11/1658
- G06F11/2025
- G06F11/2028
- G06F11/203
- IPC, 1
- G06F11 20
