IP Library Granted Patent US 9,098,439
Granted Patent B2
US 9,098,439 · App. 13/344,313 · Granted Aug 4, 2015

Providing a fault tolerant system in a loosely-coupled cluster environment using application checkpoints and logs

View Patent ↗
Loading inventors, assignments & file history…
Monitor This Case
Get email alerts when status or documents change.
Order Certified Copies
Most orders are placed with the USPTO same day — all within 24 business hours.
Order via The Patent Place →
Pre-filled with this patent's details
Quick Facts
Patent No.
US 9,098,439
App. No.
13/344,313
Granted
Aug 4, 2015
Kind
B2
Abstract

An approach to providing failure protection in a loosely-coupled cluster environment. A node in the cluster generates checkpoints of application data in a consistent state for an application that is running on a first node in the cluster. The node sends the checkpoint to one or more of the other nodes in the cluster. The node may also generate log entries of changes in the application data that occur between checkpoints of the application data. The node may send the log entries to other nodes in the cluster. The node may similarly receive external checkpoints and external log entries from other nodes in the cluster. In response to a node failure, the node may start an application on the failed node and recover the application using the external checkpoints and external log entries for the application.

Claims (45)

1. A computer program product stored on a non-transitory computer-readable storage medium, the computer program product comprising instructions for:

generating a checkpoint of application data that is in a consistent state on a first node in a loosely-coupled cluster comprising a plurality of nodes separate from the first node, wherein the application data is data of an application running on the first node;

sending the checkpoint to a plurality of nodes of the loosely-couple cluster, the plurality of nodes comprising a plurality of failover nodes for the first node, wherein the checkpoint is sent directly from the first node to the plurality of failover nodes;

generating one or more log entries of changes in the application data between checkpoints of the application data;

sending the log entries directly to the plurality of failover nodes from the first node;

detecting a failure associated with the application running on the first node;

determining a failover node that comprises a most current checkpoint data for the application from among the plurality of failover nodes that contain checkpoint data for the application in response to detecting the failure on the first node; and

recovering the application running on the first node on the determined failover node that comprises the most current checkpoint data for the application.

2. The computer program product of claim 1 , further comprising receiving one or more external checkpoints for one or more second applications running on one or more of the plurality of nodes in the loosely-coupled cluster.

3. The computer program product of claim 1 , further comprising receiving one or more external log entries for one or more second applications running on one or more of the plurality of nodes in the loosely-coupled cluster.

4. The computer program product of claim 1 , further comprising monitoring for node failures on the plurality of nodes in the loosely-coupled cluster resulting in a failed node.

5. The computer program product of claim 4 , further comprising determining that the first node is a failover node for the failed node.

6. The computer program product of claim 5 , further comprising recovering, on the first node, one or more second applications on the failed node from an external checkpoint received by the first node.

7. The computer program product of claim 6 , further comprising further recovering the one or more second applications on the failed node from external log entries received by the first node.

8. The computer program product of claim 1 , further comprising creating one or more consistency groups that comprise a plurality of applications, and generating a checkpoint of consistency group data that is in consistent state for the plurality of applications in the consistency group.

9. The computer program product of claim 1 , further comprising registering the application running on the first node with an event manager running on the first node.

10. A system comprising:

a first recovery apparatus running on a first node in a loosely-coupled cluster comprising a plurality of nodes separate from the first node, the second node plurality of nodes comprising a plurality of failover nodes for the first node, the first recovery apparatus:

generating a checkpoint of application data that is in a consistent state on the first node, wherein the application data is data of the first application;

generating one or more log entries of changes in the application data between checkpoints of the application data;

a first replicating service running on the first node, the first replicating service:

sending the checkpoint directly to the plurality of failover nodes from the first node for storage in local storage of the failover nodes;

sending the log entries directly to the plurality of failover nodes from the first node for storage in the local storage of the failover nodes; and

a second recovery apparatus running on the second node in the loosely-coupled cluster, the second recovery apparatus:

monitoring for failure of the first node;

determining a second node that comprises a most current checkpoint data for the application from among the plurality of failover nodes that contain checkpoint data for the application in response to detecting a failure on the first node;

recovering, on the second node, the first application from the checkpoint sent to the second node; and

further recovering, on the second node, the first application from the log entries sent to the second node.

11. The system of claim 10 , wherein the second recovery apparatus is further configured to failback the first application to the first node in response to the first node recovering from the failure.

12. The system of claim 10 , wherein the first node is a first computer and the second node is a second computer separate from the first computer and connected to the first computer by a local area network connection.

13. The system of claim 10 , wherein the first recovery apparatus and the second recovery apparatus are registered to the loosely-coupled cluster.

14. A computer-implemented method comprising:

generating a checkpoint of application data that is in a consistent state on a first node in a loosely-coupled cluster comprising a plurality of nodes separate from the first node, wherein the application data is data of an application running on the first node;

sending the checkpoint to a plurality of nodes of the loosely-couple cluster, the plurality of nodes comprising a plurality of failover nodes for the first node, wherein the checkpoint is sent directly from the first node to the plurality of failover nodes, the checkpoint being saved in local storage of the failover nodes;

generating one or more log entries of changes in the application data between checkpoints of the application data;

sending the log entries directly to the plurality of failover nodes from the first node, wherein the log entries are saved in local storage of the failover nodes;

detecting a failure associated with the application running on the first node;

determining a failover node that comprises a most current checkpoint data for the application from among the plurality of failover nodes that contain checkpoint data for the application in response to detecting the failure on the first node; and

recovering the application running on the first node on the determined failover node that comprises the most current checkpoint data for the application.

15. The method of claim 14 , further comprising generating a second checkpoint of the application data that is in a consistent state on the first node.

16. The method of claim 15 , further comprising deleting the log entries in response the second checkpoint being generated.

17. The method of claim 14 , further comprising receiving one or more external checkpoints for one or more second applications running on one or more of the plurality of nodes in the loosely-coupled cluster.

18. The method of claim 17 , further comprising receiving one or more shared log entries for the one or more second applications running on the one or more of the plurality of nodes in the loosely-coupled cluster.

19. The method of claim 18 , further comprising determining that the first node is a failover node for a failed node.

20. The method of claim 19 , further comprising recovering, on the first node, one or more second applications on the failed node from the shared checkpoint and the one or more external log entries.

Assignments (6)
RELEASE OF SECURITY INTEREST Recorded Nov 20, 2020
From: WILMINGTON TRUST, NATIONAL ASSOCIATION
To: GLOBALFOUNDRIES INC.
Reel/Frame 054636/0001 →
SECURITY AGREEMENT Recorded Nov 29, 2018
From: GLOBALFOUNDRIES INC.
To: WILMINGTON TRUST, NATIONAL ASSOCIATION
Reel/Frame 049490/0001 →
CORRECTIVE ASSIGNMENT TO CORRECT THE RECEIVING PARTY DATA PREVIOUSLY RECORDED ON REEL 036331 FRAME 0044. ASSIGNOR(S) HEREBY CONFIRMS THE ASSIGNMENT. Recorded Oct 21, 2015
From: INTERNATIONAL BUSINESS MACHINES CORPORATION
To: GLOBALFOUNDRIES U.S.2 LLC
Reel/Frame 036953/0823 →
ASSIGNMENT OF ASSIGNOR'S INTEREST Recorded Oct 5, 2015
From: GLOBALFOUNDRIES U.S. 2 LLC; GLOBALFOUNDRIES U.S. INC.
To: GLOBALFOUNDRIES INC.
Reel/Frame 036779/0001 →
ASSIGNMENT OF ASSIGNOR'S INTEREST Recorded Aug 14, 2015
From: INTERNATIONAL BUSINESS MACHINES CORPORATION
To: GLOBALFOUNDRIES U.S. 2 LLC COMPANY
Reel/Frame 036331/0044 →
ASSIGNMENT OF ASSIGNOR'S INTEREST Recorded Jan 5, 2012
From: CHIU, LAWRENCE Y.; FAN, SHAN; LIU, YANG; MEI, MEI; MUENCH, PAUL H.
To: INTERNATIONAL BUSINESS MACHINES CORPORATION
Reel/Frame 027487/0896 →